This is an automated email from the ASF dual-hosted git repository.
mobuchowski 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 53862c82e99 Fix GCS operators emitting invalid OpenLineage events with
no dataset name (#73246)
53862c82e99 is described below
commit 53862c82e99218e8956cf01e34119219fffe5b51
Author: Kacper Muda <[email protected]>
AuthorDate: Wed Sep 16 19:25:10 2026 +0200
Fix GCS operators emitting invalid OpenLineage events with no dataset name
(#73246)
* Fix openlineage iterating over missing objects
* Also guard destination path when it resolves to None
Addresses review feedback on apache/airflow#73246: a None destination
path crashed the OpenLineage hook before building any dataset, instead
of the empty OperatorLineage the bucket check already returns.
Co-Authored-By: Claude Sonnet 5 <[email protected]>
---------
Co-authored-by: Claude Sonnet 5 <[email protected]>
---
.../providers/google/cloud/operators/gcs.py | 4 +++
.../google/cloud/transfers/local_to_gcs.py | 7 +++-
.../tests/unit/google/cloud/operators/test_gcs.py | 8 +++++
.../google/cloud/transfers/test_local_to_gcs.py | 39 ++++++++++++++++++++++
4 files changed, 57 insertions(+), 1 deletion(-)
diff --git
a/providers/google/src/airflow/providers/google/cloud/operators/gcs.py
b/providers/google/src/airflow/providers/google/cloud/operators/gcs.py
index 9c0d9758a79..2cc26532148 100644
--- a/providers/google/src/airflow/providers/google/cloud/operators/gcs.py
+++ b/providers/google/src/airflow/providers/google/cloud/operators/gcs.py
@@ -355,6 +355,9 @@ class GCSDeleteObjectsOperator(GoogleCloudBaseOperator):
from airflow.providers.google.cloud.openlineage.utils import
extract_ds_name_from_gcs_path
from airflow.providers.openlineage.extractors import OperatorLineage
+ if not self.bucket_name:
+ return OperatorLineage()
+
objects = []
if self.objects is not None:
objects = self.objects
@@ -378,6 +381,7 @@ class GCSDeleteObjectsOperator(GoogleCloudBaseOperator):
},
)
for object_name in objects
+ if object_name
]
return OperatorLineage(inputs=input_datasets)
diff --git
a/providers/google/src/airflow/providers/google/cloud/transfers/local_to_gcs.py
b/providers/google/src/airflow/providers/google/cloud/transfers/local_to_gcs.py
index 17b0d3443ed..9128d46e81d 100644
---
a/providers/google/src/airflow/providers/google/cloud/transfers/local_to_gcs.py
+++
b/providers/google/src/airflow/providers/google/cloud/transfers/local_to_gcs.py
@@ -143,6 +143,9 @@ class LocalFilesystemToGCSOperator(BaseOperator):
from airflow.providers.google.cloud.openlineage.utils import WILDCARD,
extract_ds_name_from_gcs_path
from airflow.providers.openlineage.extractors import OperatorLineage
+ if not self.bucket or self.dst is None:
+ return OperatorLineage()
+
source_facets = {}
if isinstance(self.src, str): # Single path provided, possibly
relative or with wildcard
original_src = f"{self.src}"
@@ -165,6 +168,8 @@ class LocalFilesystemToGCSOperator(BaseOperator):
dest_object = self.dst if os.path.basename(self.dst) else
extract_ds_name_from_gcs_path(self.dst)
return OperatorLineage(
- inputs=[Dataset(namespace="file", name=src, facets=source_facets)
for src in source_objects],
+ inputs=[
+ Dataset(namespace="file", name=src, facets=source_facets) for
src in source_objects if src
+ ],
outputs=[Dataset(namespace=f"gs://{self.bucket}",
name=dest_object)],
)
diff --git a/providers/google/tests/unit/google/cloud/operators/test_gcs.py
b/providers/google/tests/unit/google/cloud/operators/test_gcs.py
index 44ceea9aa8c..ba1c8ef9e79 100644
--- a/providers/google/tests/unit/google/cloud/operators/test_gcs.py
+++ b/providers/google/tests/unit/google/cloud/operators/test_gcs.py
@@ -261,6 +261,7 @@ class TestGCSDeleteObjectsOperator:
(None, "pre", ["/"]),
(None, "dir/pre*", ["dir"]),
(None, "*", ["/"]),
+ (["folder/a.txt", None, "", "b.json"], None, ["folder/a.txt",
"b.json"]),
),
ids=(
"objects",
@@ -273,6 +274,7 @@ class TestGCSDeleteObjectsOperator:
"prefix with no ending slash",
"directory with prefix with wildcard",
"just wildcard",
+ "objects with None and empty entries",
),
)
def test_get_openlineage_facets_on_start(self, objects, prefix, inputs):
@@ -305,6 +307,12 @@ class TestGCSDeleteObjectsOperator:
print("EXPECTED:", expected_inputs)
print("ACTUAL:", lineage.inputs)
+ def test_get_openlineage_facets_on_start_no_bucket_name(self):
+ operator = GCSDeleteObjectsOperator(task_id=TASK_ID, bucket_name=None,
objects=["a.txt"])
+ lineage = operator.get_openlineage_facets_on_start()
+ assert lineage.inputs == []
+ assert lineage.outputs == []
+
class TestGoogleCloudStorageListOperator:
@mock.patch("airflow.providers.google.cloud.operators.gcs.GCSHook")
diff --git
a/providers/google/tests/unit/google/cloud/transfers/test_local_to_gcs.py
b/providers/google/tests/unit/google/cloud/transfers/test_local_to_gcs.py
index 6e32ce62971..b1631d69d9c 100644
--- a/providers/google/tests/unit/google/cloud/transfers/test_local_to_gcs.py
+++ b/providers/google/tests/unit/google/cloud/transfers/test_local_to_gcs.py
@@ -236,6 +236,45 @@ class TestFileToGcsOperator:
assert all(inp.name in expected_inputs for inp in result.inputs)
assert all(inp.namespace == "file" for inp in result.inputs)
+ def
test_get_openlineage_facets_on_start_with_none_and_empty_src_entries(self):
+ operator = LocalFilesystemToGCSOperator(
+ task_id="gcs_to_file_sensor",
+ dag=self.dag,
+ src=[f"{self.tmpdir_posix}/fake1.csv", None, "",
f"{self.tmpdir_posix}/fake2.csv"],
+ dst="test/",
+ **self._config,
+ )
+ result = operator.get_openlineage_facets_on_start()
+ assert len(result.inputs) == 2
+ assert {inp.name for inp in result.inputs} == {
+ f"{self.tmpdir_posix}/fake1.csv",
+ f"{self.tmpdir_posix}/fake2.csv",
+ }
+
+ def test_get_openlineage_facets_on_start_no_bucket(self):
+ operator = LocalFilesystemToGCSOperator(
+ task_id="gcs_to_file_sensor",
+ dag=self.dag,
+ src=[f"{self.tmpdir_posix}/fake1.csv"],
+ dst="test/",
+ **{**self._config, "bucket": None},
+ )
+ result = operator.get_openlineage_facets_on_start()
+ assert result.inputs == []
+ assert result.outputs == []
+
+ def test_get_openlineage_facets_on_start_no_dst(self):
+ operator = LocalFilesystemToGCSOperator(
+ task_id="gcs_to_file_sensor",
+ dag=self.dag,
+ src=[f"{self.tmpdir_posix}/fake1.csv"],
+ dst=None,
+ **self._config,
+ )
+ result = operator.get_openlineage_facets_on_start()
+ assert result.inputs == []
+ assert result.outputs == []
+
# Return value tests
@mock.patch("airflow.providers.google.cloud.transfers.local_to_gcs.GCSHook",
autospec=True)
def test_execute_returns_list_of_destination_uris_single_file(self,
mock_hook):