Class SpannerChangeStreamSource<T>
java.lang.Object
io.github.flink.gcp.connector.spanner.source.SpannerChangeStreamSource<T>
- Type Parameters:
T- the record type produced
- All Implemented Interfaces:
Serializable,org.apache.flink.api.connector.source.Source<T,,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState> org.apache.flink.api.connector.source.SourceReaderFactory<T,,ChangeStreamPartitionSplit> org.apache.flink.api.java.typeutils.ResultTypeQueryable<T>
@PublicEvolving
public final class SpannerChangeStreamSource<T>
extends Object
implements org.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit,SpannerChangeStreamEnumeratorState>, org.apache.flink.api.java.typeutils.ResultTypeQueryable<T>
FLIP-27 source for Cloud Spanner Change Streams.
- See Also:
-
Method Summary
Modifier and TypeMethodDescriptionstatic <T> SpannerChangeStreamSourceBuilder<T>builder()Returns a builder for a Cloud Spanner Change Streams source.org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,SpannerChangeStreamEnumeratorState> createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context) org.apache.flink.api.connector.source.SourceReader<T,ChangeStreamPartitionSplit> createReader(org.apache.flink.api.connector.source.SourceReaderContext context) org.apache.flink.api.connector.source.Boundednessorg.apache.flink.core.io.SimpleVersionedSerializer<SpannerChangeStreamEnumeratorState>org.apache.flink.api.common.typeinfo.TypeInformation<T>org.apache.flink.core.io.SimpleVersionedSerializer<ChangeStreamPartitionSplit>org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,SpannerChangeStreamEnumeratorState> restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context, SpannerChangeStreamEnumeratorState checkpoint) Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.flink.api.connector.source.Source
declareWatermarks
-
Method Details
-
builder
Returns a builder for a Cloud Spanner Change Streams source.- Type Parameters:
T- the record type produced- Returns:
- the builder
-
getBoundedness
public org.apache.flink.api.connector.source.Boundedness getBoundedness()- Specified by:
getBoundednessin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState>
-
createReader
public org.apache.flink.api.connector.source.SourceReader<T,ChangeStreamPartitionSplit> createReader(org.apache.flink.api.connector.source.SourceReaderContext context) throws Exception - Specified by:
createReaderin interfaceorg.apache.flink.api.connector.source.SourceReaderFactory<T,ChangeStreamPartitionSplit> - Throws:
Exception
-
createEnumerator
public org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,SpannerChangeStreamEnumeratorState> createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context) - Specified by:
createEnumeratorin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState>
-
restoreEnumerator
public org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,SpannerChangeStreamEnumeratorState> restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context, SpannerChangeStreamEnumeratorState checkpoint) - Specified by:
restoreEnumeratorin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState>
-
getSplitSerializer
public org.apache.flink.core.io.SimpleVersionedSerializer<ChangeStreamPartitionSplit> getSplitSerializer()- Specified by:
getSplitSerializerin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState>
-
getEnumeratorCheckpointSerializer
public org.apache.flink.core.io.SimpleVersionedSerializer<SpannerChangeStreamEnumeratorState> getEnumeratorCheckpointSerializer()- Specified by:
getEnumeratorCheckpointSerializerin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, SpannerChangeStreamEnumeratorState>
-
getProducedType
- Specified by:
getProducedTypein interfaceorg.apache.flink.api.java.typeutils.ResultTypeQueryable<T>
-