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
>
>

Reply via email to