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(

Reply via email to