java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.io.streaming.UnboundedSourceDataset

public final class UnboundedSourceDataset extends Object
Translator facing entry point turning a Beam UnboundedSource into a streaming Spark Dataset of rows, with the DataSourceV2 micro-batch glue as nested classes.

The dataset has two columns, "payload" of type BINARY holding the element encoded with the supplied WindowedValue coder, and "eventTimestamp" of type TIMESTAMP holding the event timestamp of that element.

The event time watermark is declared here and only here. Spark 4 rejects a second withWatermark further down the plan, so downstream translators must never call it again.

  • Field Details

  • Method Details

    • of

      public static <T, CheckpointMarkT extends UnboundedSource.CheckpointMark> org.apache.spark.sql.Dataset<org.apache.spark.sql.Row> of(org.apache.spark.sql.SparkSession session, UnboundedSource<T,CheckpointMarkT> source, Coder<WindowedValue<T>> windowedValueCoder, SparkStructuredStreamingPipelineOptions options, String transformName)
      Builds the streaming Dataset for source with the event time watermark applied.
      Type Parameters:
      T - the element type of the source
      CheckpointMarkT - the checkpoint mark type of the source
      Parameters:
      session - the active Spark session
      source - the Beam unbounded source to read
      windowedValueCoder - the coder of the "payload" column
      options - the pipeline options, supplying the watermark delay and the micro-batch limits
      transformName - the full name of the read transform, used for naming only