This is an automated email from the ASF dual-hosted git repository.
github-actions[bot] pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from c1017953ac3 Adds support for reading at a given Delta Lake version or
timestamp (#39758)
add 2577a6cff4a Install Cloud Spanner emulator component in Go PreCommit
CI (#39850)
add 2e445307153 Add Kafka Streams runner skeleton module and portable
entry points
add 61c891a69ca Address review notes on KafkaStreamsPipelineResult and
state dir
add cef6544e792 Address review feedback on Kafka Streams Runner skeleton
add 0d445f79e0c Drop setRunner(null) suppression; make applicationId
required
add d4c7dae98b8 Catch Exception in KafkaStreamsRunner.run() to avoid
job-server leak
add 64ad482f9de Merge pull request #38534: [GSoC 2026] Kafka Streams
runner skeleton module + portable entry points
add faefc95e9f0 [GSoC 2026] Kafka Streams runner — translation framework +
Impulse translator (#38689)
add 0f875a07162 Add temporary feature-branch CI for Kafka Streams runner
(#38725)
add ff554f7438d [GSoC 2026] Kafka Streams runner — ExecutableStage
(stateless ParDo) translator (#38764)
add 4c02a7789a2 [GSoC 2026] Kafka Streams runner — Redistribute translator
+ ExecutableStage type-agnostic edge (#38843)
add 4636b1f333e #38957: Add in-memory WatermarkManager core
(per-source-partition tracking
add 1058e94027e [GSoC 2026] Kafka Streams runner #38987: Wire
WatermarkManager into ExecutableStageProcessor
add 3e963a68343 [GSoC 2026] Kafka Streams runner #39051: Add
KStreamsPayload Serde for crossing topic boundaries
add 5e65d4772fb [GSoC 2026] Kafka Streams runner #39141: Add GroupByKey
(GlobalWindow, fire at watermark)
add 4d5847f6bfa [GSoC 2026] Kafka Streams runner #39211: Add
KafkaStreamsTestRunner test harness
add 27c4522ae5e [GSoC 2026] Kafka Streams runner #39249: Support Create
add ff323ebd6db [GSoC 2026] Kafka Streams runner #39273: Kafka Streams
runner: Flatten support
add 75adf49ee8d [GSoC 2026] Kafka Streams runner: surface SDK-harness
metrics as MetricResults (#39341)
add 87911f00973 [GSoC 2026] Kafka Streams runner:
TestPipeline-dispatchable test runner (PAssert works) (#39362)
add d9cea589e5b [GSoC 2026] Kafka Streams runner: validatesRunner task;
Create and Flatten suites green (#39380)
add fca14356614 [GSoC 2026] Kafka Streams runner: multi-output executable
stages (#39410)
add b615ae85aed [GSoC 2026] Kafka Streams runner: enable ParDoTest in the
ValidatesRunner suite (#39451)
add 47614a943a7 [GSoC 2026] Kafka Streams runner: windowed GroupByKey via
ReduceFnRunner (#39494)
add 27198768f89 [GSoC 2026] Kafka Streams runner: run on a real broker,
correctly across partitions (#39546)
add cf1f10bdf82 [GSoC 2026] Kafka Streams runner: bound a bundle by
element count (#39578)
add fc36301e6c2 [GSoC 2026] Kafka Streams runner: CombineTest coverage and
two review follow-ups (#39610)
add 5861f31e8ac [GSoC 2026] Kafka Streams runner: read unbounded sources
(#39611)
add cb30afd092e [GSoC 2026] Kafka Streams runner: user documentation,
marked experimental (#39627)
add e051e06bf83 [GSoC 2026] Kafka Streams runner: Python wrapper that
starts its own job server (#39680)
add 4ff618047c3 [GSoC 2026] Kafka Streams runner: terminate a bounded
pipeline when it is drained (#39700)
add 5f1658f8b5e [GSoC 2026] Kafka Streams runner: portable ValidatesRunner
suite for Python (#39736)
add 511a40e4f20 [GSoC 2026] Kafka Streams runner: separate the source's
poll size from the bundle size, and expose the session timeout (#39748)
add 65a2e400c73 [GSoC 2026] Kafka Streams runner: bound a source poll in
time, not only in elements (#39761)
add 104dc272e8d [GSoC 2026] Kafka Streams runner: ask for primitive reads
in the Java wrapper (#39766)
add 49d459a243e [GSoC 2026] Kafka Streams runner: put the runner behind an
opt-in build flag (#39762)
add 10ff557d616 [GSoC 2026] Kafka Streams runner: an application for
measuring instances coming and going (#39752)
add 5b3702484f4 [GSoC 2026] Kafka Streams runner: license header and
Python formatting for master CI
add b12bfc159ba [GSoC 2026] Kafka Streams runner: shorten the explanation
comments
add 8d5f6511a6b Merge pull request #39781: [GSoC 2026] Kafka Streams
runner: shorten comments, and fix the license header and Python formatting
add e2681085c96 Build Kafka Streams runner during javaPreCommit (#18479)
add 8b3b08dc799 Merge pull request #39784: Build Kafka Streams runner
during javaPreCommit (#18479)
add b67cf234ead [GSoC 2026] Kafka Streams runner: update the CHANGES.md
entry
add 52dadbac433 Merge pull request #39786: [GSoC 2026] Kafka Streams
runner: update the CHANGES.md entry
add 89617f59b45 Merge branch 'master' of https://github.com/apache/beam
into feat/18479-kafka-streams-runner-skeleton
add 6b75a44f0da Removed feature branch build
add 52d6eeab74e [GSoC 2026] Kafka Streams runner: say what the Python side
is, and is not
add 473a5dd649c Merge pull request #39847: [GSoC 2026] Kafka Streams
runner: say what the Python side is, and is not
add e10e8c969b0 Merge pull request #39785: Merge Kafka Streams Runner
skeleton (#18479)
No new revisions were added by this update.
Summary of changes:
... beam_PreCommit_Java_Kafka_Streams_Runner.json} | 0
.github/workflows/README.md | 1 +
.github/workflows/beam_PreCommit_Go.yml | 2 +
.github/workflows/beam_PreCommit_Java.yml | 1 +
...> beam_PreCommit_Java_Kafka_Streams_Runner.yml} | 72 +++-
CHANGES.md | 1 +
build.gradle.kts | 4 +
runners/kafka-streams/build.gradle | 201 +++++++++
runners/kafka-streams/job-server/build.gradle | 88 ++++
runners/kafka-streams/measurement/build.gradle | 78 ++++
.../kafka-streams/measurement/docker-compose.yml | 44 ++
.../streams/measurement/RescalingMeasurement.java | 246 +++++++++++
.../kafka/streams/measurement/package-info.java | 26 ++
runners/kafka-streams/proto/build.gradle | 34 ++
.../src/main/proto/kafka_streams_payload.proto | 58 +++
.../kafka/streams/KafkaStreamsJobInvoker.java | 96 +++++
.../kafka/streams/KafkaStreamsJobServerDriver.java | 106 +++++
.../kafka/streams/KafkaStreamsPipelineOptions.java | 151 +++++++
.../kafka/streams/KafkaStreamsPipelineResult.java | 77 ++++
.../kafka/streams/KafkaStreamsPipelineRunner.java | 185 ++++++++
.../KafkaStreamsPortablePipelineResult.java | 154 +++++++
.../runners/kafka/streams/KafkaStreamsRunner.java | 138 ++++++
.../kafka/streams/KafkaStreamsRunnerRegistrar.java | 48 +++
.../kafka/streams/KafkaStreamsTopicManager.java | 171 ++++++++
.../beam/runners/kafka/streams/package-info.java | 20 +
.../streams/translation/EmptyBoundedSource.java | 89 ++++
.../translation/ExecutableStageProcessor.java | 365 ++++++++++++++++
.../translation/ExecutableStageTranslator.java | 133 ++++++
.../streams/translation/FlattenProcessor.java | 124 ++++++
.../streams/translation/FlattenTranslator.java | 95 +++++
.../GroupByKeyBroadcastPartitioner.java | 70 ++++
.../streams/translation/GroupByKeyTranslator.java | 215 ++++++++++
.../streams/translation/ImpulseProcessor.java | 145 +++++++
.../streams/translation/ImpulseTranslator.java | 77 ++++
.../kafka/streams/translation/KStreamsPayload.java | 184 ++++++++
.../streams/translation/KStreamsPayloadSerde.java | 123 ++++++
.../KafkaStreamsExecutableStageContextFactory.java | 66 +++
.../KafkaStreamsPipelineTranslator.java | 208 +++++++++
.../translation/KafkaStreamsStateInternals.java | 463 +++++++++++++++++++++
.../translation/KafkaStreamsTimerInternals.java | 258 ++++++++++++
.../KafkaStreamsTranslationContext.java | 179 ++++++++
.../streams/translation/PTransformTranslator.java | 41 ++
.../kafka/streams/translation/ReadProcessor.java | 206 +++++++++
.../kafka/streams/translation/ReadTranslator.java | 272 ++++++++++++
.../translation/RedistributeTranslator.java | 57 +++
.../streams/translation/ShuffleByKeyProcessor.java | 138 ++++++
.../streams/translation/StageOutputProcessor.java | 103 +++++
.../kafka/streams/translation/StoreKeys.java | 105 +++++
.../streams/translation/TerminationReporter.java | 105 +++++
.../streams/translation/TerminationTracker.java | 176 ++++++++
.../translation/UnboundedReadProcessor.java | 338 +++++++++++++++
.../streams/translation/WatermarkAggregator.java | 98 +++++
.../streams/translation/WatermarkManager.java | 136 ++++++
.../streams/translation/WatermarkPayload.java | 49 +++
.../translation/WindowedGroupByKeyProcessor.java | 337 +++++++++++++++
.../kafka/streams/translation/package-info.java | 20 +
.../streams/KafkaStreamsJobServerDriverTest.java | 71 ++++
.../streams/KafkaStreamsPipelineOptionsTest.java | 90 ++++
.../KafkaStreamsPipelineRunnerConfigTest.java | 86 ++++
.../KafkaStreamsPortablePipelineResultTest.java | 96 +++++
.../kafka/streams/KafkaStreamsRunnerBrokerIT.java | 372 +++++++++++++++++
.../kafka/streams/KafkaStreamsRunnerTest.java | 159 +++++++
.../kafka/streams/KafkaStreamsTestRunner.java | 149 +++++++
.../kafka/streams/MultiOutputStageTest.java | 87 ++++
.../kafka/streams/TestKafkaStreamsRunner.java | 151 +++++++
.../kafka/streams/TestKafkaStreamsRunnerTest.java | 82 ++++
.../streams/translation/BundleBoundaryTest.java | 118 ++++++
.../translation/ChainedExecutableStageTest.java | 163 ++++++++
.../kafka/streams/translation/CreateTest.java | 73 ++++
.../ExecutableStageProcessorWatermarkTest.java | 145 +++++++
.../translation/ExecutableStageTranslatorTest.java | 82 ++++
.../translation/FixedWindowGroupByKeyTest.java | 106 +++++
.../translation/FlattenParallelismTest.java | 141 +++++++
.../kafka/streams/translation/FlattenTest.java | 212 ++++++++++
.../kafka/streams/translation/GroupByKeyTest.java | 93 +++++
.../streams/translation/ImpulseTranslatorTest.java | 144 +++++++
.../translation/KStreamsPayloadSerdeTest.java | 97 +++++
.../KafkaStreamsPipelineTranslatorTest.java | 146 +++++++
.../KafkaStreamsTimerInternalsTest.java | 208 +++++++++
.../translation/MetricsAcrossBundlesTest.java | 81 ++++
.../kafka/streams/translation/MetricsTest.java | 79 ++++
.../kafka/streams/translation/ReadTest.java | 75 ++++
.../streams/translation/SharedTestCollector.java | 92 ++++
.../translation/ShuffleByKeyProcessorTest.java | 125 ++++++
.../translation/StageOutputProcessorTest.java | 108 +++++
.../StandardWindowFnTranslationTest.java | 152 +++++++
.../translation/TerminationTrackerTest.java | 172 ++++++++
.../streams/translation/UnboundedReadTest.java | 263 ++++++++++++
.../translation/WatermarkAggregatorTest.java | 160 +++++++
.../streams/translation/WatermarkManagerTest.java | 155 +++++++
.../translation/WatermarkPropagationTest.java | 97 +++++
.../python/apache_beam/options/pipeline_options.py | 27 ++
.../kafka_streams_java_job_server_test.py | 140 +++++++
.../runners/portability/kafka_streams_runner.py | 137 ++++++
.../portability/kafka_streams_runner_test.py | 285 +++++++++++++
sdks/python/test-suites/portable/common.gradle | 39 ++
sdks/python/tox.ini | 5 +
settings.gradle.kts | 12 +
.../en/documentation/runners/kafkastreams.md | 248 +++++++++++
website/www/site/data/capability_matrix.yaml | 155 +++++++
.../layouts/partials/section-menu/en/runners.html | 1 +
101 files changed, 12634 insertions(+), 22 deletions(-)
copy .github/trigger_files/{IO_Iceberg_Integration_Tests_Dataflow.json =>
beam_PreCommit_Java_Kafka_Streams_Runner.json} (100%)
copy .github/workflows/{beam_PreCommit_Go.yml =>
beam_PreCommit_Java_Kafka_Streams_Runner.yml} (56%)
create mode 100644 runners/kafka-streams/build.gradle
create mode 100644 runners/kafka-streams/job-server/build.gradle
create mode 100644 runners/kafka-streams/measurement/build.gradle
create mode 100644 runners/kafka-streams/measurement/docker-compose.yml
create mode 100644
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
create mode 100644
runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java
create mode 100644 runners/kafka-streams/proto/build.gradle
create mode 100644
runners/kafka-streams/proto/src/main/proto/kafka_streams_payload.proto
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobInvoker.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobServerDriver.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineResult.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineRunner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPortablePipelineResult.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerRegistrar.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTopicManager.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/package-info.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/EmptyBoundedSource.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/FlattenTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyBroadcastPartitioner.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ImpulseTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadSerde.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsExecutableStageContextFactory.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsPipelineTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsStateInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternals.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTranslationContext.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/PTransformTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ReadTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/RedistributeTranslator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/StoreKeys.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationReporter.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/TerminationTracker.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/UnboundedReadProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkAggregator.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkManager.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WatermarkPayload.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/WindowedGroupByKeyProcessor.java
create mode 100644
runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/package-info.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsJobServerDriverTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptionsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineRunnerConfigTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPortablePipelineResultTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerBrokerIT.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsTestRunner.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/MultiOutputStageTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/TestKafkaStreamsRunner.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/TestKafkaStreamsRunnerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/BundleBoundaryTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ChainedExecutableStageTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/CreateTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FixedWindowGroupByKeyTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FlattenParallelismTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/FlattenTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/GroupByKeyTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ImpulseTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KStreamsPayloadSerdeTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsPipelineTranslatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/KafkaStreamsTimerInternalsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsAcrossBundlesTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ReadTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/SharedTestCollector.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ShuffleByKeyProcessorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/StageOutputProcessorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/StandardWindowFnTranslationTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/TerminationTrackerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/UnboundedReadTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkAggregatorTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkManagerTest.java
create mode 100644
runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/WatermarkPropagationTest.java
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_java_job_server_test.py
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_runner.py
create mode 100644
sdks/python/apache_beam/runners/portability/kafka_streams_runner_test.py
create mode 100644
website/www/site/content/en/documentation/runners/kafkastreams.md