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"}