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());

Reply via email to