schizophrenicmaniac commented on code in PR #73626:
URL: https://github.com/apache/airflow/pull/73626#discussion_r4121056316
##########
shared/observability/src/airflow_shared/observability/traces/__init__.py:
##########
@@ -248,6 +249,23 @@ def _load_exporter_from_env() -> SpanExporter:
return ep.load()()
+class _ForkSafeTracerProvider(TracerProvider):
+ """
+ ``TracerProvider`` that survives ``os.fork()``.
+
+ The SDK's ``after_in_child`` handler refreshes process-dependent resource
attributes on a
+ ``ThreadPoolExecutor``. A blocked detector is not bounded by the pool's
timeout -- leaving
+ ``get_aggregated_resources`` joins the pool threads -- and the default
service-instance
+ detector takes a module-level lock that a forked child inherits locked
with no owner. The
+ child then never returns from ``os.fork()``, which is how OTel stalls
every LocalExecutor
+ worker and DagFileProcessor child Airflow forks. Only the tracer lock is
refreshed here;
+ the child keeps the parent's process-dependent resource attributes rather
than deadlocking.
+ """
+
+ def _handle_fork(self) -> None:
+ self._tracers_lock = threading.Lock()
Review Comment:
Good catch, thanks. I reproduced it with a referenced MeterProvider and the
child hangs the same way. Added a `_ForkSafeMeterProvider` that
`get_otel_logger()` now installs, with the same tests as the tracer side.
##########
shared/observability/tests/observability/test_traces.py:
##########
@@ -360,3 +369,88 @@ def test_roundtrip_via_carrier(self):
span = trace.get_current_span(ctx)
assert get_task_span_detail_level(span) == 3
+
+
+_FORK_SCENARIO = textwrap.dedent(
+ """
+ import os
+ import signal
+ import time
+
+ from opentelemetry.sdk import resources as otel_resources
+
+ from airflow_shared.observability.traces import _ForkSafeTracerProvider
+
+ _ForkSafeTracerProvider() # registers the after_in_child handler under
test
+
+ with otel_resources._service_instance_id_lock:
+ pid = os.fork()
+ if pid == 0:
+ os._exit(0)
+
+ deadline = time.monotonic() + 10
+ status = None
+ while time.monotonic() < deadline:
+ waited_pid, waited_status = os.waitpid(pid, os.WNOHANG)
+ if waited_pid == pid:
+ status = waited_status
+ break
+ time.sleep(0.05)
+
+ if status is None:
+ os.kill(pid, signal.SIGKILL)
+ os.waitpid(pid, 0)
+ raise SystemExit("child never returned from os.fork()")
+ """
+)
+
+
+class TestForkSafeTracerProvider:
+ def test_handle_fork_refreshes_the_tracer_lock(self):
+ provider = _ForkSafeTracerProvider()
+ inherited_lock = provider._tracers_lock
+
+ provider._handle_fork()
+
+ assert provider._tracers_lock is not inherited_lock
+
+ @mock.patch("opentelemetry.sdk.trace._get_process_dependent_resource")
Review Comment:
Done, thanks.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]