类 FlinkSourceEnumerator<SplitT extends org.apache.seatunnel.api.source.SourceSplit,​EnumStateT>

  • 类型参数:
    SplitT - The generic type of source split
    EnumStateT - 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 of SplitEnumerator, used for proxy all SourceSplitEnumerator in flink.
    • 构造器详细资料

      • 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
      • 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
      • notifyCheckpointAborted

        public void notifyCheckpointAborted​(long checkpointId)
                                     throws Exception
        指定者:
        notifyCheckpointAborted 在接口中 org.apache.flink.api.common.state.CheckpointListener
        抛出:
        Exception