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

public final class BeamSourceCheckpoint extends Object
Durable state of one Beam unbounded source under the per source checkpoint location Spark hands to toMicroBatchStream.

<location>/splits pins the split list, written once by the driver. <location>/marks/<epoch>/<splitId> holds the coded checkpoint mark of a split at the end of the batch ending at that epoch. The epoch Spark last committed is read from Spark's own commits and offsets logs two levels up. All IO goes through Spark's CheckpointFileManager, writes are atomic renames.

  • Constructor Details

    • BeamSourceCheckpoint

      public BeamSourceCheckpoint(String checkpointLocation, org.apache.hadoop.conf.Configuration hadoopConf)
  • Method Details

    • location

      public String location()
    • readSplits

      public @Nullable List<UnboundedSource<?,?>> readSplits() throws IOException
      The pinned split list, or null if none was pinned yet.
      Throws:
      IOException
    • writeSplits

      public void writeSplits(List<? extends UnboundedSource<?,?>> splits) throws IOException
      Pins the split list, fails if one is pinned already.
      Throws:
      IOException
    • prepareEpoch

      public void prepareEpoch(long epoch) throws IOException
      Creates the mark directory of epoch, the driver calls this once per batch.
      Throws:
      IOException
    • writeMark

      public void writeMark(int splitId, long epoch, byte[] codedMark) throws IOException
      Writes the mark, creating the epoch directory if a manager without parent creation needs it.
      Throws:
      IOException
    • readMark

      public byte @Nullable [] readMark(int splitId, long epoch) throws IOException
      The coded mark of a split at an epoch, or null if absent.
      Throws:
      IOException
    • readSparkCommittedEpoch

      public long readSparkCommittedEpoch()
      The end epoch of this source in the last batch Spark committed, or -1 if there is none or the logs cannot be read. The location is <root>/sources/<index>, the batch id is the highest entry of <root>/commits and its epoch is line index after the version and metadata lines of <root>/offsets/<id>.
    • purgeMarksBelow

      public void purgeMarksBelow(long epoch) throws IOException
      Deletes the marks of every epoch strictly below epoch, one recursive delete per epoch directory. Lists the marks directory once, later calls delete the range above the previous floor only. Idempotent.
      Throws:
      IOException