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):