Skip to content

Stop forked children re-exporting the parent's OTel metrics - #71855

Open
dstandish wants to merge 6 commits into
apache:mainfrom
astronomer:otel-fork-flush
Open

dstandish wants to merge 6 commits into
apache:mainfrom
astronomer:otel-fork-flush

Conversation

@dstandish

@dstandish dstandish commented Aug 19, 2026 •

Copy link
Copy Markdown
Contributor

The OTel metrics pipeline can end up with two writers on one cumulative series — one climbing,
one frozen — in three independent ways. This fixes all of them.

Several pipelines coexist in one process. Every MeterProvider owns a
PeriodicExportingMetricReader whose constructor starts a daemon export thread, and
shutdown_on_exit=False means nothing reaps one that gets replaced. Its instruments stop being
recorded to, so it republishes frozen totals under a different start_time_unix_nano for the life
of the process. Two separate things build more than one. stats.initialize() is called repeatedly:
BaseExecutor.__init__ and SchedulerJobRunner._execute each call it, with a
stats.incr("schedulerjob_start") in between that materialises the first pipeline. And
shared/observability is symlinked into both airflow-core and task-sdk, so Python loads it twice —
as airflow._shared.observability and airflow.sdk._shared.observability — and each copy holds
its own _factory/_backend. Importing airflow.sdk.serde initialises the second copy, which
airflow-core does during deserialization. A scheduler reaches three live pipelines on every start.

Fixed by building the pipeline once per process. Ownership is recorded on the provider itself, in
the OTel global registry, rather than in module state: a module-level guard only deduplicates
within its own copy, so it would leave two pipelines rather than one, and would leave the fork
cleanup below working only in whichever copy happened to build the pipeline. The OTel registry is
process-wide, which makes the provider the only place both copies can agree on. The marker records
the pid that built it, which is also what keeps a forked child from adopting a pipeline it does not
own.

A forked child re-exports the pipeline it inherited. fork() copies the parent's provider, and
the SDK registers register_at_fork(after_in_child=...) for every PeriodicExportingMetricReader,
so the child restarts the export thread behind it. Nothing in the child records to that pipeline,
so it republishes the totals the parent held at the instant of the fork, once per export interval,
for as long as the child lives. Airflow forks constantly and the long-lived children make it
permanent: LocalExecutor pool workers, the OpenLineage dag-state-change ProcessPoolExecutor, and
the scheduler's log and health-check servers all fork from a scheduler whose pipeline is already
live. The stale pid on the inherited provider is what identifies it. Stopping the reader is not
enough on its own — the revived ticker does one final collect on its way out — so the collect
callback is dropped too.

A MeterProvider the SDK built for the deployment, from OTEL_CONFIG_FILE or an
opentelemetry-instrument agent, reaches the child the same way, carrying the atexit shutdown that
shutdown_on_exit=True registered for it. Its readers are left running, since on that path they
are the only pipeline the child has, but the inherited copy of that hook is dropped so the child
cannot dump the parent's state on the way out. Only a provider's own _metric_readers are ever
stopped, never the class-level _all_metric_readers every provider shares, so an agent's pipeline
keeps working in forked children.

Two smaller corrections fall out of the above. A provider's readers are found via
_metric_readers where the SDK has it and _sdk_config.metric_readers before 1.44, because
Airflow accepts opentelemetry-api>=1.27.0 with no upper bound and reading only the newer
attribute would leave the fix silently inert on a supported install. And flush_otel_metrics() no
longer calls force_flush() on whatever happens to be globally installed, which raised
AttributeError: '_ProxyMeterProvider' object has no attribute 'force_flush' when no provider had
been set.

The fork tests in this file were also not running at all: run_service_helper imported the test
module as tests.observability.metrics.test_otel_logger, and a regular tests package installed
in site-packages shadows the local namespace package regardless of cwd or PYTHONPATH. They now
load by file path, so the fork behaviour above is verified rather than assumed.

related: #71800 — that PR fixes the first defect above by shutting down the replaced provider
instead of never building a second one. The two are alternatives for that half; shutting down and
replacing force-flushes the old pipeline and restarts the series under a new
start_time_unix_nano, which building once avoids. Note that shutting down a provider is no longer
safe once it is deliberately shared between both module copies, so the two approaches now conflict
rather than merely overlapping.

supersedes: #71804, which carried an earlier version of this work.


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (Opus 5), for this description and the verification runs

Generated-by: Claude Code (Opus 5) following the guidelines

A forked child inherited the atexit flush hook, the MeterProvider, and — because
the SDK restarts its reader threads in the child — a live pipeline over the
parent's accumulated state, so it re-exported metrics it never recorded. Building
the pipeline once per process rather than once per call closes the same
duplication on the re-initialization path.
The idle child shows the inherited pipeline goes quiet, but every worker that emits
metrics reaches a different state: it stops the pipeline it inherited and then exports
through one of its own, on the same interval, for the rest of its life. Nothing pinned
that the second half still works -- dropping the reset of the module's provider
reference leaves the child recording into the pipeline it just silenced, and its own
metrics disappear with no test noticing.
``_metric_readers`` is only an attribute of ``MeterProvider`` from opentelemetry-sdk 1.44;
before that the same per-provider list lives on ``_sdk_config.metric_readers``. Airflow
asks for ``opentelemetry-api>=1.27.0`` with no upper bound, so a supported install can
easily be on a version where reaching only for the newer attribute finds nothing and the
child goes on exporting its parent's totals.

That failure is silent apart from one warning per fork, which is what makes it worth
guarding: nothing else about a pipeline nobody stopped looks wrong.

Found by building this onto Astro Runtime 3.3-2, which ships opentelemetry-sdk 1.42.1:
the child kept re-exporting the parent's metrics with the fix in place.
The observability package is symlinked into both airflow-core and task-sdk, so
Python loads it twice and each copy gets its own module globals. A guard held in
one of those globals can only ever deduplicate within its own copy, so a
scheduler still ended up with two pipelines exporting alongside each other --
one per copy -- rather than the one intended.

The same split affects the fork cleanup: with the handle in module state, only
whichever copy happened to build the pipeline could recognise and stop it in the
child. Recording ownership on the provider instead puts it in the one place both
copies already share, the OTel global provider registry, so any copy can tell
whether the pipeline belongs to this process.

Measured over the scheduler's startup order, this takes the process from three
live MeterProviders and three exporter threads to one of each.
The surrounding reasoning is framed entirely around cumulative
temporality, so a reader running delta could conclude the fork
handling does not concern them.
@junghoo-de

Copy link
Copy Markdown

Confirming this on Airflow 3.3.2 in production (CeleryExecutor + KubernetesExecutor, 3 scheduler replicas), exporting OTLP with cumulative temporality to Google Cloud Monitoring (Telemetry API).

  • Each scheduler process ends up with more than one live MeterProvider under the same service.instance.id. In 3.3.2 the two stats.initialize() calls are executors/base_executor.py:213 and jobs/scheduler_job_runner.py:1643; the replaced provider keeps exporting (shutdown_on_exit=False).
  • Observable effect: the scheduler's serde.load_serializers count series receives ~80 points per 20 minutes instead of 40 at a 30s export interval, with alternating ~5s/25s gaps (two writers with a fixed phase offset). The backend rejects part of these points as duplicate time series or for exceeding its maximum sampling period.
  • Series first created after the scheduler loop starts (e.g. dagrun.duration.*) had a single writer and stayed accurate — 18/18 DAGs matched the DagRun Finished scheduler log lines over a 24-minute window.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants