public class AlignedSourceReader extends FileStoreSourceReader implements org.apache.flink.api.connector.source.ExternallyInducedSourceReader<org.apache.flink.table.data.RowData,FileStoreSourceSplit>
FileStoreSourceReader
is that only after the allocated splits are
fully consumed, checkpoints will be made and the next batch of splits will be requested.Constructor and Description |
---|
AlignedSourceReader(org.apache.flink.api.connector.source.SourceReaderContext readerContext,
TableRead tableRead,
FileStoreSourceReaderMetrics metrics,
IOManager ioManager,
Long limit,
org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue<org.apache.flink.connector.base.source.reader.RecordsWithSplitIds<org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData>>> elementsQueue,
NestedProjectedRowData rowData) |
Modifier and Type | Method and Description |
---|---|
void |
handleSourceEvents(org.apache.flink.api.connector.source.SourceEvent sourceEvent) |
protected void |
onSplitFinished(Map<String,FileStoreSourceSplitState> finishedSplitIds) |
Optional<Long> |
shouldTriggerCheckpoint() |
close, initializedState, start, toSplitType
addSplits, getNumberOfCurrentlyAssignedSplits, isAvailable, notifyNoMoreSplits, pauseOrResumeSplits, pollNext, snapshotState
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
addSplits, isAvailable, notifyCheckpointComplete, notifyNoMoreSplits, pauseOrResumeSplits, pollNext, snapshotState, start
close
public AlignedSourceReader(org.apache.flink.api.connector.source.SourceReaderContext readerContext, TableRead tableRead, FileStoreSourceReaderMetrics metrics, IOManager ioManager, @Nullable Long limit, org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue<org.apache.flink.connector.base.source.reader.RecordsWithSplitIds<org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData>>> elementsQueue, @Nullable NestedProjectedRowData rowData)
public void handleSourceEvents(org.apache.flink.api.connector.source.SourceEvent sourceEvent)
handleSourceEvents
in interface org.apache.flink.api.connector.source.SourceReader<org.apache.flink.table.data.RowData,FileStoreSourceSplit>
handleSourceEvents
in class org.apache.flink.connector.base.source.reader.SourceReaderBase<org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator<org.apache.flink.table.data.RowData>,org.apache.flink.table.data.RowData,FileStoreSourceSplit,FileStoreSourceSplitState>
protected void onSplitFinished(Map<String,FileStoreSourceSplitState> finishedSplitIds)
onSplitFinished
in class FileStoreSourceReader
public Optional<Long> shouldTriggerCheckpoint()
shouldTriggerCheckpoint
in interface org.apache.flink.api.connector.source.ExternallyInducedSourceReader<org.apache.flink.table.data.RowData,FileStoreSourceSplit>
Copyright © 2023–2025 The Apache Software Foundation. All rights reserved.