Skip creating an empty hook lineage collector for OpenLineage events - #74190
Conversation
When an operator has no extractor, or its extractor finds no inputs or outputs, OpenLineage falls back to the hook lineage collector. Hooks report lineage through that process-wide collector, so if the task process never created it, nothing was collected. The fallback still created it, and on Airflow 3 creating the collector initializes the asset URI handlers of every installed provider, which imports heavy client libraries such as the Google Cloud ones. Each task event is emitted from a forked child by default, so a child that created the collector threw it away on exit and the next event paid again. On a released Airflow 3.1 image that is about a second per event on an idle machine, and far more under load: on CI runners in the OpenLineage e2e compat tests a trivial task's forked child took about 23 seconds, more on Python 3.11 than on 3.10, so Dags with long task chains hit their dagrun_timeout. The fallback now returns no hook lineage when the collector was never created. The check looks the collector up the same way the common.compat getter does and reads its singleton state; any state it does not recognize counts as created, so unexpected cases keep the previous behaviour.
Clearing the process-wide cache would let whichever test calls the getter next, possibly with no lineage readers registered, cache a NoOpCollector for every test that runs afterwards. A test-local cached function exercises the same functools.cache behaviour without touching shared state.
7668ecd to
0dbe545
Compare
|
On the open question of why this got worse with the Python 3.11 image: I traced it, and it comes from Measured in the published
Why
A sub-second cost locally becomes the ~23s per event seen in CI once it runs in every forked event child on runners capped at about 1.2 CPU. A newer Beyond fixing CI, this is a nice optimisation in its own right. Most tasks never report hook lineage, yet on Airflow 2.11–3.1 every OpenLineage event (START and COMPLETE, for every task) paid for building an empty collector. That meant importing the asset URI handlers of every installed provider and the heavy client libraries they pull in, only to find nothing and discard it. With this change that work happens only when a hook actually collected lineage. Every deployment running OpenLineage on those versions gets cheaper, faster event emission, on any Python version: per your capped measurement, the fork-to-emit gap dropped from about 22.5s to 0.11s. It also makes OpenLineage less sensitive to how many providers are installed and to changes in third-party import-time behaviour, like the Drafted-by: Claude Code (Opus 5.5); reviewed by @potiuk before posting |
related: #74151
This PR is intended to solve this CI error that I've encountered when upgrading the base image to Python 3.11: https://github.com/apache/airflow/actions/runs/37183552752/job/111391234299?pr=74151
It avoids unnecessarily initializing the hook lineage collector when no hook lineage was collected. Its initialization loads provider asset URI handlers and expensive client libraries in every forked OpenLineage event process, causing significant overhead that became particularly visible with the Python 3.11 CI image and led to Dag timeouts.
I'm still investigating why it got worse in the 3.11 image.
AI Summary
Each OpenLineage task event is emitted from a forked child. If the operator has no extractor, or its extractor finds no inputs or outputs, the child falls back to the hook lineage collector. Creating that collector initializes the asset URI handlers of every installed provider, which imports heavy client libraries such as the Google Cloud ones. On released Airflow 3.1.x the task process never creates the collector unless a hook reports lineage. So every forked child created it from scratch, found it empty, and threw it away on exit.
On an idle machine that costs about a second per event. On loaded CI runners it took about 23 seconds per event with Python 3.11, and less with 3.10, because the image ships no bytecode for installed packages. That is why the OpenLineage e2e compat job for 3.1.8 started failing on #74151, which moves CI to Python 3.11. Trivial tasks took 20–50 seconds, and Dags with long task chains ran past their 5-minute
dagrun_timeout. The harness reruns failed Dags under a new run id, and that rerun then reported doubled events ("Expected 1 events ... got 2").Hooks report lineage through the process-wide collector. So if it was never created in this process, nothing was collected, and
get_hook_lineagecan returnNonewithout creating it. The check looks the collector up the same way the common.compat getter does:airflow.sdk.lineagewith itsfunctools.cachegetter on Airflow 3.2+;airflow.lineage.hookmodule global on 2.11–3.1.Any state it does not recognize counts as created, for example a getter replaced by a test mock. In those cases the previous behaviour is kept.
An earlier approach built the collector once in the task process before forking. It only halved the per-task cost, and the capped e2e below still failed with it, so it was dropped.
Checks run:
extractors/manager.pyandtest_manager.pypasses. OpenLineagetest_manager.py,test_listener.pyandtest_base.py: 149 passed.test_manager.pygives 52 passed, 4 skipped on each. The skips are the tests for other Airflow versions.Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Opus 5.5) following the guidelines
🤖 Generated with Claude Code
https://claude.ai/code/session_014q2ntWTPsnfsF1xpoN3XiR