Class SparkStructuredStreamingRunner

java.lang.Object
org.apache.beam.sdk.PipelineRunner<SparkStructuredStreamingPipelineResult>
org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingRunner

public final class SparkStructuredStreamingRunner extends PipelineRunner<SparkStructuredStreamingPipelineResult>
A Spark runner build on top of Spark's SQL Engine (Structured Streaming framework).

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.

To run a Beam pipeline with the default options using Spark's local mode, we would do the following:


 Pipeline p = [logic for pipeline creation]
 PipelineResult result = p.run();
 

To create a pipeline runner to run against a different spark cluster, with a custom master url we would do the following:


 Pipeline p = [logic for pipeline creation]
 SparkCommonPipelineOptions options = p.getOptions.as(SparkCommonPipelineOptions.class);
 options.setSparkMaster("spark://host:port");
 PipelineResult result = p.run();