This is an automated email from the ASF dual-hosted git repository.

eladkal 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 c58a05d6e90 Improve BigQuery-to-SQL column lineage rendering (#70816)
c58a05d6e90 is described below

commit c58a05d6e90328b32f1ff71cbf3b30839bdc9c40
Author: Ulada Zakharava <[email protected]>
AuthorDate: Mon Aug 17 09:53:08 2026 +0200

    Improve BigQuery-to-SQL column lineage rendering (#70816)
---
 .../google/cloud/transfers/bigquery_to_sql.py      | 57 ++++++++++++++++++----
 .../cloud/transfers/test_bigquery_to_mssql.py      |  8 ++-
 .../cloud/transfers/test_bigquery_to_mysql.py      |  8 ++-
 .../cloud/transfers/test_bigquery_to_postgres.py   | 13 ++++-
 4 files changed, 73 insertions(+), 13 deletions(-)

diff --git 
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
 
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
index 5c188fe8d04..5e4cc7407fa 100644
--- 
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
+++ 
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
@@ -151,7 +151,12 @@ class BigQueryToSqlBaseOperator(BaseOperator):
         operators. Children still provide a concrete SQL hook via
         ``get_sql_hook()`` and may override behavior if needed.
         """
-        from airflow.providers.common.compat.openlineage.facet import Dataset
+        from airflow.providers.common.compat.openlineage.facet import (
+            Dataset,
+            DatasetFacet,
+            SchemaDatasetFacet,
+            SchemaDatasetFacetFields,
+        )
         from airflow.providers.google.cloud.openlineage.utils import (
             BIGQUERY_NAMESPACE,
             get_facets_from_bq_table_for_given_fields,
@@ -182,20 +187,40 @@ class BigQueryToSqlBaseOperator(BaseOperator):
 
         if self.selected_fields:
             if isinstance(self.selected_fields, str):
-                bigquery_field_names = list(self.selected_fields)
+                transferred_field_names = [
+                    field.strip() for field in self.selected_fields.split(",") 
if field.strip()
+                ]
             else:
-                bigquery_field_names = self.selected_fields
+                transferred_field_names = list(self.selected_fields)
         else:
-            bigquery_field_names = [f.name for f in getattr(table_obj, 
"schema", [])]
+            transferred_field_names = [field.name for field in 
getattr(table_obj, "schema", [])]
 
         input_dataset = Dataset(
             namespace=BIGQUERY_NAMESPACE,
             name=self.source_project_dataset_table,
-            facets=get_facets_from_bq_table_for_given_fields(table_obj, 
bigquery_field_names),
+            facets=get_facets_from_bq_table_for_given_fields(table_obj, 
selected_fields=None),
+        )
+        input_schema_facet = input_dataset.facets.get("schema") if 
input_dataset.facets else None
+        transferred_field_names_set = set(transferred_field_names)
+        lineage_input_schema_facet = (
+            SchemaDatasetFacet(
+                fields=[
+                    field
+                    for field in getattr(input_schema_facet, "fields", [])
+                    if field.name in transferred_field_names_set
+                ]
+            )
+            if input_schema_facet
+            else None
+        )
+        output_schema_facet = SchemaDatasetFacet(
+            fields=[
+                SchemaDatasetFacetFields(name=field_name, type=None) for 
field_name in transferred_field_names
+            ]
         )
 
         sql_hook = self.get_sql_hook()
-        db_info = sql_hook.get_openlineage_database_info(sql_hook.get_conn())
+        db_info = sql_hook.get_openlineage_database_info(sql_hook.connection)
         if db_info is None:
             self.log.debug("OpenLineage: could not get database info from SQL 
hook %s", type(sql_hook))
             return OperatorLineage()
@@ -228,11 +253,25 @@ class BigQueryToSqlBaseOperator(BaseOperator):
                 else:
                     output_name = f"{table_part}"
 
+        # The identity helper requires the source schema to be a subset of the 
destination schema.
+        lineage_input_dataset = Dataset(
+            namespace=input_dataset.namespace,
+            name=input_dataset.name,
+            facets={"schema": lineage_input_schema_facet} if 
lineage_input_schema_facet else {},
+        )
         column_lineage_facet = get_identity_column_lineage_facet(
-            bigquery_field_names, input_datasets=[input_dataset]
+            transferred_field_names, input_datasets=[lineage_input_dataset]
         )
 
-        output_facets = column_lineage_facet or {}
-        output_dataset = Dataset(namespace=namespace, name=output_name, 
facets=output_facets)
+        output_facets: dict[str, DatasetFacet] = {"schema": 
output_schema_facet}
+
+        if column_lineage_facet:
+            output_facets.update(column_lineage_facet)
+
+        output_dataset = Dataset(
+            namespace=namespace,
+            name=output_name,
+            facets=output_facets,
+        )
 
         return OperatorLineage(inputs=[input_dataset], 
outputs=[output_dataset])
diff --git 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
index 0295f7060f2..9ce37ef445f 100644
--- 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
+++ 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
@@ -199,6 +199,9 @@ class TestBigQueryToMsSqlOperator:
         output_ds = result.outputs[0]
         assert output_ds.namespace == "mssql://localhost:1433"
         assert output_ds.name == "mydb.dbo.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", 
"name", "value"}
+        assert all(field.type is None for field in output_schema_fields)
 
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
@@ -239,11 +242,14 @@ class TestBigQueryToMsSqlOperator:
         assert "schema" in input_ds.facets
 
         schema_fields = [f.name for f in input_ds.facets["schema"].fields]
-        assert set(schema_fields) == {"id", "name"}
+        assert set(schema_fields) == {"id", "name", "value"}
 
         output_ds = result.outputs[0]
         assert output_ds.namespace == "mssql://server.example:1433"
         assert output_ds.name == "mydb.dbo.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", "name"}
+        assert all(field.type is None for field in output_schema_fields)
 
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
diff --git 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mysql.py 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mysql.py
index 1865efe7cfc..5e749f68856 100644
--- 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mysql.py
+++ 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mysql.py
@@ -106,6 +106,9 @@ class TestBigQueryToMySqlOperator:
         output_ds = result.outputs[0]
         assert output_ds.namespace == "mysql://localhost:3306"
         assert output_ds.name == "mydb.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", 
"name", "value"}
+        assert all(field.type is None for field in output_schema_fields)
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
         assert set(col_lineage.fields.keys()) == {"id", "name", "value"}
@@ -144,11 +147,14 @@ class TestBigQueryToMySqlOperator:
         assert input_ds.name == 
f"{TEST_PROJECT}.{TEST_DATASET}.{TEST_TABLE_ID}"
         assert "schema" in input_ds.facets
         schema_fields = [f.name for f in input_ds.facets["schema"].fields]
-        assert set(schema_fields) == {"id", "name"}
+        assert set(schema_fields) == {"id", "name", "value"}
 
         output_ds = result.outputs[0]
         assert output_ds.namespace == "mysql://localhost:3306"
         assert output_ds.name == "mydb.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", "name"}
+        assert all(field.type is None for field in output_schema_fields)
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
         assert set(col_lineage.fields.keys()) == {"id", "name"}
diff --git 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
index b44cefb1f30..0dac90c9025 100644
--- 
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
+++ 
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_postgres.py
@@ -226,6 +226,9 @@ class TestBigQueryToPostgresOperator:
         op.execute(context=context)
 
         result = 
op.get_openlineage_facets_on_complete(task_instance=MagicMock())
+        
mock_postgres_hook.get_openlineage_database_info.assert_called_once_with(
+            mock_postgres_hook.connection
+        )
         assert len(result.inputs) == 1
         assert len(result.outputs) == 1
 
@@ -239,6 +242,9 @@ class TestBigQueryToPostgresOperator:
         output_ds = result.outputs[0]
         assert output_ds.namespace == "postgres://localhost:5432"
         assert output_ds.name == "postgresdb.postgres-schema.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", 
"name", "value"}
+        assert all(field.type is None for field in output_schema_fields)
 
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
@@ -261,7 +267,7 @@ class TestBigQueryToPostgresOperator:
             task_id=TASK_ID,
             dataset_table=f"{TEST_DATASET}.{TEST_TABLE_ID}",
             target_table_name="destination",
-            selected_fields=["id", "name"],
+            selected_fields=" id, name, ",
             database="postgresdb",
         )
         op.bigquery_hook = mock_bq_hook
@@ -279,11 +285,14 @@ class TestBigQueryToPostgresOperator:
         assert input_ds.name == 
f"{TEST_PROJECT}.{TEST_DATASET}.{TEST_TABLE_ID}"
         assert "schema" in input_ds.facets
         schema_fields = [f.name for f in input_ds.facets["schema"].fields]
-        assert set(schema_fields) == {"id", "name"}
+        assert set(schema_fields) == {"id", "name", "value"}
 
         output_ds = result.outputs[0]
         assert output_ds.namespace == "postgres://localhost:5432"
         assert output_ds.name == "postgresdb.postgres-schema.destination"
+        output_schema_fields = output_ds.facets["schema"].fields
+        assert {field.name for field in output_schema_fields} == {"id", "name"}
+        assert all(field.type is None for field in output_schema_fields)
         assert "columnLineage" in output_ds.facets
         col_lineage = output_ds.facets["columnLineage"]
         assert set(col_lineage.fields.keys()) == {"id", "name"}

Reply via email to