Package org.apache.beam.runners.spark.structuredstreaming.io.streaming
package org.apache.beam.runners.spark.structuredstreaming.io.streaming
-
ClassesClassDescriptionExecutor side cache of live Beam
UnboundedSource.UnboundedReaders keyed by checkpoint location and split.A live reader with the epoch it is positioned at and the coded mark taken there.Durable state of one Beam unbounded source under the per source checkpoint location Spark hands totoMicroBatchStream.Translator facing entry point turning a BeamUnboundedSourceinto a streaming SparkDatasetof rows, with the DataSourceV2 micro-batch glue as nested classes.