Class UnboundedSourceDataset
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.io.streaming.UnboundedSourceDataset
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 Summary
Fields -
Method Summary
Modifier and TypeMethodDescriptionstatic <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 streamingDatasetforsourcewith the event time watermark applied.
-
Field Details
-
COL_PAYLOAD
- See Also:
-
COL_EVENT_TS
- See Also:
-
SCHEMA
public static final org.apache.spark.sql.types.StructType SCHEMA
-
-
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 streamingDatasetforsourcewith the event time watermark applied.- Type Parameters:
T- the element type of the sourceCheckpointMarkT- the checkpoint mark type of the source- Parameters:
session- the active Spark sessionsource- the Beam unbounded source to readwindowedValueCoder- the coder of the "payload" columnoptions- the pipeline options, supplying the watermark delay and the micro-batch limitstransformName- the full name of the read transform, used for naming only
-