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

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


The following commit(s) were added to refs/heads/master by this push:
     new d7557af6a55 [Spark 4] Add streaming dispatch seam and lifecycle hooks 
(#39906)
d7557af6a55 is described below

commit d7557af6a557c927168ac9dc21c64bfd44b76536
Author: Tobias Kaymak <[email protected]>
AuthorDate: Fri Aug 28 03:53:03 2026 +0200

    [Spark 4] Add streaming dispatch seam and lifecycle hooks (#39906)
    
    Prepares the structured streaming runner for a Spark 4 streaming
    translator without changing batch behavior.
    
    - Add PipelineTranslatorFactory, the seam the Spark 4 module shadows to
      dispatch streaming pipelines. The shared base rejects streaming with a
      clear message instead of the previous generic checkArgument.
    - Open up EvaluationContext (non final, protected ctor, leaves()) and add
      a no-op stop() so a streaming context can override evaluation.
    - Add a createEvaluationContext hook to PipelineTranslator and skip the
      persist and lineage breaking optimizations for streaming datasets.
    - Plumb the EvaluationContext into SparkStructuredStreamingPipelineResult
      so cancel() can stop a running streaming query.
    - Default the state store provider to RocksDB, required by Spark 4
      transformWithState and inert for batch.
    - Add watermarkDelayMillis, maxRecordsPerMicroBatch, maxBatchDurationMillis
      and streamingStopAfterIdleBatches options.
    - Narrow the Spark 4 test source override exclude to the legacy DStream
      package. The previous glob also matched structuredstreaming and would
      have silently dropped the structured streaming tests.
---
 runners/spark/4/build.gradle                       |  8 +++--
 .../SparkStructuredStreamingPipelineOptions.java   | 34 ++++++++++++++++++
 .../SparkStructuredStreamingPipelineResult.java    | 16 +++++++--
 .../SparkStructuredStreamingRunner.java            | 26 +++++++++-----
 .../translation/EvaluationContext.java             | 18 ++++++++--
 .../translation/PipelineTranslator.java            | 22 ++++++++++--
 .../translation/PipelineTranslatorFactory.java     | 41 ++++++++++++++++++++++
 .../translation/SparkSessionFactory.java           |  6 ++++
 8 files changed, 150 insertions(+), 21 deletions(-)

diff --git a/runners/spark/4/build.gradle b/runners/spark/4/build.gradle
index f2746b06158..ae8750ec5b1 100644
--- a/runners/spark/4/build.gradle
+++ b/runners/spark/4/build.gradle
@@ -53,12 +53,14 @@ tasks.validatesStructuredStreamingRunnerBatch {
 
 tasks.validatesRunner.dependsOn(validatesStructuredStreamingRunnerBatch)
 
-// Exclude DStream-based streaming tests from the shared-base copy: the Spark 
4 module
-// supports only structured streaming (batch) and does not include legacy 
DStream support.
+// Exclude legacy DStream-based streaming tests (under 
runners/spark/translation/streaming)
+// from the shared-base copy: the Spark 4 module does not include legacy 
DStream support.
 // Streaming test utilities also depend on kafka.server.KafkaServerStartable 
which was
 // removed in Kafka 2.8.0 (the first Kafka version with a _2.13 artifact).
+// Note: structured streaming tests 
(**/structuredstreaming/translation/streaming/**) are
+// intentionally NOT excluded.
 tasks.named("copyTestSourceOverrides") {
-  exclude "**/translation/streaming/**"
+  exclude "**/runners/spark/translation/streaming/**"
 }
 
 // Spark 4 uses org.lz4:lz4-java instead of at.yawk.lz4:lz4-java
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
index 3371a403b2c..29cc4cb99cf 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
@@ -39,4 +39,38 @@ public interface SparkStructuredStreamingPipelineOptions 
extends SparkCommonPipe
   boolean getUseActiveSparkSession();
 
   void setUseActiveSparkSession(boolean value);
+
+  @Description(
+      "Watermark delay in milliseconds applied to event timestamps of 
streaming sources "
+          + "(streaming mode only).")
+  @Default.Long(0)
+  long getWatermarkDelayMillis();
+
+  void setWatermarkDelayMillis(long value);
+
+  // Note: deliberately NOT named getMaxRecordsPerBatch. The legacy Spark 
runner's
+  // SparkPipelineOptions already declares Long getMaxRecordsPerBatch(); a 
same-name getter with a
+  // different return type breaks proxy generation for every registered 
PipelineOptions interface.
+  @Description(
+      "Maximum number of records to read per micro-batch from a streaming 
source "
+          + "(streaming mode only).")
+  @Default.Integer(1000)
+  int getMaxRecordsPerMicroBatch();
+
+  void setMaxRecordsPerMicroBatch(int value);
+
+  @Description(
+      "Maximum duration in milliseconds of a micro-batch trigger interval 
(streaming mode only).")
+  @Default.Long(500)
+  long getMaxBatchDurationMillis();
+
+  void setMaxBatchDurationMillis(long value);
+
+  @Description(
+      "Test-oriented: gracefully stop streaming queries after this many 
consecutive empty "
+          + "micro-batches. Disabled if negative (streaming mode only).")
+  @Default.Integer(-1)
+  int getStreamingStopAfterIdleBatches();
+
+  void setStreamingStopAfterIdleBatches(int value);
 }
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java
index 9d3419e1947..b592b6fb742 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java
@@ -25,27 +25,33 @@ import java.util.concurrent.ExecutionException;
 import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
-import javax.annotation.Nullable;
+import java.util.function.Supplier;
 import 
org.apache.beam.runners.spark.structuredstreaming.metrics.MetricsAccumulator;
+import 
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.PipelineResult;
 import org.apache.beam.sdk.metrics.MetricResults;
 import org.apache.beam.sdk.util.UserCodeException;
 import org.apache.spark.SparkException;
+import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Duration;
 
 public class SparkStructuredStreamingPipelineResult implements PipelineResult {
 
   private final Future<?> pipelineExecution;
+  // Supplies the context of the translated pipeline, null until translation 
has completed.
+  private final Supplier<? extends @Nullable EvaluationContext> 
evaluationContext;
   private final MetricsAccumulator metrics;
-  private @Nullable final Runnable onTerminalState;
+  private final @Nullable Runnable onTerminalState;
   private PipelineResult.State state;
 
   SparkStructuredStreamingPipelineResult(
       Future<?> pipelineExecution,
+      Supplier<? extends @Nullable EvaluationContext> evaluationContext,
       MetricsAccumulator metrics,
-      @Nullable final Runnable onTerminalState) {
+      final @Nullable Runnable onTerminalState) {
     this.pipelineExecution = pipelineExecution;
+    this.evaluationContext = evaluationContext;
     this.metrics = metrics;
     this.onTerminalState = onTerminalState;
     // pipelineExecution is expected to have started executing eagerly.
@@ -113,6 +119,10 @@ public class SparkStructuredStreamingPipelineResult 
implements PipelineResult {
 
   @Override
   public PipelineResult.State cancel() throws IOException {
+    EvaluationContext ctx = evaluationContext.get();
+    if (ctx != null) {
+      ctx.stop();
+    }
     pipelineExecution.cancel(true);
     offerNewState(PipelineResult.State.CANCELLED);
     return state;
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java
index 96717f29e87..f78026847fa 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java
@@ -17,12 +17,11 @@
  */
 package org.apache.beam.runners.spark.structuredstreaming;
 
-import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
-
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
 import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.atomic.AtomicReference;
 import javax.annotation.Nullable;
 import org.apache.beam.runners.core.metrics.MetricsPusher;
 import org.apache.beam.runners.core.metrics.NoOpMetricsSink;
@@ -30,8 +29,8 @@ import 
org.apache.beam.runners.spark.structuredstreaming.metrics.MetricsAccumula
 import 
org.apache.beam.runners.spark.structuredstreaming.metrics.SparkBeamMetricSource;
 import 
org.apache.beam.runners.spark.structuredstreaming.translation.EvaluationContext;
 import 
org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslator;
+import 
org.apache.beam.runners.spark.structuredstreaming.translation.PipelineTranslatorFactory;
 import 
org.apache.beam.runners.spark.structuredstreaming.translation.SparkSessionFactory;
-import 
org.apache.beam.runners.spark.structuredstreaming.translation.batch.PipelineTranslatorBatch;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.PipelineRunner;
 import org.apache.beam.sdk.metrics.MetricsEnvironment;
@@ -54,9 +53,9 @@ import org.slf4j.LoggerFactory;
  * 
href="https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html";>Structured
  * Streaming framework</a>).
  *
- * <p><b>This runner is experimental, its coverage of the Beam model is still 
partial. Due to
- * limitations of the Structured Streaming framework (e.g. lack of support for 
multiple stateful
- * operators), streaming mode is not yet supported by this runner. </b>
+ * <p><b>This runner is experimental, its coverage of the Beam model is still 
partial. Streaming
+ * mode requires the Spark 4 module (beam-runners-spark-4); the shared Spark 3 
module supports batch
+ * pipelines only. </b>
  *
  * <p>The runner translates transforms defined on a Beam pipeline to Spark 
`Dataset` transformations
  * (leveraging the high level Dataset API) and then submits these to Spark to 
be executed.
@@ -145,17 +144,26 @@ public final class SparkStructuredStreamingRunner
             + " It is still experimental, its coverage of the Beam model is 
partial. ***");
 
     PipelineTranslator.detectStreamingMode(pipeline, options);
-    checkArgument(!options.isStreaming(), "Streaming is not supported.");
 
     final SparkSession sparkSession = 
SparkSessionFactory.getOrCreateSession(options);
     final MetricsAccumulator metrics = 
MetricsAccumulator.getInstance(sparkSession);
 
+    // Set once the pipeline is translated, so the result can stop an ongoing 
(streaming)
+    // evaluation on cancel. Remains null until translation completes.
+    final AtomicReference<EvaluationContext> ctxRef = new AtomicReference<>();
+
     final Future<?> submissionFuture =
-        runAsync(() -> translatePipeline(sparkSession, pipeline).evaluate());
+        runAsync(
+            () -> {
+              EvaluationContext ctx = translatePipeline(sparkSession, 
pipeline);
+              ctxRef.set(ctx);
+              ctx.evaluate();
+            });
 
     final SparkStructuredStreamingPipelineResult result =
         new SparkStructuredStreamingPipelineResult(
             submissionFuture,
+            ctxRef::get,
             metrics,
             sparkStopFn(sparkSession, options.getUseActiveSparkSession()));
 
@@ -186,7 +194,7 @@ public final class SparkStructuredStreamingRunner
 
     PipelineTranslator.replaceTransforms(pipeline, options);
 
-    PipelineTranslator pipelineTranslator = new PipelineTranslatorBatch();
+    PipelineTranslator pipelineTranslator = 
PipelineTranslatorFactory.create(options.isStreaming());
     return pipelineTranslator.translate(pipeline, sparkSession, options);
   }
 
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/EvaluationContext.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/EvaluationContext.java
index b8448567eaf..0e677051fb6 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/EvaluationContext.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/EvaluationContext.java
@@ -40,10 +40,10 @@ import org.slf4j.LoggerFactory;
  */
 @SuppressWarnings("Slf4jDoNotLogMessageOfExceptionExplicitly")
 @Internal
-public final class EvaluationContext {
+public class EvaluationContext {
   private static final Logger LOG = 
LoggerFactory.getLogger(EvaluationContext.class);
 
-  interface NamedDataset<T> {
+  public interface NamedDataset<T> {
     String name();
 
     @Nullable
@@ -53,11 +53,16 @@ public final class EvaluationContext {
   private final Collection<? extends NamedDataset<?>> leaves;
   private final SparkSession session;
 
-  EvaluationContext(Collection<? extends NamedDataset<?>> leaves, SparkSession 
session) {
+  protected EvaluationContext(Collection<? extends NamedDataset<?>> leaves, 
SparkSession session) {
     this.leaves = leaves;
     this.session = session;
   }
 
+  /** The leaf datasets of the translated pipeline that require evaluation. */
+  protected Collection<? extends NamedDataset<?>> leaves() {
+    return leaves;
+  }
+
   /** Trigger evaluation of all leaf datasets. */
   public void evaluate() {
     for (NamedDataset<?> ds : leaves) {
@@ -113,6 +118,13 @@ public final class EvaluationContext {
     }
   }
 
+  /**
+   * Stops any ongoing streaming execution triggered by this context.
+   *
+   * <p>This is a no-op for batch pipelines.
+   */
+  public void stop() {}
+
   public SparkSession getSparkSession() {
     return session;
   }
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java
index a681dea2fde..8470c968550 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java
@@ -25,6 +25,7 @@ import static 
org.apache.beam.sdk.values.PCollection.IsBounded.UNBOUNDED;
 import java.io.IOException;
 import java.io.Serializable;
 import java.util.ArrayList;
+import java.util.Collection;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -128,7 +129,20 @@ public abstract class PipelineTranslator {
     TranslatingVisitor translator = new TranslatingVisitor(session, options, 
dependencies.results);
     pipeline.traverseTopologically(translator);
 
-    return new EvaluationContext(translator.leaves, session);
+    return createEvaluationContext(translator.leaves, session, options);
+  }
+
+  /**
+   * Creates the {@link EvaluationContext} for the translated pipeline.
+   *
+   * <p>Subclasses may override this to return a specialized context, e.g. to 
evaluate streaming
+   * pipelines.
+   */
+  protected EvaluationContext createEvaluationContext(
+      Collection<? extends EvaluationContext.NamedDataset<?>> leaves,
+      SparkSession session,
+      SparkCommonPipelineOptions options) {
+    return new EvaluationContext(leaves, session);
   }
 
   /**
@@ -311,12 +325,14 @@ public abstract class PipelineTranslator {
       TranslationResult<?, T> result = getResult(pCollection);
       result.dataset = dataset;
 
-      if (cache && result.usages() > 1) {
+      // Caching and lineage breaking are batch-only optimizations, streaming 
datasets must pass
+      // through untouched.
+      if (cache && result.usages() > 1 && !dataset.isStreaming()) {
         LOG.info("Dataset {} will be cached for reuse.", result.name);
         dataset.persist(storageLevel); // use NONE to disable
       }
 
-      if (result.estimatePlanComplexity() > PLAN_COMPLEXITY_THRESHOLD) {
+      if (!dataset.isStreaming() && result.estimatePlanComplexity() > 
PLAN_COMPLEXITY_THRESHOLD) {
         // Break linage of dataset to limit planning overhead for complex 
query plans.
         LOG.info("Breaking linage of dataset {} to limit complexity of query 
plan.", result.name);
         result.dataset = sparkSession.createDataset(dataset.rdd(), 
dataset.encoder());
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslatorFactory.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslatorFactory.java
new file mode 100644
index 00000000000..4609d1d380e
--- /dev/null
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslatorFactory.java
@@ -0,0 +1,41 @@
+/*
+ * 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.spark.structuredstreaming.translation;
+
+import 
org.apache.beam.runners.spark.structuredstreaming.translation.batch.PipelineTranslatorBatch;
+import org.apache.beam.sdk.annotations.Internal;
+
+/**
+ * Factory to create the {@link PipelineTranslator} matching the execution 
mode of the pipeline.
+ *
+ * <p>This shared base version only supports batch mode. The Spark 4 module 
shadows this file to
+ * additionally dispatch to a streaming translator.
+ */
+@Internal
+public final class PipelineTranslatorFactory {
+  private PipelineTranslatorFactory() {}
+
+  /** Creates a {@link PipelineTranslator} for the given execution mode. */
+  public static PipelineTranslator create(boolean streaming) {
+    if (streaming) {
+      throw new UnsupportedOperationException(
+          "Streaming pipelines require the Spark 4 runner 
(beam-runners-spark-4).");
+    }
+    return new PipelineTranslatorBatch();
+  }
+}
diff --git 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkSessionFactory.java
 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkSessionFactory.java
index 5b9e5b6fae8..822d1871b12 100644
--- 
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkSessionFactory.java
+++ 
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkSessionFactory.java
@@ -175,6 +175,12 @@ public class SparkSessionFactory {
       sparkConf.setIfMissing("spark.sql.shuffle.partitions", 
Integer.toString(partitions));
     }
 
+    // Spark 4 transformWithState (used for streaming pipelines) requires the 
RocksDB state store.
+    // This is a harmless, inert configuration for batch pipelines on Spark 3.
+    sparkConf.setIfMissing(
+        "spark.sql.streaming.stateStore.providerClass",
+        
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider");
+
     return SparkSession.builder().config(sparkConf);
   }
 

Reply via email to