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__':

Reply via email to