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