Skip to content

Skip creating an empty hook lineage collector for OpenLineage events - #74190

Merged
potiuk merged 2 commits into
apache:mainfrom
shahar1:openlineage-skip-unbuilt-hook-lineage-collector
Oct 4, 2026
Merged

potiuk merged 2 commits into
apache:mainfrom
shahar1:openlineage-skip-unbuilt-hook-lineage-collector

Conversation

@shahar1

@shahar1 shahar1 commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

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_lineage can return None without creating it. The check looks the collector up the same way the common.compat getter does:

  • airflow.sdk.lineage with its functools.cache getter on Airflow 3.2+;
  • the airflow.lineage.hook module 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:

  • OpenLineage e2e compat (Airflow 3.1.8, Python 3.11), local docker-compose: Airflow services were capped at 1.2 CPU each to match CI. Without this change the child gap averaged 22.5s, the same as CI. Results:
    • without the change: run failed; mean task span 65s;
    • with the change: passed in 4m04s with no Dag reruns; per-event fork-to-emit gap 0.11s mean; mean task span 7.9s.
  • Breeze, Python 3.11: mypy on extractors/manager.py and test_manager.py passes. OpenLineage test_manager.py, test_listener.py and test_base.py: 149 passed.
  • Provider wheels from this branch on Airflow 3.1.8 and 2.11.0 (provider-compat setup): test_manager.py gives 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?
  • Yes — Claude Code (Opus 5.5)

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

🤖 Generated with Claude Code

https://claude.ai/code/session_014q2ntWTPsnfsF1xpoN3XiR

@shahar1
shahar1 marked this pull request as ready for review October 4, 2026 14:47
@shahar1
shahar1 requested a review from mobuchowski as a code owner October 4, 2026 14:47
@shahar1
shahar1 requested review from ashb and potiuk October 4, 2026 15:20

@SameerMesiah97 SameerMesiah97 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Just one nit.

Comment thread providers/openlineage/tests/unit/openlineage/extractors/test_manager.py Outdated
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.
@shahar1
shahar1 force-pushed the openlineage-skip-unbuilt-hook-lineage-collector branch from 7668ecd to 0dbe545 Compare October 4, 2026 16:11
@shahar1
shahar1 requested a review from SameerMesiah97 October 4, 2026 16:12
@shahar1 shahar1 added the ready for maintainer review Set after triaging when all criteria pass. label Oct 4, 2026
@potiuk

potiuk commented Oct 4, 2026

Copy link
Copy Markdown
Member

On the open question of why this got worse with the Python 3.11 image: I traced it, and it comes from google-api-core hitting a change in importlib.metadata between 3.10 and 3.11. It isn't the interpreter or missing bytecode.

Measured in the published apache/airflow:3.1.8-python3.10 and -python3.11 images, which the compat job builds on:

  • Same work on both sides: neither image ships .pyc files, both interpreters are built the same way (PGO + LTO, -O3) and run equally fast, and collector creation loads the same ~1,070 modules with the same Google, protobuf (upb) and grpc versions.
  • Creating the hook lineage collector on a warm run takes 0.74s on 3.11 vs 0.41s on 3.10, and 2.28s vs 1.48s cold. Almost all of the extra time is in the __init__ of google.api_core (51 → 208 ms) and of google.cloud.secretmanager_v1 (51 → 204 ms).
  • Both run check_python_version() when imported. In google-api-core 2.30.0, the version in the 3.1.8 image, that function always calls importlib.metadata.packages_distributions(), uncached, just to build a package name for a possible warning. The two calls made during collector creation take 0.41s on 3.11 vs 0.10s on 3.10. That covers nearly all of the gap.

Why packages_distributions() is slower on 3.11:

  • On 3.10 it reads only top_level.txt for each distribution and skips any that don't have one.
  • On 3.11, if top_level.txt is missing, it falls back to _top_level_inferred(). That loads dist.files, which means parsing the whole RECORD file.
  • top_level.txt is written only by setuptools-built wheels. In the 3.11 image, 160 of 431 distributions don't have one (flit/hatch/maturin wheels, apache-airflow-core itself, litellm, pandas, scipy, openai and others), so 3.11 parses about 27,000 extra RECORD rows on every call. It finds 334 import names instead of 228.
  • This came from importlib_metadata v4.7.0 (Right way to run AdHoc Pipeline #330: "In packages_distributions, now infer top-level names from .files() when a top-level.txt (Setuptools-specific metadata) is not present"). It reached CPython in 3.11.0a4, via bpo-44893 / bpo-44893: Implement EntryPoint as simple class with attributes. python/cpython#30150, which "Syncs with importlib_metadata 4.8.1". Python 3.10's stdlib is still at importlib_metadata 4.6.

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 google-api-core helps only partly. main pins 2.38.0, which caches packages_distributions(), so calls drop from 2 to 1. But each forked child is a fresh process: in my test the remaining call still took 0.20s vs 0.05s on 3.10, and the compat image stays on 2.30.0 anyway.

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 google-api-core / importlib.metadata combination above.


Drafted-by: Claude Code (Opus 5.5); reviewed by @potiuk before posting

@potiuk
potiuk merged commit 8055f6a into apache:main Oct 4, 2026
103 checks passed
@shahar1
shahar1 deleted the openlineage-skip-unbuilt-hook-lineage-collector branch October 5, 2026 04:03
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:providers provider:openlineage AIP-53 ready for maintainer review Set after triaging when all criteria pass.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants