Skip to content

Commit df5c833

Browse files
beaver-felixDullaho
authored andcommitted
[SPARK-56565][CONNECT][TESTS] Retry flaky python foreachBatch termination test
### What changes were proposed in this pull request? This PR un-skips the Connect parity test `test_streaming_foreach_batch_graceful_stop` and applies the retry-timeout infrastructure introduced in SPARK-52843 (`eventually` and `timeout`). Additionally, it replaces the py4j-specific `_jvm.java.lang.Thread.sleep()` call with standard Python `time.sleep()`. ### Why are the changes needed? This resolves [SPARK-56565](https://issues.apache.org/jira/browse/SPARK-56565). The test for `foreachBatch` graceful stop was previously skipped for Connect because it inherently relied on Py4J's pinned thread execution model and caused flaky timeout behaviors. With the new retry infrastructure, we can safely re-enable it for the Connect module. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Verified by standard GitHub Actions on personal fork. ### Was this patch authored or co-authored using generative AI tooling? No. Closes #55481 from Feelitx/master. Lead-authored-by: Felix <160706963+Feelitx@users.noreply.github.com> Co-authored-by: Luqina <cuong.m9325@gmail.com> Signed-off-by: Hyukjin Kwon <gurwls223@apache.org>
1 parent 6f26070 commit df5c833

1 file changed

Lines changed: 12 additions & 4 deletions

File tree

‎python/pyspark/sql/tests/connect/streaming/test_parity_foreach_batch.py‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,10 +15,10 @@
1515
# limitations under the License.
1616
#
1717

18-
import unittest
19-
18+
import time
2019
from pyspark.sql.tests.streaming.test_streaming_foreach_batch import StreamingTestsForeachBatchMixin
2120
from pyspark.testing.connectutils import ReusedConnectTestCase, should_test_connect
21+
from pyspark.testing.utils import eventually, timeout
2222
from pyspark.errors import PySparkPicklingError
2323

2424
if should_test_connect:
@@ -29,9 +29,17 @@ class StreamingForeachBatchParityTests(StreamingTestsForeachBatchMixin, ReusedCo
2929
def test_streaming_foreach_batch_propagates_python_errors(self):
3030
super().test_streaming_foreach_batch_propagates_python_errors()
3131

32-
@unittest.skip("This seems specific to py4j and pinned threads. The intention is unclear")
32+
@eventually(timeout=180, catch_timeout=True)
33+
@timeout(timeout=60)
3334
def test_streaming_foreach_batch_graceful_stop(self):
34-
super().test_streaming_foreach_batch_graceful_stop()
35+
# SPARK-39218: Make foreachBatch streaming query stop gracefully
36+
def func(batch_df, _):
37+
time.sleep(10)
38+
39+
q = self.spark.readStream.format("rate").load().writeStream.foreachBatch(func).start()
40+
time.sleep(3) # 'rowsPerSecond' defaults to 1. Waits 3 secs out for the input.
41+
q.stop()
42+
self.assertIsNone(q.exception(), "No exception has to be propagated.")
3543

3644
def test_nested_dataframes(self):
3745
def curried_function(df):

0 commit comments

Comments
 (0)