diff --git a/runners/spark/4/build.gradle b/runners/spark/4/build.gradle index f2746b06158e..ae8750ec5b11 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 3371a403b2c9..29cc4cb99cfb 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 9d3419e19473..b592b6fb742d 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.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 MetricResults metrics() { @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 96717f29e87f..f78026847fad 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.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 @@ * href="https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html">Structured * Streaming framework). * - *
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. + *
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. * *
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 SparkStructuredStreamingPipelineResult run(final Pipeline pipeline) {
+ " 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 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 a681dea2fde5..8470c968550b 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 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 EvaluationContext translate(
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.
+ *
+ * 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 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 5b9e5b6fae86..822d1871b12e 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 @@ private static SparkSession.Builder sessionBuilder(
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);
}