This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 8055f6a8717 Skip creating an empty hook lineage collector for
OpenLineage events (#74190)
8055f6a8717 is described below
commit 8055f6a871763c8efd1c17843235cd7192991d09
Author: Shahar Epstein <[email protected]>
AuthorDate: Mon Oct 5 02:20:33 2026 +0300
Skip creating an empty hook lineage collector for OpenLineage events
(#74190)
* Skip creating an empty hook lineage collector for OpenLineage events
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.
* Avoid clearing the real hook lineage collector cache in tests
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.
---
.../providers/openlineage/extractors/manager.py | 31 ++++++++++
.../unit/openlineage/extractors/test_manager.py | 71 +++++++++++++++++++++-
2 files changed, 100 insertions(+), 2 deletions(-)
diff --git
a/providers/openlineage/src/airflow/providers/openlineage/extractors/manager.py
b/providers/openlineage/src/airflow/providers/openlineage/extractors/manager.py
index 75090449580..c6b936970a1 100644
---
a/providers/openlineage/src/airflow/providers/openlineage/extractors/manager.py
+++
b/providers/openlineage/src/airflow/providers/openlineage/extractors/manager.py
@@ -16,6 +16,8 @@
# under the License.
from __future__ import annotations
+import importlib
+import types
from collections.abc import Iterator
from typing import TYPE_CHECKING
@@ -45,6 +47,33 @@ if TYPE_CHECKING:
from airflow.providers.common.compat.sdk import BaseOperator
+def _is_hook_lineage_collector_created() -> bool:
+ """
+ Return False only if the hook lineage collector was certainly never
created in this process.
+
+ Hooks report lineage through this process-wide collector, so if it was
never created, nothing was
+ collected. Creating it just to read it back empty imports every provider's
asset URI handlers,
+ which takes seconds. The check reads Airflow internals, so any state it
does not recognize, such
+ as a getter that is not a plain function, counts as created.
+ """
+ try:
+ from airflow.sdk import lineage
+ except ImportError:
+ # Airflow < 3.2 keeps the collector in a module global instead of a
cached getter.
+ try:
+ hook = importlib.import_module("airflow.lineage.hook")
+ except ImportError:
+ return True
+ if getattr(hook, "_hook_lineage_collector", True) is not None:
+ return True
+ return type(getattr(hook, "get_hook_lineage_collector", None)) is not
types.FunctionType
+
+ try:
+ return lineage.get_hook_lineage_collector.cache_info().currsize != 0
+ except AttributeError:
+ return True
+
+
def _iter_extractor_types() -> Iterator[type[BaseExtractor]]:
if PythonExtractor is not None:
yield PythonExtractor
@@ -267,6 +296,8 @@ class ExtractorManager(LoggingMixin):
except ImportError:
return None
+ if not _is_hook_lineage_collector_created():
+ return None
collector = get_hook_lineage_collector()
if not hasattr(collector, "has_collected"):
return None
diff --git
a/providers/openlineage/tests/unit/openlineage/extractors/test_manager.py
b/providers/openlineage/tests/unit/openlineage/extractors/test_manager.py
index d026c1b5039..d153b5b7895 100644
--- a/providers/openlineage/tests/unit/openlineage/extractors/test_manager.py
+++ b/providers/openlineage/tests/unit/openlineage/extractors/test_manager.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import functools
import tempfile
from typing import TYPE_CHECKING, Any
from unittest import mock
@@ -35,13 +36,16 @@ from airflow.providers.common.compat.lineage.entities
import Column, File, Table
from airflow.providers.common.compat.sdk import BaseOperator, Context,
ObjectStoragePath
from airflow.providers.common.sql.hooks.lineage import SqlJobHookLineageExtra
from airflow.providers.openlineage.extractors import OperatorLineage
-from airflow.providers.openlineage.extractors.manager import ExtractorManager
+from airflow.providers.openlineage.extractors.manager import (
+ ExtractorManager,
+ _is_hook_lineage_collector_created,
+)
from airflow.providers.openlineage.utils.utils import Asset
from airflow.utils.state import State, TaskInstanceState
from tests_common.test_utils.compat import DateTimeSensor, PythonOperator
from tests_common.test_utils.markers import
skip_if_force_lowest_dependencies_marker
-from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS
+from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS,
AIRFLOW_V_3_2_PLUS
if TYPE_CHECKING:
try:
@@ -617,6 +621,69 @@ def
test_get_hook_lineage_returns_none_when_nothing_collected(hook_lineage_colle
mock_sql_fn.assert_not_called()
+@patch("airflow.providers.common.compat.lineage.hook.get_hook_lineage_collector",
autospec=True)
+@patch(
+
"airflow.providers.openlineage.extractors.manager._is_hook_lineage_collector_created",
+ autospec=True,
+ return_value=False,
+)
+def test_get_hook_lineage_does_not_create_collector(mock_created,
mock_get_collector):
+ """Creating the collector only to read it back empty imports every
provider's asset URI handlers."""
+ result = ExtractorManager().get_hook_lineage(
+ task_instance=MagicMock(spec=TaskInstance),
task_instance_state=TaskInstanceState.SUCCESS
+ )
+
+ assert result is None
+ mock_get_collector.assert_not_called()
+
+
[email protected](not AIRFLOW_V_3_2_PLUS, reason="Airflow 3.2+ caches the
collector getter")
+class TestIsHookLineageCollectorCreatedAirflow32:
+ def test_real_getter_is_cached(self):
+ from airflow.sdk import lineage
+
+ assert hasattr(lineage.get_hook_lineage_collector, "cache_info")
+
+ def test_tracks_whether_cached_getter_was_called(self):
+ # A test-local cache instead of clearing the real one: clearing it
would let the next caller,
+ # possibly under ``mock_plugin_manager`` with no readers, cache a
NoOpCollector for later tests.
+ getter = functools.cache(lambda: object())
+ with patch("airflow.sdk.lineage.get_hook_lineage_collector",
new=getter):
+ assert _is_hook_lineage_collector_created() is False
+ getter()
+ assert _is_hook_lineage_collector_created() is True
+
+ def test_true_when_getter_is_not_cached(self):
+ with patch("airflow.sdk.lineage.get_hook_lineage_collector",
new=lambda: None):
+ assert _is_hook_lineage_collector_created() is True
+
+
[email protected](AIRFLOW_V_3_2_PLUS, reason="Airflow before 3.2 keeps the
collector in a module global")
+class TestIsHookLineageCollectorCreatedBeforeAirflow32:
+ def test_false_when_collector_never_created(self):
+ with patch("airflow.lineage.hook._hook_lineage_collector", None):
+ assert _is_hook_lineage_collector_created() is False
+
+ def test_true_when_collector_created(self):
+ with patch("airflow.lineage.hook._hook_lineage_collector", object()):
+ assert _is_hook_lineage_collector_created() is True
+
+ def test_true_when_getter_is_replaced(self):
+ from airflow.lineage import hook
+
+ # On Airflow 3.0 and 3.1 the shared ``hook_lineage_collector`` fixture
replaces the getter with a
+ # mock. A spec'd mock also passes ``isinstance`` checks against a
function, so it is the harder
+ # case to treat as created.
+ with (
+ patch("airflow.lineage.hook._hook_lineage_collector", None),
+ patch(
+ "airflow.lineage.hook.get_hook_lineage_collector",
+ new=MagicMock(spec=hook.get_hook_lineage_collector),
+ ),
+ ):
+ assert _is_hook_lineage_collector_created() is True
+
+
def test_get_hook_lineage_passes_failed_state(hook_lineage_collector):
hook = MagicMock()
hook_lineage_collector.add_extra(