Hi Holden,

> Have you looked at how Beam/Flink do their Python stream testing?

I have not looked into how Beam/Flink perform their python testing, however
my proposal is not creating a new framework, it is simply extending the
existing StreamTest framework that exists in Scala to be usable in Python /
Pyspark.

On Wed, Sep 9, 2026 at 11:31 PM Holden Karau <[email protected]> wrote:

> I can take a look next week if we’re not in a rush, otherwise I can take a
> look over the weekend. In general I like the idea, stream testing is
> painful today and making it easier is an A+. Have you looked at how
> Beam/Flink do their Python stream testing?
>
>
> Twitter: https://twitter.com/holdenkarau
> Fight Health Insurance: https://www.fighthealthinsurance.com/
> <https://www.fighthealthinsurance.com/?q=hk_email>
> Books (Learning Spark, High Performance Spark, etc.):
> https://amzn.to/2MaRAG9  <https://amzn.to/2MaRAG9>
> YouTube Live Streams: https://www.youtube.com/user/holdenkarau
> Pronouns: she/her
>
> On Thu, Sep 10, 2026 at 2:28 PM Hyukjin Kwon <[email protected]> wrote:
>
>> +1
>>
>> On Fri, 8 May 2026 at 04:21, Jerry Peng <[email protected]>
>> wrote:
>>
>>>
>>> Hi all,
>>>
>>> I'd like to propose adding a Python-native port of the Scala
>>> `StreamTest` framework
>>> (`sql/core/src/test/scala/org/apache/spark/sql/streaming/StreamTest.scala`)
>>> to PySpark, and invite review of a 4-PR stack that delivers it.
>>>
>>> *Why*
>>>
>>> Scala's `StreamTest` provides a declarative, action-based testing
>>> paradigm used by ~38 Structured Streaming test suites in the Apache Spark
>>> codebase. PySpark has no equivalent today: streaming tests are written
>>> ad-hoc with manual `writeStream` / `processAllAvailable` / `time.sleep`
>>> loops, which makes them verbose, fragile, and slow to develop. As PySpark
>>> streaming usage grows -- including `transformWithState`, Python data
>>> sources, RTM -- the lack of a structured testing framework is increasingly
>>> a productivity tax on PySpark contributors and downstream users writing
>>> PySpark streaming applications.
>>>
>>> *What*
>>>
>>> A new test-only utility package `pyspark.testing.streaming` plus a small
>>> set of test-only helpers in `PythonSQLUtils.scala`. Tests are written as a
>>> sequence of actions:
>>>
>>>   class WordCountTest(StreamTest):
>>>
>>>
>>>
>>>       def test_count(self):
>>>
>>>
>>>
>>>           source = MemoryStream(self.spark, "string")
>>>
>>>           counts = source.to_df().groupBy("value").count()
>>>
>>>
>>>
>>>           self.run_stream_test(
>>>
>>>               counts,
>>>
>>>
>>>
>>>               AddData(source, "a", "b", "a"),
>>>
>>>               CheckAnswer(("a", 2), ("b", 1)),
>>>
>>>
>>>
>>>               output_mode="complete",
>>>
>>>
>>>
>>>           )
>>>
>>> The framework supports the full action vocabulary you'd expect from the
>>> Scala equivalent -- `StartStream`, `StopStream`, `AddData`,
>>> `ProcessAllAvailable`, `CheckAnswer` / `CheckLastBatch` / `CheckNewAnswer`
>>> / `CheckAnswerByFunc`, `Assert` / `AssertOnQuery` / `Execute`,
>>> `ExpectFailure` -- plus Real-Time Mode (`LowLatencyMemoryStream`,
>>> `ContinuousMemorySink`, `Trigger.RealTime`) via polling Check actions.
>>>
>>> *How*
>>>
>>> The strategy is to keep the JVM streaming runtime as the source of truth
>>> and put a thin Python layer on top of it. Rather than reimplement memory
>>> source and sink in python, the framework wraps the existing JVM
>>> `MemoryStream[Row]` / `MemorySink` / `LowLatencyMemoryStream` /
>>> `ContinuousMemorySink` through a handful of test-only helpers added to
>>> `PythonSQLUtils.scala`, and exposes a Python `unittest.TestCase` subclass
>>> (`StreamTest`) whose `run_stream_test()` driver walks an action list and
>>> dispatches against those helpers. The result is action-for-action parity
>>> with Scala `StreamTest` for the common case, idiomatic Python at the call
>>> site (flexible expected-row shapes, dataclass-style actions, type hints),
>>> and -- because the JVM does the heavy lifting -- minimal risk of behavioral
>>> drift in testing semantics.
>>>
>>>
>>> *PR stack*
>>>
>>>   
>>> ┌─────┬────────────────────────────┬────────────────────────────────────────────────────────────────────────┬───────┐
>>>
>>>
>>>   │  #  │           Branch           │
>>> Scope                                  │ Tests │
>>>
>>>   
>>> ├─────┼────────────────────────────┼────────────────────────────────────────────────────────────────────────┼───────┤
>>>
>>>
>>>   │ 1   │ streamtest-py-1-bridge     │ JVM bridge + MemoryStream Python
>>> wrapper                               │ 11    │
>>>
>>>
>>> ├─────┼────────────────────────────┼────────────────────────────────────────────────────────────────────────┼───────┤
>>>
>>>   │ 2   │ streamtest-py-2-base       │ StreamTest base +
>>> lifecycle/assertion actions                          │ +12   │
>>>
>>>
>>>   
>>> ├─────┼────────────────────────────┼────────────────────────────────────────────────────────────────────────┼───────┤
>>>
>>>
>>>   │ 3   │ streamtest-py-3-checks     │ CheckAnswer family +
>>> flexible-type expected values                     │ +26   │
>>>
>>>
>>>   
>>> ├─────┼────────────────────────────┼────────────────────────────────────────────────────────────────────────┼───────┤
>>>
>>>
>>>   │ 4   │ streamtest-py-4-rtm        │ Real-Time Mode
>>> (LowLatencyMemoryStream, ContinuousMemorySink, polling) │ +13   │
>>>
>>>   
>>> ├─────┼────────────────────────────┼────────────────────────────────────────────────────────────────────────┼───────┤
>>>
>>>
>>>
>>>
>>> The full design doc, including the architecture diagram and usage
>>> examples, lives at:
>>>
>>>
>>> https://github.com/jerrypeng/spark/blob/stack/streamtest-py-5-design-doc/python/pyspark/testing/streaming/DESIGN.md
>>>
>>> *Asks*
>>>
>>> I would appreciate feedback on this effort.  I did mostly use claude to
>>> help me do this conversation.
>>>
>>> It would be great if I can get reviews on PRs for this starting from the
>>> first PR:
>>>
>>> https://github.com/apache/spark/pull/55656
>>>
>>> Thanks,
>>> Jerry
>>>
>>

Reply via email to