This is an automated email from the ASF dual-hosted git repository.
tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new c74609d71ee [Python] Support Lakehouse runtime catalog tables in
ReadFromBigQuery… (#40225)
c74609d71ee is described below
commit c74609d71eeae8127cd8fa6b1c06a2c9e92422a5
Author: claudevdm <[email protected]>
AuthorDate: Mon Oct 5 18:38:58 2026 -0400
[Python] Support Lakehouse runtime catalog tables in ReadFromBigQuery…
(#40225)
* [Python] Support Lakehouse runtime catalog tables in ReadFromBigQuery
DIRECT_READ
Port of the Java change in #39597 to the Python SDK.
- parse_table_reference now assigns segments by explicit splitting
(mirroring
BigQueryHelpers.parseTableSpec) so 4-part
`project.catalog.namespace.table`
and `project:catalog.namespace.table` specs map to a composite
`catalog.namespace` dataset id. Previously the dotted form raised and the
colon form silently bound the project to `project:catalog`.
- Storage Read split tolerates tables that report no numBytes by requesting
zero streams so the Read API picks the stream count; the export source's
size estimate returns 0 instead of raising.
- Unit tests for the new spec forms, domain-scoped projects, invalid specs,
and the numBytes-missing paths.
* comments
* comments
---
CHANGES.md | 1 +
sdks/python/apache_beam/io/gcp/bigquery.py | 46 +++++++---
sdks/python/apache_beam/io/gcp/bigquery_test.py | 52 +++++++++++
sdks/python/apache_beam/io/gcp/bigquery_tools.py | 100 +++++++++++++++++----
.../apache_beam/io/gcp/bigquery_tools_test.py | 63 +++++++++++++
5 files changed, 234 insertions(+), 28 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index dbe5a054455..4682f88612d 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -71,6 +71,7 @@
* (Go) Added `wait.On`, which delays each input window until the corresponding
windows in its signal PCollections have closed
([#39909](https://github.com/apache/beam/issues/39909)).
* (Python) Expanded the SDK worker heap dump
(`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC
stats, and glibc `mallinfo2` native-heap/fragmentation stats to help
distinguish native-heap from Python-object memory growth
([#39244](https://github.com/apache/beam/issues/39244)).
* The `disableCounterMetrics`, `disableStringSetMetrics` and
`disableBoundedTrieMetrics` experiments are now honored by the Python SDK, as
they already were in Java (Python)
([#38746](https://github.com/apache/beam/issues/38746)).
+* ReadFromBigQuery now supports Lakehouse runtime catalog (BigLake metastore)
tables with `method=DIRECT_READ`, using 4-part
`project.catalog.namespace.table` identifiers. Previously
`project:catalog.namespace.table` was silently mis-parsed and tables that
report no `numBytes` failed to split (Python)
([#39597](https://github.com/apache/beam/issues/39597)).
## Breaking Changes
diff --git a/sdks/python/apache_beam/io/gcp/bigquery.py
b/sdks/python/apache_beam/io/gcp/bigquery.py
index 38acd29da7d..0caac19fb79 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery.py
@@ -29,7 +29,11 @@ creating the sources or sinks respectively).
Also, for programming convenience, instances of TableReference and TableSchema
have a string representation that can be used for the corresponding arguments:
- - TableReference can be a PROJECT:DATASET.TABLE or DATASET.TABLE string.
+ - TableReference can be a PROJECT:DATASET.TABLE, PROJECT.DATASET.TABLE,
+ DATASET.TABLE or, for Lakehouse runtime catalog tables,
+ PROJECT.CATALOG.NAMESPACE.TABLE, which maps to a composite
+ CATALOG.NAMESPACE dataset id. A Lakehouse reference must include the
+ project id.
- TableSchema can be a NAME:TYPE{,NAME:TYPE}* string
(e.g. 'month:STRING,event_count:INTEGER').
@@ -789,6 +793,10 @@ class _CustomBigQuerySource(BoundedSource):
table_ref.projectId = self._get_project()
table = bq.get_table(
table_ref.projectId, table_ref.datasetId, table_ref.tableId)
+ if table.numBytes is None:
+ # Some tables don't report storage statistics, e.g. Lakehouse runtime
+ # catalog tables.
+ return None
return int(table.numBytes)
elif self.query is not None and self.query.is_accessible():
project = self._get_project()
@@ -1111,8 +1119,24 @@ class _CustomBigQueryStorageSource(BoundedSource):
if table_reference.projectId else self._get_parent_project())
table = bq.get_table(
project, table_reference.datasetId, table_reference.tableId)
+ # None for tables that don't report storage statistics, e.g. Lakehouse
+ # runtime catalog tables.
return table.numBytes
+ def _get_stream_count(self, bq, desired_bundle_size):
+ """Number of streams to request; 0 lets the Storage Read API decide."""
+ table_size = self._get_table_size(bq, self.table_reference)
+ if table_size is None:
+ # A table that reports no size, e.g. a Lakehouse runtime catalog table,
+ # cannot be split by size.
+ return 0
+ stream_count = 0
+ if desired_bundle_size > 0:
+ stream_count = min(
+ int(table_size / desired_bundle_size),
+ _CustomBigQueryStorageSource.MAX_SPLIT_COUNT)
+ return max(stream_count, _CustomBigQueryStorageSource.MIN_SPLIT_COUNT)
+
def _get_bq_metadata(self):
if not self.bq_io_metadata:
self.bq_io_metadata = create_bigquery_io_metadata(self._step_name)
@@ -1256,14 +1280,7 @@ class _CustomBigQueryStorageSource(BoundedSource):
requested_session.read_options.row_restriction = self.row_restriction
storage_client = bq_storage.BigQueryReadClient()
- stream_count = 0
- if desired_bundle_size > 0:
- table_size = self._get_table_size(bq, self.table_reference)
- stream_count = min(
- int(table_size / desired_bundle_size),
- _CustomBigQueryStorageSource.MAX_SPLIT_COUNT)
- stream_count = max(
- stream_count, _CustomBigQueryStorageSource.MIN_SPLIT_COUNT)
+ stream_count = self._get_stream_count(bq, desired_bundle_size)
parent = 'projects/{}'.format(self.table_reference.projectId)
read_session = storage_client.create_read_session(
@@ -2953,10 +2970,13 @@ class ReadFromBigQuery(PTransform):
table (str, callable, ValueProvider): The ID of the table, or a callable
that returns it. If dataset argument is :data:`None` then the table
argument must contain the entire table reference specified as:
- ``'DATASET.TABLE'`` or ``'PROJECT:DATASET.TABLE'``. If it's a callable,
- it must receive one argument representing an element to be written to
- BigQuery, and return a TableReference, or a string table name as
specified
- above.
+ ``'DATASET.TABLE'``, ``'PROJECT:DATASET.TABLE'``,
+ ``'PROJECT.DATASET.TABLE'`` or, for Lakehouse runtime catalog tables,
+ ``'PROJECT.CATALOG.NAMESPACE.TABLE'``. A Lakehouse reference must
+ include the project id.
+ If it's a callable, it must receive one argument representing an element
+ to be written to BigQuery, and return a TableReference, or a string table
+ name as specified above.
dataset (str): The ID of the dataset containing this table or
:data:`None` if the table reference is specified entirely by the table
argument.
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_test.py
index 874ae68362e..d7d4cab09d9 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_test.py
@@ -754,6 +754,58 @@ class TestReadFromBigQuery(unittest.TestCase):
Lineage.query(p.result.metrics(), Lineage.SOURCE),
set(["bigquery:project.dataset.table"]))
+ @parameterized.expand([
+ # Tables without storage statistics (e.g. Lakehouse runtime catalog
+ # tables) let the Storage Read API pick the stream count.
+ param(num_bytes=None, expected_max_stream_count=0),
+ param(
+ num_bytes=5,
+ expected_max_stream_count=beam_bq._CustomBigQueryStorageSource.
+ MIN_SPLIT_COUNT),
+ ])
+ def test_direct_read_split_stream_count(
+ self, num_bytes, expected_max_stream_count):
+ class DummyTable:
+ numBytes = num_bytes
+
+ with mock.patch.object(BigQueryWrapper, '_bigquery_client'), \
+ mock.patch.object(BigQueryWrapper, 'get_table',
+ return_value=DummyTable()), \
+ mock.patch.object(bq_storage.BigQueryReadClient,
+ 'create_read_session') as mock_create_session:
+ mock_create_session.return_value = mock.Mock(streams=[])
+ source = beam_bq._CustomBigQueryStorageSource(
+ method=ReadFromBigQuery.Method.DIRECT_READ,
+ table='project.catalog.namespace.table',
+ pipeline_options=PipelineOptions(['--project=project']))
+ self.assertEqual(source.estimate_size(), num_bytes)
+ self.assertEqual(list(source.split(desired_bundle_size=1 << 20)), [])
+
+ _, kwargs = mock_create_session.call_args
+ self.assertEqual(kwargs['max_stream_count'], expected_max_stream_count)
+ self.assertEqual(
+ kwargs['read_session'].table,
+ 'projects/project/datasets/catalog.namespace/tables/table')
+
+ @parameterized.expand([
+ # A table without storage statistics has an unknown size, not a size
+ # of zero.
+ param(num_bytes=None, expected_size=None),
+ param(num_bytes=5, expected_size=5),
+ ])
+ def test_export_estimate_size(self, num_bytes, expected_size):
+ class DummyTable:
+ numBytes = num_bytes
+
+ with mock.patch.object(BigQueryWrapper, '_bigquery_client'), \
+ mock.patch.object(BigQueryWrapper, 'get_table',
+ return_value=DummyTable()):
+ source = beam_bq._CustomBigQuerySource(
+ method=ReadFromBigQuery.Method.EXPORT,
+ table='project.catalog.namespace.table',
+ pipeline_options=PipelineOptions(['--project=project']))
+ self.assertEqual(source.estimate_size(), expected_size)
+
def test_read_all_lineage(self):
# TODO(https://github.com/apache/beam/issues/34549): This test relies on
# lineage metrics which Prism doesn't seem to handle correctly. Defaulting
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
index abe204c237a..cc1a16a659f 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_tools.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_tools.py
@@ -106,8 +106,12 @@ UNKNOWN_MIME_TYPE = 'application/octet-stream'
BQ_STREAMING_INSERT_TIMEOUT_SEC = 120
_PROJECT_PATTERN = r'([a-z0-9.-]+:)?[a-z][a-z0-9-]*[a-z0-9]'
-_DATASET_PATTERN = r'\w{1,1024}'
+# Dots are allowed because Lakehouse runtime catalog tables are addressed with
+# a composite 'catalog.namespace' dataset id.
+_DATASET_PATTERN = r'[-\w.]{1,1024}'
_TABLE_PATTERN = r'[\p{L}\p{M}\p{N}\p{Pc}\p{Pd}\p{Zs}$]{1,1024}'
+# A single segment of a project id (no dots or colons).
+_PROJECT_NAME_SEGMENT_PATTERN = r'[-a-z0-9]*[a-z0-9]'
# For CI: temp dataset of this name are automatically deleted
_TEMP_DATASET_PREFIX = 'beam_temp_dataset_'
@@ -252,9 +256,15 @@ def parse_table_reference(table, dataset=None,
project=None):
table: The ID of the table. The ID must contain only letters
(a-z, A-Z), numbers (0-9), connectors (-_). If dataset argument is None
then the table argument must contain the entire table reference:
- 'DATASET.TABLE' or 'PROJECT:DATASET.TABLE'. This argument can be a
- TableReference instance in which case dataset and project are
- ignored and the reference is returned as a result. Additionally, for
date
+ 'DATASET.TABLE', 'PROJECT:DATASET.TABLE' or 'PROJECT.DATASET.TABLE'.
+ Lakehouse runtime catalog tables use four parts,
+ 'PROJECT.CATALOG.NAMESPACE.TABLE' (the canonical spelling) or
+ 'PROJECT:CATALOG.NAMESPACE.TABLE', which parse to a composite
+ 'CATALOG.NAMESPACE' dataset id. The project id is required for these,
+ because a composite dataset id without one would be ambiguous with
+ 'PROJECT.DATASET.TABLE'. This argument can be a TableReference instance
+ in which case dataset and project are ignored and the reference is
+ returned as a result. Additionally, for date
partitioned tables, appending '$YYYYmmdd' to the table name is supported,
e.g. 'DATASET.TABLE$YYYYmmdd'.
dataset: The ID of the dataset containing this table or null if the table
@@ -288,17 +298,10 @@ def parse_table_reference(table, dataset=None,
project=None):
# table argument will contain a full table reference instead of just a
# table name.
if dataset is None:
- pattern = (
- f'((?P<project>{_PROJECT_PATTERN})[:\\.])?'
- f'(?P<dataset>{_DATASET_PATTERN})\\.(?P<table>{_TABLE_PATTERN})')
- match = regex.fullmatch(pattern, table)
- if not match:
- raise ValueError(
- 'Expected a table reference (PROJECT:DATASET.TABLE or '
- 'DATASET.TABLE) instead of %s.' % table)
- table_reference.projectId = match.group('project')
- table_reference.datasetId = match.group('dataset')
- table_reference.tableId = match.group('table')
+ project_id, dataset_id, table_id = _split_table_spec(table)
+ table_reference.projectId = project_id
+ table_reference.datasetId = dataset_id
+ table_reference.tableId = table_id
else:
table_reference.projectId = project
table_reference.datasetId = dataset
@@ -306,6 +309,73 @@ def parse_table_reference(table, dataset=None,
project=None):
return table_reference
+def _invalid_table_spec(table_spec):
+ return ValueError(
+ 'Expected a table reference (PROJECT:DATASET.TABLE, '
+ 'PROJECT.DATASET.TABLE, DATASET.TABLE, '
+ 'PROJECT:CATALOG.NAMESPACE.TABLE or PROJECT.CATALOG.NAMESPACE.TABLE) '
+ 'instead of %s.' % table_spec)
+
+
+def _split_table_spec(table_spec):
+ """Splits a table spec string into (project, dataset, table).
+
+ The regex only validates the character set; segment assignment is done by
+ explicit splitting so that composite Lakehouse dataset ids ('catalog.ns')
+ and domain-scoped project ids ('example.com:proj') are both handled. This
+ mirrors BigQueryHelpers.parseTableSpec in the Java SDK.
+ """
+ pattern = (
+ f'((?P<project>{_PROJECT_PATTERN})[:\\.])?'
+ f'(?P<dataset>{_DATASET_PATTERN})\\.(?P<table>{_TABLE_PATTERN})')
+ if not regex.fullmatch(pattern, table_spec):
+ raise _invalid_table_spec(table_spec)
+
+ # Table ids cannot contain '.', so the table is always the last segment.
+ last_dot = table_spec.rfind('.')
+ table = table_spec[last_dot + 1:]
+ prefix = table_spec[:last_dot]
+
+ colon_count = prefix.count(':')
+ if colon_count == 0:
+ # Purely dotted form: 'd.t', 'p.d.t', 'p.catalog.ns.t'. A dataset id may
+ # itself contain dots (a composite 'catalog.namespace'), so omitting the
+ # project id would be ambiguous with 'PROJECT.DATASET.TABLE'. Only the
+ # plain two-segment 'DATASET.TABLE' may leave the project id out.
+ first_dot = prefix.find('.')
+ if first_dot < 0:
+ return None, prefix, table
+ project = prefix[:first_dot]
+ dataset = prefix[first_dot + 1:]
+ if not regex.fullmatch(_PROJECT_PATTERN, project) or not dataset:
+ raise _invalid_table_spec(table_spec)
+ return project, dataset, table
+
+ if colon_count == 1:
+ # 'p:d.t', 'p:catalog.ns.t', 'example.com:proj.ds.t'. If the part before
+ # the colon is dotted it is the domain half of a legacy domain-scoped
+ # project id; the first segment after the colon completes the project
+ # and any remaining segments are the (possibly composite) dataset.
+ colon = prefix.find(':')
+ project = prefix[:colon]
+ dataset = prefix[colon + 1:]
+ first_dot = dataset.find('.')
+ if (first_dot >= 0 and first_dot < len(dataset) - 1 and '.' in project and
+ regex.fullmatch(_PROJECT_NAME_SEGMENT_PATTERN, dataset[:first_dot])):
+ project = project + ':' + dataset[:first_dot]
+ dataset = dataset[first_dot + 1:]
+ return project, dataset, table
+
+ if colon_count == 2:
+ # The last colon terminates a domain-scoped project id spelled with its
+ # own colon: 'example.com:proj:ds.t' or 'example.com:proj:catalog.ns.t'.
+ last_colon = prefix.rfind(':')
+ return prefix[:last_colon], prefix[last_colon + 1:], table
+
+ # Project ids contain at most one colon, so more than two is never valid.
+ raise _invalid_table_spec(table_spec)
+
+
# -----------------------------------------------------------------------------
# BigQueryWrapper.
diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
b/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
index c2a5c73f675..9f8b78ee5c3 100644
--- a/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
+++ b/sdks/python/apache_beam/io/gcp/bigquery_tools_test.py
@@ -136,6 +136,46 @@ class TestTableReferenceParser(unittest.TestCase):
('project:dataset.test- table', 'project', 'dataset', 'test- table'),
('project.dataset. test_table', 'project', 'dataset', ' test_table'),
('project.dataset.test$table', 'project', 'dataset', 'test$table'),
+ ('project.dataset.table', 'project', 'dataset', 'table'),
+ ('my-project.my_dataset.table', 'my-project', 'my_dataset', 'table'),
+ # Lakehouse runtime catalog tables: composite 'catalog.namespace'
+ # dataset id.
+ (
+ 'project.catalog.namespace.table',
+ 'project',
+ 'catalog.namespace',
+ 'table'),
+ (
+ 'project:catalog.namespace.table',
+ 'project',
+ 'catalog.namespace',
+ 'table'),
+ (
+ 'my-project.my-catalog.ns.events$20240101',
+ 'my-project',
+ 'my-catalog.ns',
+ 'events$20240101'),
+ # Domain-scoped project ids, in both spellings.
+ (
+ 'example.com:proj.dataset.table',
+ 'example.com:proj',
+ 'dataset',
+ 'table'),
+ (
+ 'example.com:proj:dataset.table',
+ 'example.com:proj',
+ 'dataset',
+ 'table'),
+ (
+ 'example.com:proj.catalog.namespace.table',
+ 'example.com:proj',
+ 'catalog.namespace',
+ 'table'),
+ (
+ 'example.com:proj:catalog.namespace.table',
+ 'example.com:proj',
+ 'catalog.namespace',
+ 'table'),
])
def test_calling_with_fully_qualified_table_ref(
self,
@@ -159,10 +199,33 @@ class TestTableReferenceParser(unittest.TestCase):
self.assertEqual(parsed_ref.datasetId, datasetId)
self.assertEqual(parsed_ref.tableId, tableId)
+ def test_composite_dataset_requires_project(self):
+ # A composite 'catalog.namespace' dataset id is only recognised with an
+ # explicit project id. Without one it would be ambiguous with
+ # 'PROJECT.DATASET.TABLE', so it is rejected rather than guessed at.
+ self.assertRaises(
+ ValueError, parse_table_reference, 'my_catalog.namespace.test_table')
+
def test_calling_with_insufficient_table_ref(self):
table = 'test_table'
self.assertRaises(ValueError, parse_table_reference, table)
+ @parameterized.expand([
+ ('a:b:c:d.table', ),
+ ('project:.table', ),
+ ('project:dataset', ),
+ ('MyProject.MyDataset.table', ),
+ ('1project.dataset.table', ),
+ ('c.namespace.table', ),
+ ('.dataset.table', ),
+ ('project..table', ),
+ ('a.b.c.d.e.f.g', ),
+ ])
+ def test_calling_with_invalid_table_ref(self, table):
+ # The dotted specs here have a leading segment that cannot be a project
+ # id; they are rejected rather than bound as a composite dataset id.
+ self.assertRaises(ValueError, parse_table_reference, table)
+
def test_calling_with_all_arguments(self):
projectId = 'test_project'
datasetId = 'test_dataset'