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'

Reply via email to