From 4ce21ca50cd270c802c1b9eeb2e91d1a71c77a34 Mon Sep 17 00:00:00 2001 From: Tian Gao Date: Thu, 23 Jul 2026 13:14:48 -0700 Subject: [PATCH] Reduce the expected progress for streaming data source tests --- .../pyspark/sql/tests/test_python_streaming_datasource.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/python/pyspark/sql/tests/test_python_streaming_datasource.py b/python/pyspark/sql/tests/test_python_streaming_datasource.py index f6bfbdc655f58..704a2986ded88 100644 --- a/python/pyspark/sql/tests/test_python_streaming_datasource.py +++ b/python/pyspark/sql/tests/test_python_streaming_datasource.py @@ -305,7 +305,7 @@ def check_batch(df, batch_id): assertDataFrameEqual(df, [Row(batch_id * 2), Row(batch_id * 2 + 1)]) q = df.writeStream.foreachBatch(check_batch).start() - wait_for_condition(q, lambda query: len(query.recentProgress) >= 10) + wait_for_condition(q, lambda query: len(query.recentProgress) >= 5) q.stop() q.awaitTermination() self.assertIsNone(q.exception(), "No exception has to be propagated.") @@ -400,7 +400,7 @@ def check_batch(df, batch_id): assertDataFrameEqual(df, [Row(batch_id * 2), Row(batch_id * 2 + 1)]) q = df.writeStream.foreachBatch(check_batch).start() - wait_for_condition(q, lambda query: len(query.recentProgress) >= 10) + wait_for_condition(q, lambda query: len(query.recentProgress) >= 5) q.stop() q.awaitTermination() self.assertIsNone(q.exception(), "No exception has to be propagated.") @@ -456,7 +456,7 @@ def check_batch(df, batch_id): assertDataFrameEqual(df, [Row(batch_id * 2), Row(batch_id * 2 + 1)]) q = df.writeStream.foreachBatch(check_batch).start() - wait_for_condition(q, lambda query: len(query.recentProgress) >= 10) + wait_for_condition(q, lambda query: len(query.recentProgress) >= 5) q.stop() q.awaitTermination() self.assertIsNone(q.exception(), "No exception has to be propagated.")