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.")