Class BeamReaderCache
UnboundedSource.UnboundedReaders keyed by checkpoint location and split.
An entry records the epoch its reader is positioned at and the mark taken there. A batch starting at that epoch reuses the reader and finalizes the pending mark, the start epoch of a batch is always committed by Spark. Any other start epoch, or a reader that moved without completing its batch, closes the entry without finalizing and restores the reader from the durable mark at the start epoch.
A sweeper thread closes readers idle for longer than their idle timeout, finalizing marks
whose epoch Spark committed, see BeamSourceCheckpoint.readSparkCommittedEpoch(), and
dropping the others, the source redelivers. The timeout must exceed the longest gap between two
micro-batches of one split. Under spark.sql.streaming.asyncProgressTrackingEnabled the
commit log lags, idle readers then drop their marks. Speculative execution can leave a losing
attempt's mark finalized on another executor, this source is not safe under
spark.speculation with sources whose reads are not deterministic.
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic final classA live reader with the epoch it is positioned at and the coded mark taken there. -
Method Summary
Modifier and TypeMethodDescriptionstatic <T> BeamReaderCache.CachedReader<T> acquire(String key, long startEpoch, UnboundedSource<T, ?> source, PipelineOptions options, long idleTimeoutMillis, LongSupplier committedEpoch, org.apache.beam.runners.spark.structuredstreaming.io.streaming.BeamReaderCache.MarkRestorer restorer) Returns the reader forkeypositioned atstartEpoch, reusing the cached one if it is there, restoring from the durable mark otherwise.static voidinvalidate(String key) Closes and forgets the reader ofkey, nothing is finalized.static voidCloses and forgets every cached reader.static String
-
Method Details
-
key
-
acquire
public static <T> BeamReaderCache.CachedReader<T> acquire(String key, long startEpoch, UnboundedSource<T, ?> source, PipelineOptions options, long idleTimeoutMillis, LongSupplier committedEpoch, org.apache.beam.runners.spark.structuredstreaming.io.streaming.BeamReaderCache.MarkRestorer restorer) throws IOExceptionReturns the reader forkeypositioned atstartEpoch, reusing the cached one if it is there, restoring from the durable mark otherwise. A zero length durable mark means a fresh start.committedEpochsupplies the epoch Spark last committed, -1 if unknown.- Throws:
IllegalStateException- ifstartEpoch > 0and no durable mark existsIOException
-
invalidate
Closes and forgets the reader ofkey, nothing is finalized. -
invalidateAll
public static void invalidateAll()Closes and forgets every cached reader.
-