@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()