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

  • 类型参数:
    SplitT - The generic type of source split
    EnumStateT - The generic type of enumerator state
    所有已实现的接口:
    Serializable, org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT>,​EnumStateT>, org.apache.flink.api.java.typeutils.ResultTypeQueryable<org.apache.seatunnel.api.table.type.SeaTunnelRow>

    public class FlinkSource<SplitT extends org.apache.seatunnel.api.source.SourceSplit,​EnumStateT extends Serializable>
    extends Object
    implements org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT>,​EnumStateT>, org.apache.flink.api.java.typeutils.ResultTypeQueryable<org.apache.seatunnel.api.table.type.SeaTunnelRow>
    The source implementation of Source, used for proxy all SeaTunnelSource in flink.
    另请参阅:
    序列化表格
    • 构造器详细资料

      • FlinkSource

        public FlinkSource​(org.apache.seatunnel.api.source.SeaTunnelSource<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitT,​EnumStateT> source,
                           org.apache.seatunnel.shade.com.typesafe.config.Config envConfig)
    • 方法详细资料

      • getBoundedness

        public org.apache.flink.api.connector.source.Boundedness getBoundedness()
        指定者:
        getBoundedness 在接口中 org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>,​EnumStateT extends Serializable>
      • createReader

        public org.apache.flink.api.connector.source.SourceReader<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT>> createReader​(org.apache.flink.api.connector.source.SourceReaderContext readerContext)
                                                                                                                                                     throws Exception
        指定者:
        createReader 在接口中 org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>,​EnumStateT extends Serializable>
        抛出:
        Exception
      • getSplitSerializer

        public org.apache.flink.core.io.SimpleVersionedSerializer<SplitWrapper<SplitT>> getSplitSerializer()
        指定者:
        getSplitSerializer 在接口中 org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>,​EnumStateT extends Serializable>
      • getEnumeratorCheckpointSerializer

        public org.apache.flink.core.io.SimpleVersionedSerializer<EnumStateT> getEnumeratorCheckpointSerializer()
        指定者:
        getEnumeratorCheckpointSerializer 在接口中 org.apache.flink.api.connector.source.Source<org.apache.seatunnel.api.table.type.SeaTunnelRow,​SplitWrapper<SplitT extends org.apache.seatunnel.api.source.SourceSplit>,​EnumStateT extends Serializable>
      • getProducedType

        public org.apache.flink.api.common.typeinfo.TypeInformation<org.apache.seatunnel.api.table.type.SeaTunnelRow> getProducedType()
        指定者:
        getProducedType 在接口中 org.apache.flink.api.java.typeutils.ResultTypeQueryable<SplitT extends org.apache.seatunnel.api.source.SourceSplit>