SGTM. Rather have it on master than developed on a branch further. Nice job setting it up so it won't be released before it is ready. We could include it in snapshots if you want something obviously-less-stable that you could point someone to for testing.
Kenn On Tue, Aug 18, 2026 at 6:37 AM Jan Lukavský <[email protected]> wrote: > Hi, > > TL;DR summary: > > The runner is a working skeleton with a well defined set of supported > features. Because the set of features is quite narrow and it is currently > not extensively tested against real-world use-cases, it is currently built > only with a specific profile (-Pwith-kafka-streams-runner), which prevents > us from releasing not well-tested or incomplete runner. The purpose of > merging is to enable broader (sub)community to eventually emerge and > develop/maintain it. This could include cooperation with the Apache Kafka > community, which can be pulled-in once the skeleton is merged. > > Jan > On 8/18/26 12:24, Junaid wrote: > > Hi all, > > I would like to propose merging the Kafka Streams runner from its feature > branch into master. It has been developed this summer as a Google Summer of > Code project under the tracking issue [1] > <https://github.com/apache/beam/issues/18479>, mentored by Jan Lukavský. > The pull request is at [2] <https://github.com/apache/beam/pull/39785>. > > What it is. A portable runner that translates a Beam pipeline into a Kafka > Streams topology and executes user code over the Fn API. What makes it > different from the other runners is that Kafka Streams is a library, not a > cluster: there is no job manager and no resource manager to operate. A > pipeline is an ordinary JVM process that reads from and writes to Kafka, > and you scale it by starting more copies of that process. State, fault > tolerance and exactly-once come from Kafka itself, through consumer groups, > changelog topics and transactions. > > Current state. It is a skeleton. It runs a real subset of the model: > bounded and unbounded reads, stateless ParDo including multiple outputs, > GroupByKey and Combine, global, fixed and sliding windows with the default > trigger and allowed lateness, Flatten, Redistribute, metrics, and > exactly-once via Kafka transactions. That subset is covered by 105 unit > tests, 59 of Beam's own @ValidatesRunner tests, and integration tests > against a real broker. Beam's Python portable suite also runs against it, > so it is exercised from a non-Java SDK as well. > > It also has real gaps, each tracked: side inputs (#39628 > <https://github.com/apache/beam/issues/39628>), stateful ParDo and user > timers (#39629 <https://github.com/apache/beam/issues/39629>), merging > windows and custom WindowFns (#39630 > <https://github.com/apache/beam/issues/39630>), splittable DoFn (#39631 > <https://github.com/apache/beam/issues/39631>), TestStream (#39632 > <https://github.com/apache/beam/issues/39632>), and reading a source in > parallel (#39626 <https://github.com/apache/beam/issues/39626>). > > And it has at least one known bug rather than a missing feature: bundles > are not closed after a bounded time (#39633 > <https://github.com/apache/beam/issues/39633>). The maxBundleTimeMs > option is accepted and has no effect, because closing a bundle from a > wall-clock punctuator duplicated output against a real broker and the cause > is not yet understood. There may be others we have not found. > > Why I propose merging now. The motivation was to build a skeleton that can > be developed further by several contributors, and master is where that can > happen: the runner can be built, run and worked on by anyone interested, > and the gaps above are well-defined pieces of work someone could pick up. > > To make that safe, the runner is not part of the standard build. Its > subprojects are only included when -Pwith-kafka-streams-runner is passed: > > ./gradlew -Pwith-kafka-streams-runner :runners:kafka-streams:build > > Without the flag they are not in the build at all, so nothing reaches > users who have not asked for it, and no release artifact contains it. It is > built in the Java precommit, so it cannot rot unnoticed. If the runner > becomes stable enough the flag comes off and it is built like any other > runner; if it does not, it can be dropped again without affecting anyone, > because no release ever shipped it. > > What is the potential of the runner. Two things, one operational and one > about recovery. > > The operational one is that there is nothing to operate. If you already > run Kafka, a Beam pipeline becomes an ordinary application you deploy like > any other — no job manager, no resource manager, no second distributed > system to size, upgrade and keep alive. Scaling up or down is starting or > stopping a process, and the consumer group redistributes the work. > > The recovery one is that Kafka Streams reassigns partitions and restores > state from a changelog, where a checkpoint-based engine restarts a job from > its last checkpoint. On a Mac, with one broker and two instances, killing > the instance holding the source read with kill -9 and timing until the > other took the work over: > > session.timeout.ms = 6000 -> 8.6, 8.7, 9.1, 10.0 s > session.timeout.ms = 45000 -> 55.7, 55.8 s (Kafka's default) > > Handover is dominated by how long the consumer group takes to notice, > which is session.timeout.ms and is configurable; the recovery work itself > is the remainder. These are laptop numbers meant to show the shape of the > thing, not a benchmark against other runners. > > This lazy consensus request will be open for at least 72 hours. If there > are no objections by then the consensus will pass. Any comments or > objections are welcome, here or on the pull request. > > Thanks, > Junaid Shaukat > https://github.com/junaiddshaukat > > [1] https://github.com/apache/beam/issues/18479 > [2] https://github.com/apache/beam/pull/39785 > >
