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

public final class BeamReaderCache extends Object
Executor side cache of live Beam 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.

  • Method Details

    • key

      public static String key(String checkpointLocation, int splitId)
    • 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 IOException
      Returns the reader for key positioned at startEpoch, reusing the cached one if it is there, restoring from the durable mark otherwise. A zero length durable mark means a fresh start. committedEpoch supplies the epoch Spark last committed, -1 if unknown.
      Throws:
      IllegalStateException - if startEpoch > 0 and no durable mark exists
      IOException
    • invalidate

      public static void invalidate(String key)
      Closes and forgets the reader of key, nothing is finalized.
    • invalidateAll

      public static void invalidateAll()
      Closes and forgets every cached reader.