This is an automated email from the ASF dual-hosted git repository. Amar3tto pushed a commit to branch leaderboard-test in repository https://gitbox.apache.org/repos/asf/beam.git
commit 9ca8a321fe0b2c3f5c1392bd0a9dc39b6ada5275 Author: Vitaly Terentyev <[email protected]> AuthorDate: Wed Sep 30 16:22:42 2026 +0400 delay inject --- .../examples/complete/game/leader_board_it_test.py | 20 +++++++++++++++----- 1 file changed, 15 insertions(+), 5 deletions(-) diff --git a/sdks/python/apache_beam/examples/complete/game/leader_board_it_test.py b/sdks/python/apache_beam/examples/complete/game/leader_board_it_test.py index 218d4e4a778..7f861cc759e 100644 --- a/sdks/python/apache_beam/examples/complete/game/leader_board_it_test.py +++ b/sdks/python/apache_beam/examples/complete/game/leader_board_it_test.py @@ -34,6 +34,7 @@ Usage: # pytype: skip-file import logging +import threading import time import unittest import uuid @@ -67,6 +68,7 @@ class LeaderBoardIT(unittest.TestCase): WAIT_UNTIL_FINISH_DURATION = 15 * 60 * 1000 # in milliseconds # Poll BigQuery after the pipeline wait; streaming inserts can lag. BQ_MATCHER_TIMEOUT_SECS = 10 * 60 + PUBLISH_DELAY_SECS = 2 * 60 def setUp(self): self.test_pipeline = TestPipeline(is_integration_test=True) @@ -167,14 +169,22 @@ class LeaderBoardIT(unittest.TestCase): self.addCleanup(self._cleanup_pubsub) self.addCleanup(utils.delete_bq_dataset, self.project, self.dataset_ref) - # Generate input data and inject to PubSub. - self._inject_pubsub_game_events(self.input_topic, self.DEFAULT_INPUT_COUNT) + def _publish_after_delay(): + time.sleep(self.PUBLISH_DELAY_SECS) + self._inject_pubsub_game_events( + self.input_topic, self.DEFAULT_INPUT_COUNT) - # Get pipeline options from command argument: --test-pipeline-options, - # and start pipeline job by calling pipeline main function. - leader_board.run( + publish_thread = threading.Thread(target=_publish_after_delay, daemon=True) + publish_thread.start() + + try: + # Get pipeline options from command argument: --test-pipeline-options, + # and start pipeline job by calling pipeline main function. + leader_board.run( self.test_pipeline.get_full_options_as_args(**extra_opts), save_main_session=False) + finally: + publish_thread.join(timeout=self.PUBLISH_DELAY_SECS + 60) if __name__ == '__main__':
