类 FlinkSourceReader<SplitT extends org.apache.seatunnel.api.source.SourceSplit>

  • 类型参数:
    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 of SourceReader, used for proxy all SourceReader in flink.
    • 构造器详细资料

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

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