public class UnboundedSourceImpl extends org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>
Modifier and Type | Method and Description |
---|---|
org.apache.beam.sdk.io.UnboundedSource.UnboundedReader<com.google.cloud.pubsublite.proto.SequencedMessage> |
createReader(org.apache.beam.sdk.options.PipelineOptions options,
@Nullable CheckpointMarkImpl checkpointMark) |
org.apache.beam.sdk.coders.Coder<CheckpointMarkImpl> |
getCheckpointMarkCoder() |
org.apache.beam.sdk.coders.Coder<com.google.cloud.pubsublite.proto.SequencedMessage> |
getOutputCoder() |
java.util.List<? extends org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>> |
split(int desiredNumSplits,
org.apache.beam.sdk.options.PipelineOptions options) |
public java.util.List<? extends org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>> split(int desiredNumSplits, org.apache.beam.sdk.options.PipelineOptions options) throws java.lang.Exception
split
in class org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>
java.lang.Exception
public org.apache.beam.sdk.io.UnboundedSource.UnboundedReader<com.google.cloud.pubsublite.proto.SequencedMessage> createReader(org.apache.beam.sdk.options.PipelineOptions options, @Nullable CheckpointMarkImpl checkpointMark) throws java.io.IOException
createReader
in class org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>
java.io.IOException
public org.apache.beam.sdk.coders.Coder<CheckpointMarkImpl> getCheckpointMarkCoder()
getCheckpointMarkCoder
in class org.apache.beam.sdk.io.UnboundedSource<com.google.cloud.pubsublite.proto.SequencedMessage,CheckpointMarkImpl>
public org.apache.beam.sdk.coders.Coder<com.google.cloud.pubsublite.proto.SequencedMessage> getOutputCoder()
getOutputCoder
in class org.apache.beam.sdk.io.Source<com.google.cloud.pubsublite.proto.SequencedMessage>