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>,org.apache.flink.streaming.api.lineage.LineageVertexProvider
@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>, org.apache.flink.streaming.api.lineage.LineageVertexProvider
FLIP-27 source for Cloud Spanner Change Streams.
Lineage identifies the configured Change Stream in namespace
spanner://project:instance, with name database/changeStreams/stream and resource kind
spanner-change-stream in the gcp facet. It does not infer watched tables from
filters or discover them at runtime. Extraction retains this source's boundedness and does not
open clients, resolve credentials, or invoke the deserializer. Flink 2.x extracts the metadata
automatically; Flink 1.20 supports direct inspection but not native listener delivery.
- 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.streaming.api.lineage.SourceLineageVertexorg.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>
-
getLineageVertex
public org.apache.flink.streaming.api.lineage.SourceLineageVertex getLineageVertex()- Specified by:
getLineageVertexin interfaceorg.apache.flink.streaming.api.lineage.LineageVertexProvider
-
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>
-