类 FlinkSourceReader<SplitT extends org.apache.seatunnel.api.source.SourceSplit>
- java.lang.Object
-
- org.apache.seatunnel.translation.flink.source.FlinkSourceReader<SplitT>
-
- 类型参数:
SplitT-
- 所有已实现的接口:
AutoCloseable,org.apache.flink.api.common.state.CheckpointListener,org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT>>
public class FlinkSourceReader<SplitT extends org.apache.seatunnel.api.source.SourceSplit> extends Object implements org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT>>
The implementation ofSourceReader, used for proxy allSourceReaderin flink.
-
-
构造器概要
构造器 构造器 说明 FlinkSourceReader(org.apache.seatunnel.api.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitT> sourceReader, org.apache.seatunnel.api.source.SourceReader.Context context, org.apache.seatunnel.shade.com.typesafe.config.Config envConfig)
-
方法概要
所有方法 实例方法 具体方法 修饰符和类型 方法 说明 voidaddSplits(List<SplitWrapper<SplitT>> splits)voidclose()voidhandleSourceEvents(org.apache.flink.api.connector.source.SourceEvent sourceEvent)CompletableFuture<Void>isAvailable()voidnotifyCheckpointAborted(long checkpointId)voidnotifyCheckpointComplete(long checkpointId)voidnotifyNoMoreSplits()org.apache.flink.core.io.InputStatuspollNext(org.apache.flink.api.connector.source.ReaderOutput<org.apache.seatunnel.api.table.type.SeaTunnelRow> output)List<SplitWrapper<SplitT>>snapshotState(long checkpointId)voidstart()
-
-
-
构造器详细资料
-
FlinkSourceReader
public FlinkSourceReader(org.apache.seatunnel.api.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitT> sourceReader, org.apache.seatunnel.api.source.SourceReader.Context context, org.apache.seatunnel.shade.com.typesafe.config.Config envConfig)
-
-
方法详细资料
-
start
public void start()
- 指定者:
start在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
pollNext
public org.apache.flink.core.io.InputStatus pollNext(org.apache.flink.api.connector.source.ReaderOutput<org.apache.seatunnel.api.table.type.SeaTunnelRow> output) throws Exception- 指定者:
pollNext在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>- 抛出:
Exception
-
snapshotState
public List<SplitWrapper<SplitT>> snapshotState(long checkpointId)
- 指定者:
snapshotState在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
isAvailable
public CompletableFuture<Void> isAvailable()
- 指定者:
isAvailable在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
addSplits
public void addSplits(List<SplitWrapper<SplitT>> splits)
- 指定者:
addSplits在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
notifyNoMoreSplits
public void notifyNoMoreSplits()
- 指定者:
notifyNoMoreSplits在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
handleSourceEvents
public void handleSourceEvents(org.apache.flink.api.connector.source.SourceEvent sourceEvent)
- 指定者:
handleSourceEvents在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>
-
close
public void close() throws Exception- 指定者:
close在接口中AutoCloseable- 抛出:
Exception
-
notifyCheckpointComplete
public void notifyCheckpointComplete(long checkpointId) throws Exception- 指定者:
notifyCheckpointComplete在接口中org.apache.flink.api.common.state.CheckpointListener- 指定者:
notifyCheckpointComplete在接口中org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>>- 抛出:
Exception
-
-