@Internal public class GenerateInitialPartitionsAction extends java.lang.Object
DetectNewPartitionsDoFn
.Constructor and Description |
---|
GenerateInitialPartitionsAction(ChangeStreamMetrics metrics,
ChangeStreamDao changeStreamDao) |
Modifier and Type | Method and Description |
---|---|
DoFn.ProcessContinuation |
run(DoFn.OutputReceiver<PartitionRecord> receiver,
RestrictionTracker<OffsetRange,java.lang.Long> tracker,
ManualWatermarkEstimator<Instant> watermarkEstimator,
Instant startTime)
The very first step of the pipeline when there are no partitions being streamed yet.
|
public GenerateInitialPartitionsAction(ChangeStreamMetrics metrics, ChangeStreamDao changeStreamDao)
public DoFn.ProcessContinuation run(DoFn.OutputReceiver<PartitionRecord> receiver, RestrictionTracker<OffsetRange,java.lang.Long> tracker, ManualWatermarkEstimator<Instant> watermarkEstimator, Instant startTime)
DoFn.ProcessContinuation.resume()
if the stream continues, otherwise DoFn.ProcessContinuation.stop()