This is an automated email from the ASF dual-hosted git repository.
shunping 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 dabcf50ffbf [YAML] Fix flaky ErrorHandlingTest: don't pickle Scope
into user DoFn closures (#40467)
dabcf50ffbf is described below
commit dabcf50ffbf532bb30a1233927cd1edb6ae067bb
Author: Shunping Huang <[email protected]>
AuthorDate: Thu Oct 8 13:42:48 2026 -0400
[YAML] Fix flaky ErrorHandlingTest: don't pickle Scope into user DoFn
closures (#40467)
* [YAML] Add test reproducing unpicklable provider state leaking into DoFn
closures
A user transform's DoFn that captures `self` also captures the closures
YAML patches onto it, which hold the Scope and every provider. Any
unpicklable provider state (e.g. a protobuf Descriptor) then makes
pickling fail. This test injects such state to reproduce the flaky
ErrorHandlingTest failures deterministically.
* [YAML] Don't pickle Scope (and all providers) into user DoFn closures
---
sdks/python/apache_beam/yaml/yaml_transform.py | 5 +++
.../python/apache_beam/yaml/yaml_transform_test.py | 40 ++++++++++++++++++++++
2 files changed, 45 insertions(+)
diff --git a/sdks/python/apache_beam/yaml/yaml_transform.py
b/sdks/python/apache_beam/yaml/yaml_transform.py
index a4bb9144f8e..0447478d35f 100644
--- a/sdks/python/apache_beam/yaml/yaml_transform.py
+++ b/sdks/python/apache_beam/yaml/yaml_transform.py
@@ -192,6 +192,11 @@ class Scope(LightweightScope):
self.input_providers = input_providers
self._all_followers = None
+ def __reduce__(self):
+ # Scope is construction-only, but closures attached in create_ptransform
+ # can pull it (and all providers) into user DoFn pickles. Stub it out.
+ return str, ('Pickled YAML scope stub.', )
+
def followers(self, transform_name):
if self._all_followers is None:
self._all_followers = collections.defaultdict(list)
diff --git a/sdks/python/apache_beam/yaml/yaml_transform_test.py
b/sdks/python/apache_beam/yaml/yaml_transform_test.py
index 5a1e7b9e220..ec393aefbe6 100644
--- a/sdks/python/apache_beam/yaml/yaml_transform_test.py
+++ b/sdks/python/apache_beam/yaml/yaml_transform_test.py
@@ -834,7 +834,47 @@ class YamlTransformE2ETest(unittest.TestCase):
''')
+class _UnpicklableFactory:
+ """A transform factory holding an object that cloudpickle cannot serialize.
+
+ Used to simulate provider state (e.g. lazily populated caches in the
+ process-global standard providers) that should never be pulled into the
+ pickled closure of an unrelated user transform.
+ """
+ def __init__(self):
+ from apache_beam.portability.api import beam_runner_api_pb2
+ self._descriptor = beam_runner_api_pb2.Pipeline.DESCRIPTOR
+
+ def __call__(self):
+ return beam.Map(lambda x: x)
+
+
class ErrorHandlingTest(unittest.TestCase):
+ def test_closure_does_not_capture_unrelated_providers(self):
+ # SizeLimiter's DoFn closes over `self`. Pickling that DoFn must not drag
in
+ # the YAML Scope and every provider it knows about.
+ with beam.Pipeline(options=beam.options.pipeline_options.PipelineOptions(
+ pickle_library='cloudpickle')) as p:
+ result = p | YamlTransform(
+ '''
+ type: composite
+ transforms:
+ - type: Create
+ config:
+ elements: ['a', 'b', 'biiiiig']
+ - type: SizeLimiter
+ input: Create
+ config:
+ limit: 5
+ error_handling:
+ output: errors
+ output:
+ good: SizeLimiter
+ bad: SizeLimiter.errors
+ ''',
+ providers=dict(TEST_PROVIDERS, Unpicklable=_UnpicklableFactory()))
+ assert_that(result['good'], equal_to(['a', 'b']), label="CheckGood")
+
def test_error_handling_outputs(self):
with beam.Pipeline(options=beam.options.pipeline_options.PipelineOptions(
pickle_library='cloudpickle')) as p: