From f68e589807a0434c41c43c781568df73bac89ed5 Mon Sep 17 00:00:00 2001 From: Vitaly Terentyev Date: Mon, 1 Sep 2025 11:22:47 +0400 Subject: [PATCH] Fix tests --- .../apache_beam/transforms/external_test.py | 21 +++++++++++++++++-- .../apache_beam/transforms/util_test.py | 4 +++- 2 files changed, 22 insertions(+), 3 deletions(-) diff --git a/sdks/python/apache_beam/transforms/external_test.py b/sdks/python/apache_beam/transforms/external_test.py index 2ed7d622ecd6..cead3e773037 100644 --- a/sdks/python/apache_beam/transforms/external_test.py +++ b/sdks/python/apache_beam/transforms/external_test.py @@ -33,6 +33,7 @@ from apache_beam import Pipeline from apache_beam.coders import RowCoder from apache_beam.options.pipeline_options import PipelineOptions +from apache_beam.pipeline import PipelineVisitor from apache_beam.portability.api import beam_expansion_api_pb2 from apache_beam.portability.api import external_transforms_pb2 from apache_beam.portability.api import schema_pb2 @@ -349,6 +350,22 @@ def test_external_transform_finder_non_leaf(self): self.assertTrue(pipeline.contains_external_transforms) def test_external_transform_finder_leaf(self): + def has_external(p): + class Finder(PipelineVisitor): + def __init__(self): + self.matches = [] + + def enter_composite_transform(self, node): + if isinstance(node.transform, beam.ExternalTransform): + self.matches.append(node.full_label) + + def visit_transform(self, node): + self.enter_composite_transform(node) + + finder = Finder() + p.visit(finder) + return finder.matches + pipeline = beam.Pipeline() _ = ( pipeline @@ -357,9 +374,9 @@ def test_external_transform_finder_leaf(self): 'beam:transforms:xlang:test:nooutput', ImplicitSchemaPayloadBuilder({'data': '0'}), expansion_service.ExpansionServiceServicer())) - pipeline.run().wait_until_finish() - self.assertTrue(pipeline.contains_external_transforms) + found = has_external(pipeline) + self.assertTrue(found, f"No ExternalTransform found; saw: {found}") def test_sanitize_java_traceback(self): error_string = ''' diff --git a/sdks/python/apache_beam/transforms/util_test.py b/sdks/python/apache_beam/transforms/util_test.py index b365d9b22090..051f2ab7ade3 100644 --- a/sdks/python/apache_beam/transforms/util_test.py +++ b/sdks/python/apache_beam/transforms/util_test.py @@ -762,7 +762,7 @@ def process(self, element): equal_to(expected_windows), label='before_identity', reify_windows=True) - _ = ( + after = ( before_identity | 'window' >> beam.WindowInto( beam.transforms.util._IdentityWindowFn( @@ -772,6 +772,8 @@ def process(self, element): # contain a window of None. IdentityWindowFn should # raise an exception. | 'add_timestamps2' >> beam.ParDo(AddTimestampDoFn())) + # Terminal consumer to prevent pruning: + _ = after | 'count_per_element' >> beam.combiners.Count.PerElement() class ReshuffleTest(unittest.TestCase):