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]

Reply via email to