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

Amar3tto pushed a commit to branch test-python-arm-sep-8
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/test-python-arm-sep-8 by this 
push:
     new df8d17ff7d7 Fix BigTable DirectRow
df8d17ff7d7 is described below

commit df8d17ff7d71d52993966bf1d1c88d0f7374b68d
Author: Vitaly Terentyev <[email protected]>
AuthorDate: Tue Sep 8 12:22:40 2026 +0400

    Fix BigTable DirectRow
---
 sdks/python/apache_beam/io/gcp/bigtableio.py | 26 ++++++++++++++++++++++++++
 1 file changed, 26 insertions(+)

diff --git a/sdks/python/apache_beam/io/gcp/bigtableio.py 
b/sdks/python/apache_beam/io/gcp/bigtableio.py
index cd78deb7466..9295b9e0376 100644
--- a/sdks/python/apache_beam/io/gcp/bigtableio.py
+++ b/sdks/python/apache_beam/io/gcp/bigtableio.py
@@ -62,15 +62,40 @@ try:
   from google.cloud.bigtable import Client
   from google.cloud.bigtable.batcher import MutationsBatcher
   from google.cloud.bigtable.row import Cell
+  from google.cloud.bigtable.row import DirectRow
   from google.cloud.bigtable.row import PartialRowData
 
 except ImportError:
+  DirectRow = None
   _LOGGER.warning(
       'ImportError: from google.cloud.bigtable import Client', exc_info=True)
 
 __all__ = ['WriteToBigTable', 'ReadFromBigtable']
 
 
+def _restore_direct_row_pb_mutations(row):
+  # google-cloud-bigtable >= 2.44.0 stores mutations on `_mutations`.
+  # Older MutationsBatcher reads `_pb_mutations`. Fill that attribute so
+  # pickle and worker batching both see protobuf mutations.
+  if hasattr(row, '_pb_mutations'):
+    return
+  mutations = getattr(row, '_mutations', None)
+  if mutations is None:
+    return
+  row._pb_mutations = [
+    mut._to_pb() if hasattr(mut, '_to_pb') else mut for mut in mutations
+  ]
+
+
+def _direct_row_getstate(self):
+  _restore_direct_row_pb_mutations(self)
+  return self.__dict__
+
+
+if DirectRow is not None:
+  DirectRow.__getstate__ = _direct_row_getstate
+
+
 class _BigTableWriteFn(beam.DoFn):
   """ Creates the connector can call and add_row to the batcher using each
   row in beam pipe line
@@ -168,6 +193,7 @@ class _BigTableWriteFn(beam.DoFn):
     #                     'field1',
     #                     'value1',
     #                     timestamp=datetime.now())
+    _restore_direct_row_pb_mutations(row)
     self.batcher.mutate(row)
 
   def finish_bundle(self):

Reply via email to