Hi Yi, Thanks, that gives us a practical first milestone. I agree with adding a Structured Streaming variant of the existing Spark 4 Java streaming PVR suite, using a flag such as useStructuredStreaming, and switching the default for portable streaming pipelines once that variant passes with the current exclusions.
That separates the initial backend migration from filling the remaining streaming gaps. We can then add Spark 4 PVR suites for Python and Go and expand coverage as the missing capabilities are implemented. I’ll update the proposal to reflect this sequence and distinguish the default switch from eventual removal of the old path. I’ll also review the current suite and exclusions to scope the first implementation step. Thanks, Elia On Wed, Sep 9, 2026 at 1:17 AM Yi Hu <[email protected]> wrote: > Hi Elia, > > Thanks for the proposal. I agree the dataset based runner will be the way > to go. As for logistics > > > - The existing RDD/DStream portable path stays available during > > migration and is removed only after the Dataset portable path meets > agreed > > correctness, recovery, performance, and operability gates. > > We can use the existing streaming portable runner test > > > https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_PVR_Spark4_Streaming.yml > > as the "exit gate". i.e., we can setup a structured streaming variant for > this test by feeding a flag (like useStructuredStreaming). We could flip > the default to true for streaming pipelines when the following milestone > met: structured streaming variant passed this test suite. > > Because currently this test suite almost only covers batch pipelines after > many exclude categories, this milestone should be largely feasible. > > After dataset-based runner becomes default for streaming, moving forward > we can setup Spark 4 PVR test suites for Python and Go (Currently only > Spark 3 tests are exercised) so we have a more well-defined runner > capability validate tests for each SDK and can fixing the gaps in a defined > way. > > Thanks, > > Yi > > On 2026/09/04 14:32:57 Elia Liu wrote: > > Hi Beam community, > > > > I would like to get feedback on a direction for Spark 4 portable runner > > support. > > > > Proposal: Spark 4 Portable Runner Direction > > < > https://docs.google.com/document/d/1nDHVif-XAcTk3AOjVpBiTSNfl_JytPJLxjTCodHWVvk/edit?usp=drivesdk > > > > > > The main motivation is that the current Spark portable streaming path is > > quite limited. It still relies on DStream, and several basic capabilities > > are either unsupported or have long-standing open gaps. Portable SDF > > support on Spark is still tracked in #19468 > > <https://github.com/apache/beam/issues/19468>, and the UnboundedSource > > ValidatesRunner coverage merged in #38892 > > <https://github.com/apache/beam/pull/38892> explicitly skips Spark for > that > > reason. The Spark portable streaming tracking issues for side inputs, > > state, and timers are also still open (#20395 > > <https://github.com/apache/beam/issues/20395>, #20396 > > <https://github.com/apache/beam/issues/20396>, #20397 > > <https://github.com/apache/beam/issues/20397>). There are also open > > tracking issues around portable windowing (#20429 > > <https://github.com/apache/beam/issues/20429>) and Spark streaming > > watermark propagation (#23129 < > https://github.com/apache/beam/issues/23129>). > > Separately, the DStream-based Spark streaming implementation still has > > watermark correctness bugs; a recent example is #39822 > > <https://github.com/apache/beam/issues/39822>, where an idle source can > > disappear from the stored watermark map and cause downstream timers to > > advance too early. > > > > For general portable streaming workloads, this makes the current Spark > path > > hard to use in practice. > > > > The direction I am proposing is to converge Spark 4 on one Dataset-based > > backend: > > > > - Java Pipeline and Runner API use separate planning adapters. > > - Fully bounded pipelines remain normal Spark batch jobs. > > - Pipelines containing unbounded data use Structured Streaming. > > - The existing RDD/DStream portable path stays available during > > migration and is removed only after the Dataset portable path meets > agreed > > correctness, recovery, performance, and operability gates. > > > > The main new implementation work is the Runner API-to-Dataset planning > > adapter and the ExecutableStage/Fn bridge. I tried to make that boundary > > explicit in the short proposal, rather than specifying all of the > follow-up > > streaming semantics here. > > > > At this stage I am mainly looking for feedback on the direction itself. > In > > particular, does one Dataset-based Spark 4 backend seem like the right > > long-term target, and is the proposed split between Java/Runner API > > planning adapters and batch/streaming runtime modes sensible? I would > also > > be interested in any existing Spark portable streaming gaps or migration > > constraints that should change this direction. > > > > Thanks, > > Elia > > >
