This is an automated email from the ASF dual-hosted git repository.
je-ik pushed a commit to branch feat/18479-kafka-streams-runner-skeleton
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to
refs/heads/feat/18479-kafka-streams-runner-skeleton by this push:
new 10ff557d616 [GSoC 2026] Kafka Streams runner: an application for
measuring instances coming and going (#39752)
10ff557d616 is described below
commit 10ff557d616bcc108bbf7eb659528ae2c20d12b2
Author: M Junaid Shaukat <[email protected]>
AuthorDate: Sun Aug 16 23:17:06 2026 +0500
[GSoC 2026] Kafka Streams runner: an application for measuring instances
coming and going (#39752)
* [GSoC 2026] Kafka Streams runner: an application for measuring instances
coming and going
* [GSoC 2026] Kafka Streams runner: count groups in the pipeline rather
than beside it
SpotBugs is turned off for this module. It runs the pipeline in process, so
the SDK harness and its dependencies are on the classpath and SpotBugs
reports on those instead of on the four classes here. The it/ modules do the
same for the same reason.
---
runners/kafka-streams/measurement/build.gradle | 78 ++++++
.../kafka-streams/measurement/docker-compose.yml | 27 +++
.../streams/measurement/RescalingMeasurement.java | 261 +++++++++++++++++++++
.../kafka/streams/measurement/package-info.java | 26 ++
settings.gradle.kts | 1 +
5 files changed, 393 insertions(+)
diff --git a/runners/kafka-streams/measurement/build.gradle
b/runners/kafka-streams/measurement/build.gradle
new file mode 100644
index 00000000000..5e7992a0916
--- /dev/null
+++ b/runners/kafka-streams/measurement/build.gradle
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * License); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an AS IS BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+/**
+ * An application for measuring the Kafka Streams runner's behaviour when
instances come and go.
+ *
+ * Not part of the build's verification: it is something a person runs against
a Kafka, several
+ * copies at once, and watches.
+ */
+
+apply plugin: 'org.apache.beam.module'
+apply plugin: 'application'
+mainClassName =
"org.apache.beam.runners.kafka.streams.measurement.RescalingMeasurement"
+
+applyJavaNature(
+ automaticModuleName: 'org.apache.beam.runners.kafka.streams.measurement',
+ publish: false,
+ exportJavadoc: false,
+ // This module runs the pipeline in its own process, so the SDK harness and
its dependencies are
+ // on the classpath, and SpotBugs reports on those rather than on the four
classes here — some
+ // eleven thousand warnings, none of them in this source tree. The same is
done in the it/
+ // modules, which are on the classpath of what they exercise for the same
reason. Checkstyle,
+ // ErrorProne, spotless and the nullness checker all still run.
+ enableSpotbugs: false,
+)
+
+description = "Apache Beam :: Runners :: Kafka Streams :: Measurement"
+
+def kafkaStreamsRunnerProject = ":runners:kafka-streams"
+
+evaluationDependsOn(kafkaStreamsRunnerProject)
+
+// Same pin as the runner and the job server: applyJavaNature forces the
versions in library.java,
+// which includes an older kafka-clients than the runner is compiled against.
+def kafka_version = project(kafkaStreamsRunnerProject).kafka_version
+
+configurations.configureEach {
+ resolutionStrategy.eachDependency { details ->
+ if (details.requested.group == "org.apache.kafka") {
+ details.useVersion(kafka_version)
+ details.because("Kafka Streams runner is developed against Kafka
${kafka_version}.")
+ }
+ }
+}
+
+dependencies {
+ implementation project(kafkaStreamsRunnerProject)
+ implementation project(path: ":sdks:java:core", configuration: "shadow")
+ implementation project(path: ":model:pipeline", configuration: "shadow")
+ implementation project(":runners:java-fn-execution")
+ // On the compile classpath to resolve PortablePipelineRunner, which
KafkaStreamsPipelineRunner
+ // implements; no class of it is named here, so the dependency analysis does
not see it used.
+ implementation project(":runners:java-job-service")
+ permitUnusedDeclared project(":runners:java-job-service")
+ implementation project(":runners:core-java")
+ permitUnusedDeclared project(":runners:core-java")
+ // The pipeline's own code runs in this process, so the Java SDK harness has
to be present.
+ runtimeOnly project(":sdks:java:harness")
+ implementation library.java.joda_time
+ // Without a binding the application starts and says nothing, which is
unhelpful for something
+ // whose whole purpose is to be watched while it runs.
+ runtimeOnly library.java.slf4j_simple
+}
diff --git a/runners/kafka-streams/measurement/docker-compose.yml
b/runners/kafka-streams/measurement/docker-compose.yml
new file mode 100644
index 00000000000..836f121ca33
--- /dev/null
+++ b/runners/kafka-streams/measurement/docker-compose.yml
@@ -0,0 +1,27 @@
+# One Kafka for the measurement application. One broker is enough: what gets
run several times is
+# the runner instance, not the broker.
+#
+# docker compose -f runners/kafka-streams/measurement/docker-compose.yml up
-d
+#
+# group.min.session.timeout.ms is lowered because a broker refuses a session
timeout below it, and
+# how quickly the group notices a departed instance is the floor on how
quickly its work moves. The
+# default of 6s would put a floor under every measurement of recovery.
+services:
+ kafka:
+ image: apache/kafka:4.0.0
+ container_name: ks-measurement-kafka
+ ports:
+ - "9092:9092"
+ environment:
+ KAFKA_NODE_ID: 1
+ KAFKA_PROCESS_ROLES: broker,controller
+ KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
+ KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
+ KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
+ KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
+ KAFKA_LISTENER_SECURITY_PROTOCOL_MAP:
CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
+ KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
+ KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
+ KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
+ KAFKA_GROUP_MIN_SESSION_TIMEOUT_MS: 1000
+ KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
diff --git
a/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
b/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
new file mode 100644
index 00000000000..3ba9b22ca0e
--- /dev/null
+++
b/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/RescalingMeasurement.java
@@ -0,0 +1,261 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.runners.kafka.streams.measurement;
+
+import org.apache.beam.model.pipeline.v1.RunnerApi;
+import org.apache.beam.runners.fnexecution.provisioning.JobInfo;
+import org.apache.beam.runners.kafka.streams.KafkaStreamsPipelineOptions;
+import org.apache.beam.runners.kafka.streams.KafkaStreamsPipelineRunner;
+import org.apache.beam.runners.kafka.streams.KafkaStreamsRunner;
+import org.apache.beam.sdk.Pipeline;
+import org.apache.beam.sdk.io.GenerateSequence;
+import org.apache.beam.sdk.options.Default;
+import org.apache.beam.sdk.options.Description;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.apache.beam.sdk.options.PortablePipelineOptions;
+import org.apache.beam.sdk.transforms.Count;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.transforms.MapElements;
+import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.transforms.windowing.FixedWindows;
+import org.apache.beam.sdk.transforms.windowing.Window;
+import org.apache.beam.sdk.util.construction.Environments;
+import org.apache.beam.sdk.util.construction.PipelineOptionsTranslation;
+import org.apache.beam.sdk.util.construction.PipelineTranslation;
+import org.apache.beam.sdk.util.construction.SplittableParDo;
+import org.apache.beam.sdk.values.KV;
+import org.apache.beam.sdk.values.TypeDescriptors;
+import org.joda.time.Duration;
+
+/**
+ * One instance of a streaming pipeline, run as an ordinary application, for
measuring what happens
+ * when instances are added and removed.
+ *
+ * <p>Run several of these against one Kafka. They share an application id, so
Kafka's consumer
+ * group divides the work between them, and stopping one hands its share to
the others.
+ *
+ * <p>This is an application rather than a test on purpose. The numbers only
mean something if the
+ * pipeline is doing a realistic amount of work — a grouping over thousands of
keys, fed fast enough
+ * that every partition has something to do. A pipeline that trickles produces
idle partitions, and
+ * an idle partition holds a watermark back for reasons that have nothing to
do with rescaling.
+ *
+ * <p>The source produces a fixed number of elements per second over a fixed
set of keys, so what a
+ * complete window looks like is known before the run starts: every window
should report the same
+ * number of groups. That is what makes a shortfall legible as a shortfall,
rather than as one of
+ * the many rates a pipeline could happen to be running at.
+ *
+ * <pre>
+ * docker compose -f runners/kafka-streams/measurement/docker-compose.yml up
-d
+ * ./gradlew :runners:kafka-streams:measurement:installDist
+ * </pre>
+ *
+ * <p>Then start two instances, sharing an application id and differing in
everything local to the
+ * instance. Each needs its own {@code --stateDir}: two instances sharing one
directory fail with a
+ * {@code LockException}, because Kafka Streams locks the state it keeps on
disk.
+ *
+ * <pre>
+ *
BIN=runners/kafka-streams/measurement/build/install/measurement/bin/measurement
+ * $BIN --applicationId=demo --instanceName=one --stateDir=/tmp/ks-one &
+ * $BIN --applicationId=demo --instanceName=two --stateDir=/tmp/ks-two &
+ * </pre>
+ *
+ * <p>The pipeline logs one line per key per window. Nothing is counted beside
the pipeline: the
+ * groups in a window are its own output, so the tally does not depend on how
many instances are
+ * running or on which of them happens to be doing the work.
+ *
+ * <pre>
+ * <millis> <instance> window_end=<millis> key=<key>
count=<n> skew_ms=<n>
+ * </pre>
+ *
+ * <p>Because the rate and the key space are both fixed, a complete window has
one line per key and
+ * the same count on each, so counting the lines for a window says whether the
window was complete.
+ *
+ * <p>{@code skew_ms} is the gap between the window's event time and the wall
clock when the group
+ * came out. It is what falling behind should look like: a pipeline that
cannot keep up ought to
+ * report its groups later and later while still reporting all of them, so a
climbing skew with
+ * complete windows is congestion, and missing groups are something else.
+ *
+ * <p>To watch a handover, kill one instance and watch the other's lines. The
delay before the
+ * survivor reports the killed instance's share again is dominated by {@code
--sessionTimeoutMs},
+ * which is how long the consumer group waits before deciding the instance is
gone.
+ */
+public final class RescalingMeasurement {
+
+ private RescalingMeasurement() {}
+
+ /** Whether the command line mentioned an option, so that a default is not
applied over it. */
+ private static boolean given(String[] args, String name) {
+ for (String arg : args) {
+ if (arg.equals("--" + name) || arg.startsWith("--" + name + "=")) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ /**
+ * Applies the defaults this measurement needs, where they differ from the
runner's own.
+ *
+ * <p>The runner's defaults are meant for a pipeline, not for this. Left
alone, the source is read
+ * in large enough turns that the read starves the rest of the topology and
no groups come out at
+ * all. The parallelism is raised for a related reason: a measurement of
work moving between
+ * instances needs more than the single partition the runner defaults to,
since with one partition
+ * there is nothing to divide.
+ */
+ private static void applyMeasurementDefaults(String[] args,
MeasurementOptions options) {
+ if (!given(args, "readMaxElementsPerPoll")) {
+ options.setReadMaxElementsPerPoll(200);
+ }
+ if (!given(args, "internalParallelism")) {
+ options.setInternalParallelism(3);
+ }
+ }
+
+ /** Options of the measurement itself, on top of the runner's own. */
+ public interface MeasurementOptions extends KafkaStreamsPipelineOptions {
+
+ @Description("Name for this instance in the output, so several can be told
apart.")
+ @Default.String("instance")
+ String getInstanceName();
+
+ void setInstanceName(String instanceName);
+
+ @Description(
+ "How many distinct keys the grouping runs over. Thousands, so that
every partition of the"
+ + " shuffle has work and no partition sits idle holding a
watermark back. With a"
+ + " window long enough to contain them all, this is also how many
groups a complete"
+ + " window has.")
+ @Default.Integer(2_000)
+ int getNumKeys();
+
+ void setNumKeys(int numKeys);
+
+ @Description(
+ "How many elements the source produces per second. Fixed rather than
as-fast-as-possible so"
+ + " that a window's contents are known in advance and a shortfall
is visible.")
+ @Default.Integer(20_000)
+ int getElementsPerSecond();
+
+ void setElementsPerSecond(int elementsPerSecond);
+
+ @Description("Window size in milliseconds; how often the groups are
counted and reported.")
+ @Default.Integer(1_000)
+ int getWindowMs();
+
+ void setWindowMs(int windowMs);
+ }
+
+ /**
+ * Logs each group the pipeline produces, with how far behind the wall clock
its window was.
+ *
+ * <p>One line per key per window. With a fixed rate over a fixed key space
every window holds the
+ * same groups, so counting the lines for a window says whether the window
was complete, and no
+ * counter has to be kept anywhere for that to be true — the count is the
pipeline's own output
+ * rather than a tally maintained beside it, which is what makes it
independent of how many
+ * instances are running.
+ *
+ * <p>The skew is the point of the timestamp. A pipeline that cannot keep up
should report its
+ * groups later and later rather than stop reporting them, so a skew that
climbs while the groups
+ * stay complete is the pipeline falling behind, and groups going missing is
something else.
+ */
+ private static class ReportGroupFn extends DoFn<KV<String, Long>, Void> {
+ private final String instanceName;
+
+ ReportGroupFn(String instanceName) {
+ this.instanceName = instanceName;
+ }
+
+ @ProcessElement
+ public void processElement(@Element KV<String, Long> group, BoundedWindow
window) {
+ long windowEnd = window.maxTimestamp().getMillis();
+ long now = System.currentTimeMillis();
+ System.out.printf(
+ "%d %s window_end=%d key=%s count=%d skew_ms=%d%n",
+ now, instanceName, windowEnd, group.getKey(), group.getValue(), now
- windowEnd);
+ }
+ }
+
+ public static void main(String[] args) throws Exception {
+ PipelineOptionsFactory.register(MeasurementOptions.class);
+ // Deliberately not withValidation(): that enforces the options a pipeline
needs when it is
+ // submitted to a job server, and --jobEndpoint above all, which means
nothing here because this
+ // application runs the pipeline itself.
+ MeasurementOptions options =
PipelineOptionsFactory.fromArgs(args).as(MeasurementOptions.class);
+ if (options.getApplicationId() == null ||
options.getApplicationId().isEmpty()) {
+ throw new IllegalArgumentException(
+ "--applicationId is required, and every instance of one measurement
must share it: it is"
+ + " what puts them in the same consumer group and so divides the
work between them.");
+ }
+ applyMeasurementDefaults(args, options);
+ // Pipeline.create needs a runner class even though this application never
calls pipeline.run()
+ // — it builds the pipeline proto and hands it to the runner below itself.
+ options.setRunner(KafkaStreamsRunner.class);
+ // The user code runs in this same process, so no container or separate
worker is needed.
+ options
+ .as(PortablePipelineOptions.class)
+ .setDefaultEnvironmentType(Environments.ENVIRONMENT_EMBEDDED);
+
+ // A window holds every key as long as it is long enough for the rate to
reach them all; below
+ // that the source has not got round to each key once and the window is
short by construction.
+ long elementsPerWindow = (long) options.getElementsPerSecond() *
options.getWindowMs() / 1_000L;
+ long expectedGroups = Math.min(options.getNumKeys(), elementsPerWindow);
+
+ int numKeys = options.getNumKeys();
+ Pipeline pipeline = Pipeline.create(options);
+ pipeline
+ .apply(
+ "read",
+ GenerateSequence.from(0)
+ .withRate(options.getElementsPerSecond(),
Duration.standardSeconds(1)))
+ .apply(
+ "key",
+ // numKeys is read here rather than inside the lambda: reaching
for it through options
+ // would capture the PipelineOptions in the transform, which
cannot be serialized.
+ MapElements.into(TypeDescriptors.strings()).via((Long n) -> "key-"
+ (n % numKeys)))
+ .apply("window",
Window.into(FixedWindows.of(Duration.millis(options.getWindowMs()))))
+ .apply("countPerKey", Count.perElement())
+ .apply("report", ParDo.of(new
ReportGroupFn(options.getInstanceName())));
+
+ SplittableParDo.convertReadBasedSplittableDoFnsToPrimitiveReads(pipeline);
+ RunnerApi.Pipeline proto = PipelineTranslation.toProto(pipeline);
+ JobInfo jobInfo =
+ JobInfo.create(
+ options.getApplicationId(),
+ options.getJobName(),
+ "",
+ PipelineOptionsTranslation.toProto(options));
+
+ System.out.printf(
+ "starting %s: application=%s keys=%d rate=%d/s parallelism=%d
window=%dms"
+ + " session_timeout=%dms read_per_poll=%d bundle=%d
expected_groups_per_window=%d%n",
+ options.getInstanceName(),
+ options.getApplicationId(),
+ options.getNumKeys(),
+ options.getElementsPerSecond(),
+ options.getInternalParallelism(),
+ options.getWindowMs(),
+ options.getSessionTimeoutMs(),
+ options.getReadMaxElementsPerPoll(),
+ options.getMaxBundleSize(),
+ expectedGroups);
+
+ // Blocks until the instance is stopped; a streaming pipeline has no end
of its own.
+ new KafkaStreamsPipelineRunner(options).run(proto, jobInfo);
+ }
+}
diff --git
a/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java
b/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java
new file mode 100644
index 00000000000..41579676fb8
--- /dev/null
+++
b/runners/kafka-streams/measurement/src/main/java/org/apache/beam/runners/kafka/streams/measurement/package-info.java
@@ -0,0 +1,26 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+/**
+ * An application for measuring the Kafka Streams runner's behaviour when
instances come and go.
+ *
+ * <p>Not part of the build's verification: it is something a person runs
against a Kafka, several
+ * copies at once, and watches. See {@link
+ * org.apache.beam.runners.kafka.streams.measurement.RescalingMeasurement} for
how to run it.
+ */
+package org.apache.beam.runners.kafka.streams.measurement;
diff --git a/settings.gradle.kts b/settings.gradle.kts
index 050d97dc600..de1e5ea6533 100644
--- a/settings.gradle.kts
+++ b/settings.gradle.kts
@@ -153,6 +153,7 @@ if
(startParameter.projectProperties.containsKey("with-kafka-streams-runner")) {
include(":runners:kafka-streams")
include(":runners:kafka-streams:proto")
include(":runners:kafka-streams:job-server")
+ include(":runners:kafka-streams:measurement")
}
include(":runners:local-java")
include(":runners:portability:java")