类 FlinkSourceEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
- java.lang.Object
-
- org.apache.seatunnel.translation.flink.source.FlinkSourceEnumerator<SplitT,EnumStateT>
-
- 类型参数:
SplitT- The generic type of source splitEnumStateT- The generic type of enumerator state
- 所有已实现的接口:
AutoCloseable,org.apache.flink.api.common.state.CheckpointListener,org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT>
public class FlinkSourceEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT> extends Object implements org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT>
The implementation ofSplitEnumerator, used for proxy allSourceSplitEnumeratorin flink.
-
-
构造器概要
构造器 构造器 说明 FlinkSourceEnumerator(org.apache.seatunnel.api.source.SourceSplitEnumerator<SplitT,EnumStateT> enumerator, org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext)
-
方法概要
所有方法 实例方法 具体方法 修饰符和类型 方法 说明 voidaddReader(int subtaskId)voidaddSplitsBack(List<SplitWrapper<SplitT>> splits, int subtaskId)voidclose()voidhandleSourceEvent(int subtaskId, org.apache.flink.api.connector.source.SourceEvent sourceEvent)voidhandleSplitRequest(int subtaskId, String requesterHostname)voidnotifyCheckpointAborted(long checkpointId)voidnotifyCheckpointComplete(long checkpointId)EnumStateTsnapshotState(long checkpointId)voidstart()
-
-
-
构造器详细资料
-
FlinkSourceEnumerator
public FlinkSourceEnumerator(org.apache.seatunnel.api.source.SourceSplitEnumerator<SplitT,EnumStateT> enumerator, org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext)
-
-
方法详细资料
-
start
public void start()
- 指定者:
start在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
-
handleSplitRequest
public void handleSplitRequest(int subtaskId, @Nullable String requesterHostname)- 指定者:
handleSplitRequest在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
-
addSplitsBack
public void addSplitsBack(List<SplitWrapper<SplitT>> splits, int subtaskId)
- 指定者:
addSplitsBack在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
-
addReader
public void addReader(int subtaskId)
- 指定者:
addReader在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
-
snapshotState
public EnumStateT snapshotState(long checkpointId) throws Exception
- 指定者:
snapshotState在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>- 抛出:
Exception
-
close
public void close() throws IOException- 指定者:
close在接口中AutoCloseable- 指定者:
close在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>- 抛出:
IOException
-
handleSourceEvent
public void handleSourceEvent(int subtaskId, org.apache.flink.api.connector.source.SourceEvent sourceEvent)- 指定者:
handleSourceEvent在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>
-
notifyCheckpointComplete
public void notifyCheckpointComplete(long checkpointId) throws Exception- 指定者:
notifyCheckpointComplete在接口中org.apache.flink.api.common.state.CheckpointListener- 指定者:
notifyCheckpointComplete在接口中org.apache.flink.api.connector.source.SplitEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT>- 抛出:
Exception
-
-