Thanks for the suggestion. The branch was merged now, the snapshot will be enabled in a follow-up PR.

 Jan

On 8/20/26 19:42, Kenneth Knowles wrote:
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 <http://session.timeout.ms> = 6000  ->  8.6,
    8.7, 9.1, 10.0 s
    session.timeout.ms <http://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 <http://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