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 104dc272e8d [GSoC 2026] Kafka Streams runner: ask for primitive reads
in the Java wrapper (#39766)
104dc272e8d is described below
commit 104dc272e8d4e00ac51211dd3bbe5577bb67087e
Author: M Junaid Shaukat <[email protected]>
AuthorDate: Sun Aug 16 22:14:25 2026 +0500
[GSoC 2026] Kafka Streams runner: ask for primitive reads in the Java
wrapper (#39766)
A Read expands into a splittable DoFn by default and the runner cannot
translate one, so a pipeline that merely reads failed to translate unless it
knew to convert the reads itself.
The wrapper now sets use_deprecated_read and converts the pipeline before
handing it on, so a pipeline does not have to know, and a pipeline that asks
for splittable reads still gets primitive ones rather than something that
cannot run.
---
.../runners/kafka/streams/KafkaStreamsRunner.java | 30 ++++++++++-
.../kafka/streams/KafkaStreamsRunnerTest.java | 59 ++++++++++++++++++++++
2 files changed, 88 insertions(+), 1 deletion(-)
diff --git
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
index 1b530fd22ab..924c3ac01fe 100644
---
a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
+++
b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunner.java
@@ -26,6 +26,8 @@ import org.apache.beam.sdk.PipelineRunner;
import org.apache.beam.sdk.options.ExperimentalOptions;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.util.construction.Environments;
+import org.apache.beam.sdk.util.construction.SplittableParDo;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
@@ -60,7 +62,7 @@ public class KafkaStreamsRunner extends
PipelineRunner<PipelineResult> {
@Override
public PipelineResult run(Pipeline pipeline) {
- assignPortableDefaults(pipelineOptions);
+ prepareForTranslation(pipeline, pipelineOptions);
@Nullable KafkaStreamsJobServerDriver jobServerDriver = null;
try {
if (Strings.isNullOrEmpty(pipelineOptions.getJobEndpoint())) {
@@ -89,6 +91,22 @@ public class KafkaStreamsRunner extends
PipelineRunner<PipelineResult> {
}
}
+ /**
+ * Settles the options the runner needs and rewrites the pipeline into what
it can translate.
+ *
+ * <p>The runner does not translate splittable DoFns, and a {@link
org.apache.beam.sdk.io.Read}
+ * expands into one by default, so a pipeline that merely reads would
otherwise fail to translate.
+ * Beam keeps the primitive read for exactly this case, behind an experiment
that {@link
+ * #assignPortableDefaults} sets, so a pipeline does not have to ask for it
and the proto that
+ * reaches the job server already holds primitive reads.
+ */
+ @VisibleForTesting
+ static void prepareForTranslation(
+ Pipeline pipeline, KafkaStreamsPipelineOptions pipelineOptions) {
+ assignPortableDefaults(pipelineOptions);
+
SplittableParDo.convertReadBasedSplittableDoFnsToPrimitiveReadsIfNecessary(pipeline);
+ }
+
private static void assignPortableDefaults(KafkaStreamsPipelineOptions
pipelineOptions) {
if (Strings.isNullOrEmpty(pipelineOptions.getDefaultEnvironmentType())) {
pipelineOptions.setDefaultEnvironmentType(Environments.ENVIRONMENT_LOOPBACK);
@@ -97,8 +115,18 @@ public class KafkaStreamsRunner extends
PipelineRunner<PipelineResult> {
@Nullable List<String> existingExperiments =
experimentalOptions.getExperiments();
List<String> experiments =
existingExperiments == null ? new ArrayList<>() : new
ArrayList<>(existingExperiments);
+ boolean changed = false;
if (!experiments.contains("beam_fn_api")) {
experiments.add("beam_fn_api");
+ changed = true;
+ }
+ // Splittable DoFns are not translated, so the Read that expands into one
has to stay the
+ // primitive it used to be. This is the experiment Beam looks for when
deciding that.
+ if (!experiments.contains("use_deprecated_read")) {
+ experiments.add("use_deprecated_read");
+ changed = true;
+ }
+ if (changed) {
experimentalOptions.setExperiments(experiments);
}
}
diff --git
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
index 02bc1a887c7..321b16f303c 100644
---
a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
+++
b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/KafkaStreamsRunnerTest.java
@@ -17,18 +17,28 @@
*/
package org.apache.beam.runners.kafka.streams;
+import static org.hamcrest.CoreMatchers.hasItem;
import static org.hamcrest.CoreMatchers.is;
+import static org.hamcrest.CoreMatchers.not;
import static org.hamcrest.MatcherAssert.assertThat;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.Set;
+import java.util.stream.Collectors;
import org.apache.beam.model.pipeline.v1.RunnerApi;
import org.apache.beam.runners.kafka.streams.translation.KStreamsPayload;
import org.apache.beam.sdk.Pipeline;
+import org.apache.beam.sdk.io.CountingSource;
+import org.apache.beam.sdk.io.Read;
+import org.apache.beam.sdk.options.ExperimentalOptions;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.Impulse;
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
+import org.apache.beam.sdk.util.construction.PTransformTranslation;
+import org.apache.beam.sdk.util.construction.PipelineTranslation;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.TopologyTestDriver;
import org.apache.kafka.streams.processor.api.Processor;
@@ -50,6 +60,55 @@ import org.junit.Test;
*/
public class KafkaStreamsRunnerTest {
+ /** The transform urns the pipeline would hand to the job server. */
+ private static Set<String> translatedUrns(Pipeline pipeline) {
+ return
PipelineTranslation.toProto(pipeline).getComponents().getTransformsMap().values()
+ .stream()
+ .map(transform -> transform.getSpec().getUrn())
+ .collect(Collectors.toSet());
+ }
+
+ /**
+ * A {@code Read} expands into a splittable DoFn by default, which this
runner cannot translate.
+ * The runner asks for the primitive read instead, so that a pipeline does
not have to know to.
+ */
+ @Test
+ public void
aReadReachesTheJobServerAsAPrimitiveReadRatherThanASplittableDoFn() {
+ KafkaStreamsPipelineOptions options =
+ PipelineOptionsFactory.create().as(KafkaStreamsPipelineOptions.class);
+ options.setApplicationId("read-conversion-test");
+ // Pipeline.create insists on a runner; it does not run here, the pipeline
is only translated.
+ options.setRunner(KafkaStreamsRunner.class);
+ Pipeline pipeline = Pipeline.create(options);
+ pipeline.apply("read", Read.from(CountingSource.unbounded()));
+
+ // Left alone, the read is a splittable DoFn expansion and no primitive
read is present.
+ assertThat(translatedUrns(pipeline),
not(hasItem(PTransformTranslation.READ_TRANSFORM_URN)));
+
+ KafkaStreamsRunner.prepareForTranslation(pipeline, options);
+
+ assertThat(translatedUrns(pipeline),
hasItem(PTransformTranslation.READ_TRANSFORM_URN));
+ }
+
+ /**
+ * A pipeline may ask for splittable reads outright. The runner cannot
translate them, so it asks
+ * for the primitive read anyway rather than letting a pipeline choose
something that cannot run.
+ */
+ @Test
+ public void aPipelineAskingForSplittableReadsStillGetsPrimitiveOnes() {
+ KafkaStreamsPipelineOptions options =
+ PipelineOptionsFactory.create().as(KafkaStreamsPipelineOptions.class);
+ options.setApplicationId("sdf-read-override-test");
+ options.setRunner(KafkaStreamsRunner.class);
+ options.as(ExperimentalOptions.class).setExperiments(new
ArrayList<>(List.of("use_sdf_read")));
+ Pipeline pipeline = Pipeline.create(options);
+ pipeline.apply("read", Read.from(CountingSource.unbounded()));
+
+ KafkaStreamsRunner.prepareForTranslation(pipeline, options);
+
+ assertThat(translatedUrns(pipeline),
hasItem(PTransformTranslation.READ_TRANSFORM_URN));
+ }
+
@Test
public void impulseOnlyPipelineEmitsDataAndTerminalWatermark() {
Pipeline pipeline = Pipeline.create(KafkaStreamsTestRunner.testOptions());