This is an automated email from the ASF dual-hosted git repository.

je-ik pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git

commit e10e8c969b0f31d5413902072c9bdcaba6b692d4
Merge: 2577a6cff4a 473a5dd649c
Author: Jan Lukavský <[email protected]>
AuthorDate: Sat Aug 22 20:34:26 2026 +0200

    Merge pull request #39785: Merge Kafka Streams Runner skeleton (#18479)

 .../beam_PreCommit_Java_Kafka_Streams_Runner.json  |   4 +
 .github/workflows/README.md                        |   1 +
 .github/workflows/beam_PreCommit_Java.yml          |   1 +
 .../beam_PreCommit_Java_Kafka_Streams_Runner.yml   | 119 ++++++
 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 +
 100 files changed, 12705 insertions(+)


Reply via email to