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

damccorm 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 5e5aaf3ce8f reduce LeaderBoardIT flakiness (#40241)
5e5aaf3ce8f is described below

commit 5e5aaf3ce8fece91547ba6bbca9a47a9c000167c
Author: Abdelrahman Ibrahim <[email protected]>
AuthorDate: Tue Sep 29 16:47:38 2026 +0300

    reduce LeaderBoardIT flakiness (#40241)
    
    * reduce LeaderBoardIT flakiness
    
    * increase timeout
    
    * add wait for PubSub publish
    
    * trigger workflow
    
    * trigger workflow
    
    * wait for PubSub subscription readiness
---
 .github/trigger_files/beam_PostCommit_Python.json  |  4 ++--
 .../examples/complete/game/leader_board_it_test.py | 25 +++++++++++++++++-----
 2 files changed, 22 insertions(+), 7 deletions(-)

diff --git a/.github/trigger_files/beam_PostCommit_Python.json 
b/.github/trigger_files/beam_PostCommit_Python.json
index 9bfa92ff510..e11d2260913 100644
--- a/.github/trigger_files/beam_PostCommit_Python.json
+++ b/.github/trigger_files/beam_PostCommit_Python.json
@@ -1,5 +1,5 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run.",
-  "pr": "40139",
-  "modification": 59
+  "pr": "40241",
+  "modification": 63
 }
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 4cc13171fe9..218d4e4a778 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
@@ -64,7 +64,9 @@ class LeaderBoardIT(unittest.TestCase):
   OUTPUT_TABLE_TEAMS = 'leader_board_teams'
   DEFAULT_INPUT_COUNT = 500
 
-  WAIT_UNTIL_FINISH_DURATION = 10 * 60 * 1000  # in milliseconds
+  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
 
   def setUp(self):
     self.test_pipeline = TestPipeline(is_integration_test=True)
@@ -84,6 +86,8 @@ class LeaderBoardIT(unittest.TestCase):
         name=self.sub_client.subscription_path(
             self.project, self.INPUT_SUB + _unique_id),
         topic=self.input_topic.name)
+    # New subscriptions can miss messages published immediately after create.
+    time.sleep(30)
 
     # Set up BigQuery environment
     self.dataset_ref = utils.create_bq_dataset(
@@ -97,9 +101,14 @@ class LeaderBoardIT(unittest.TestCase):
     logging.debug(
         'Injecting %d game events to topic %s', message_count, topic.name)
 
+    publish_futures = []
     for _ in range(message_count):
-      self.pub_client.publish(
-          topic.name, (self.INPUT_EVENT % 
self._test_timestamp).encode('utf-8'))
+      publish_futures.append(
+          self.pub_client.publish(
+              topic.name,
+              (self.INPUT_EVENT % self._test_timestamp).encode('utf-8')))
+    for future in publish_futures:
+      future.result()
 
   def _cleanup_pubsub(self):
     test_utils.cleanup_subscriptions(self.sub_client, [self.input_sub])
@@ -124,7 +133,10 @@ class LeaderBoardIT(unittest.TestCase):
             self.OUTPUT_TABLE_USERS,
             success_condition))
     bq_users_verifier = BigqueryMatcher(
-        self.project, users_query, self.DEFAULT_EXPECTED_CHECKSUM)
+        self.project,
+        users_query,
+        self.DEFAULT_EXPECTED_CHECKSUM,
+        timeout_secs=self.BQ_MATCHER_TIMEOUT_SECS)
 
     teams_query = (
         'SELECT total_score FROM `%s.%s.%s` '
@@ -134,7 +146,10 @@ class LeaderBoardIT(unittest.TestCase):
             self.OUTPUT_TABLE_TEAMS,
             success_condition))
     bq_teams_verifier = BigqueryMatcher(
-        self.project, teams_query, self.DEFAULT_EXPECTED_CHECKSUM)
+        self.project,
+        teams_query,
+        self.DEFAULT_EXPECTED_CHECKSUM,
+        timeout_secs=self.BQ_MATCHER_TIMEOUT_SECS)
 
     extra_opts = {
         'allow_unsafe_triggers': True,

Reply via email to