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

claudevdm 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 7d0f1bf3870 Fix PickleCoder.as_deterministic_coder() raising TypeError 
(#39943)
7d0f1bf3870 is described below

commit 7d0f1bf387016a6e6360efe31adeb08cefe32afa
Author: Amir Fathi <[email protected]>
AuthorDate: Tue Sep 22 07:05:53 2026 -0600

    Fix PickleCoder.as_deterministic_coder() raising TypeError (#39943)
    
    * Fix PickleCoder.as_deterministic_coder() TypeError
    
    PickleCoder and _MemoizingPickleCoder construct FastPrimitivesCoder with a
    requires_deterministic kwarg the constructor has never accepted, raising
    TypeError on every call instead of returning a deterministic coder. Route
    both through the existing 
_update_compatible_deterministic_fast_primitives_coder
    helper, matching FastPrimitivesCoder's own as_deterministic_coder.
    
    Fixes #39942
    
    * Warn before falling back off PickleCoder to a deterministic coder
    
    Addresses review feedback on #39943: the deterministic fallback still
    fails at encode time for arbitrary classes, so log what changed and
    what it still cannot encode before the fallback is used.
---
 CHANGES.md                                    |  1 +
 sdks/python/apache_beam/coders/coders.py      | 33 +++++++++++++++++++++++++--
 sdks/python/apache_beam/coders/coders_test.py | 15 ++++++++++++
 3 files changed, 47 insertions(+), 2 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 4866998bb05..87b8ad672eb 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -132,6 +132,7 @@
 * (Java) MongoDbIO read splitting now preserves non-ObjectId `_id` types (e.g. 
string ids) instead of failing to parse the generated range filters 
([#39900](https://github.com/apache/beam/issues/39900)).
 * (Go) Fixed GCS glob matching silently dropping objects when the glob pattern 
contains multi-byte characters 
([#39969](https://github.com/apache/beam/issues/39969)).
 * (Python) Fixed `TensorRTEngineHandlerNumPy` failing with 
`CUDA_ERROR_INVALID_VALUE` on models with a single-element input or output 
tensor ([#36306](https://github.com/apache/beam/issues/36306)).
+* (Python) Fixed 
`PickleCoder`/`_MemoizingPickleCoder.as_deterministic_coder()` raising 
`TypeError` instead of returning a working deterministic coder 
([#28558](https://github.com/apache/beam/issues/28558)).
 
 # [2.76.0] - 2026-08-31
 
diff --git a/sdks/python/apache_beam/coders/coders.py 
b/sdks/python/apache_beam/coders/coders.py
index 5f22bff5351..72d67cc9b55 100644
--- a/sdks/python/apache_beam/coders/coders.py
+++ b/sdks/python/apache_beam/coders/coders.py
@@ -38,6 +38,7 @@ encoded.
 
 import base64
 import decimal
+import logging
 import pickle
 from functools import lru_cache
 from typing import TYPE_CHECKING
@@ -114,6 +115,8 @@ __all__ = [
     'PaneInfoCoder'
 ]
 
+_LOGGER = logging.getLogger(__name__)
+
 T = TypeVar('T')
 CoderT = TypeVar('CoderT', bound='Coder')
 ProtoCoderT = TypeVar('ProtoCoderT', bound='ProtoCoder')
@@ -899,7 +902,20 @@ class _MemoizingPickleCoder(_PickleCoderBase):
     return coder_impl.CallbackCoderImpl(_nonhashable_dumps, pickler.loads)
 
   def as_deterministic_coder(self, step_label, error_message=None):
-    return FastPrimitivesCoder(self, requires_deterministic=step_label)
+    _LOGGER.warning(
+        "PickleCoder was registered for a key type in '%s', but the runner "
+        "requires a deterministic key encoding and pickle is not "
+        "deterministic. Keys in this step will be encoded with the "
+        "deterministic fallback coder instead of pickle. That coder "
+        "supports primitives, containers, protobuf messages, frozen "
+        "dataclasses, NamedTuples, enums, and classes defining both "
+        "__getstate__ and __setstate__. Any other key type will fail at "
+        "encode time. If your key type is not one of these, register a "
+        "deterministic custom Coder for it or add a type hint so the "
+        "default coder is used.",
+        step_label)
+    return _update_compatible_deterministic_fast_primitives_coder(
+        self, step_label)
 
   def to_type_hint(self):
     return Any
@@ -914,7 +930,20 @@ class PickleCoder(_PickleCoderBase):
         lambda x: dumps(x, protocol), pickle.loads)
 
   def as_deterministic_coder(self, step_label, error_message=None):
-    return FastPrimitivesCoder(self, requires_deterministic=step_label)
+    _LOGGER.warning(
+        "PickleCoder was registered for a key type in '%s', but the runner "
+        "requires a deterministic key encoding and pickle is not "
+        "deterministic. Keys in this step will be encoded with the "
+        "deterministic fallback coder instead of pickle. That coder "
+        "supports primitives, containers, protobuf messages, frozen "
+        "dataclasses, NamedTuples, enums, and classes defining both "
+        "__getstate__ and __setstate__. Any other key type will fail at "
+        "encode time. If your key type is not one of these, register a "
+        "deterministic custom Coder for it or add a type hint so the "
+        "default coder is used.",
+        step_label)
+    return _update_compatible_deterministic_fast_primitives_coder(
+        self, step_label)
 
   def to_type_hint(self):
     return Any
diff --git a/sdks/python/apache_beam/coders/coders_test.py 
b/sdks/python/apache_beam/coders/coders_test.py
index ccd947457ad..17fc46bd9f5 100644
--- a/sdks/python/apache_beam/coders/coders_test.py
+++ b/sdks/python/apache_beam/coders/coders_test.py
@@ -50,6 +50,21 @@ class PickleCoderTest(unittest.TestCase):
     self.assertNotEqual(coders.Base64PickleCoder(), coders.PickleCoder())
     self.assertNotEqual(coders.Base64PickleCoder(), object())
 
+  def test_as_deterministic_coder(self):
+    # PickleCoder.as_deterministic_coder used to construct FastPrimitivesCoder
+    # with a requires_deterministic kwarg that constructor never accepted,
+    # raising TypeError on every call (see coders.FastPrimitivesCoder).
+    v = ('a' * 10, 'b' * 90)
+    deterministic = coders.PickleCoder().as_deterministic_coder('label')
+    self.assertTrue(deterministic.is_deterministic())
+    self.assertEqual(v, deterministic.decode(deterministic.encode(v)))
+
+    memoizing_deterministic = coders._MemoizingPickleCoder(
+    ).as_deterministic_coder('label')
+    self.assertTrue(memoizing_deterministic.is_deterministic())
+    self.assertEqual(
+        v, memoizing_deterministic.decode(memoizing_deterministic.encode(v)))
+
 
 class CodersTest(unittest.TestCase):
   def test_str_utf8_coder(self):

Reply via email to