tkaymak opened a new pull request, #40281: URL: https://github.com/apache/beam/pull/40281
Part of #36841, slice 5b. Stateful ParDo, with state and event time timers, now runs on the Spark 4 Structured Streaming runner through Spark 4 `transformWithState`. Before this change, the streaming translator rejected it. **How it works** - User state reuses the legacy `SparkStateInternals`. A new `StateCells` interface lets it run on a Spark `MapState`. The existing callers keep the Table backed path. - Timers reuse the legacy `SparkTimerInternals`. They are persisted per key in a Spark `ValueState`. Each key gets one Spark timer as a wake up, at the earliest pending Beam timer plus 1 ms. Spark fires when the expiry is less than or equal to the watermark, and Beam fires only when the watermark is strictly after the timestamp, so the extra 1 ms matches the two rules. Timers fire one at a time, so a callback can still set or delete other timers. - The per task runner lifecycle is shared with the batch stateful ParDo through `StatefulTaskRunner`. That lifecycle covers setup on the executor, metrics, and teardown on task completion or failure. - These are still rejected at translation, with a message naming #36841: processing time and synchronized processing time timers, `@OnWindowExpiration`, `@RequiresTimeSortedInput`, and merging windows. **Watermark semantics** The runner's input watermark is Spark's event time watermark, as set up in the earlier slices: max event time minus `--watermarkDelayMillis`. Lateness and timer firing both follow that watermark. For out of order sources, raise the delay. Spark has no watermark hold for timer output timestamps. That is harmless while only stateless transforms follow, and it is tracked as a blocker for GroupByKey in 5c (#36841 comment). **Tests** - `StatefulParDoStreamingTest` runs a live query. It covers state across micro batches, timer firing, a callback that deletes a later due timer, and a checkpoint restart. - `PipelineTranslatorStreamingTest` has one table driven test for the rejected shapes. - The existing batch stateful ParDo tests and the legacy `SparkStateInternalsTest` still pass on Spark 3 and Spark 4. CHANGES.md follows with GroupByKey in 5c. A ValidatesRunner suite follows once Impulse and PAssert are supported. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
