Class BigtableChangeStreamSource<T>
java.lang.Object
io.github.flink.gcp.connector.bigtable.source.BigtableChangeStreamSource<T>
- Type Parameters:
T- the record type produced
- All Implemented Interfaces:
Serializable,org.apache.flink.api.connector.source.Source<T,,ChangeStreamPartitionSplit, BigtableChangeStreamEnumeratorState> 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 BigtableChangeStreamSource<T>
extends Object
implements org.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit,BigtableChangeStreamEnumeratorState>, org.apache.flink.api.java.typeutils.ResultTypeQueryable<T>, org.apache.flink.streaming.api.lineage.LineageVertexProvider
FLIP-27 source for Bigtable Change Streams.
Lineage reports the configured data table with namespace
bigtable://{project}/{instance}, name {table} and a gcp physical-resource facet.
Its boundedness follows getBoundedness(); no external metadata table is reported.
Extraction calls no user schema and opens no client. The supported Flink 2.x versions extract the
metadata automatically; Flink 1.20 supports direct inspection only.
- See Also:
-
Method Summary
Modifier and TypeMethodDescriptionstatic <T> BigtableChangeStreamSourceBuilder<T>builder()Returns a builder for a Bigtable Change Streams source.org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,BigtableChangeStreamEnumeratorState> 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<BigtableChangeStreamEnumeratorState>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,BigtableChangeStreamEnumeratorState> restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context, BigtableChangeStreamEnumeratorState checkpoint) withTableLineage(String logicalName) Returns a copy carrying the logical Table identity with the same source configuration.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
-
withTableLineage
Returns a copy carrying the logical Table identity with the same source configuration. -
getLineageVertex
public org.apache.flink.streaming.api.lineage.SourceLineageVertex getLineageVertex()- Specified by:
getLineageVertexin interfaceorg.apache.flink.streaming.api.lineage.LineageVertexProvider
-
builder
Returns a builder for a Bigtable 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, BigtableChangeStreamEnumeratorState>
-
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,BigtableChangeStreamEnumeratorState> createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context) throws Exception - Specified by:
createEnumeratorin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, BigtableChangeStreamEnumeratorState> - Throws:
Exception
-
restoreEnumerator
public org.apache.flink.api.connector.source.SplitEnumerator<ChangeStreamPartitionSplit,BigtableChangeStreamEnumeratorState> restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<ChangeStreamPartitionSplit> context, BigtableChangeStreamEnumeratorState checkpoint) throws Exception - Specified by:
restoreEnumeratorin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, BigtableChangeStreamEnumeratorState> - Throws:
Exception
-
getSplitSerializer
public org.apache.flink.core.io.SimpleVersionedSerializer<ChangeStreamPartitionSplit> getSplitSerializer()- Specified by:
getSplitSerializerin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, BigtableChangeStreamEnumeratorState>
-
getEnumeratorCheckpointSerializer
public org.apache.flink.core.io.SimpleVersionedSerializer<BigtableChangeStreamEnumeratorState> getEnumeratorCheckpointSerializer()- Specified by:
getEnumeratorCheckpointSerializerin interfaceorg.apache.flink.api.connector.source.Source<T,ChangeStreamPartitionSplit, BigtableChangeStreamEnumeratorState>
-
getProducedType
- Specified by:
getProducedTypein interfaceorg.apache.flink.api.java.typeutils.ResultTypeQueryable<T>
-