Class BeamSourceCheckpoint
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.io.streaming.BeamSourceCheckpoint
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 Summary
ConstructorsConstructorDescriptionBeamSourceCheckpoint(String checkpointLocation, org.apache.hadoop.conf.Configuration hadoopConf) -
Method Summary
Modifier and TypeMethodDescriptionlocation()voidprepareEpoch(long epoch) Creates the mark directory ofepoch, the driver calls this once per batch.voidpurgeMarksBelow(long epoch) Deletes the marks of every epoch strictly belowepoch, one recursive delete per epoch directory.byte @Nullable []readMark(int splitId, long epoch) The coded mark of a split at an epoch, or null if absent.longThe end epoch of this source in the last batch Spark committed, or -1 if there is none or the logs cannot be read.@Nullable List<UnboundedSource<?, ?>> The pinned split list, or null if none was pinned yet.voidwriteMark(int splitId, long epoch, byte[] codedMark) Writes the mark, creating the epoch directory if a manager without parent creation needs it.voidwriteSplits(List<? extends UnboundedSource<?, ?>> splits) Pins the split list, fails if one is pinned already.
-
Constructor Details
-
BeamSourceCheckpoint
public BeamSourceCheckpoint(String checkpointLocation, org.apache.hadoop.conf.Configuration hadoopConf)
-
-
Method Details
-
location
-
readSplits
The pinned split list, or null if none was pinned yet.- Throws:
IOException
-
writeSplits
Pins the split list, fails if one is pinned already.- Throws:
IOException
-
prepareEpoch
Creates the mark directory ofepoch, the driver calls this once per batch.- Throws:
IOException
-
writeMark
Writes the mark, creating the epoch directory if a manager without parent creation needs it.- Throws:
IOException
-
readMark
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>/commitsand its epoch is lineindexafter the version and metadata lines of<root>/offsets/<id>. -
purgeMarksBelow
Deletes the marks of every epoch strictly belowepoch, one recursive delete per epoch directory. Lists the marks directory once, later calls delete the range above the previous floor only. Idempotent.- Throws:
IOException
-