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,