类 FlinkSource<SplitT extends org.apache.seatunnel.api.source.SourceSplit,EnumStateT extends Serializable>
- java.lang.Object
-
- org.apache.seatunnel.translation.flink.source.FlinkSource<SplitT,EnumStateT>
-
- 类型参数:
SplitT- The generic type of source splitEnumStateT- 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 ofSource, used for proxy allSeaTunnelSourcein flink.- 另请参阅:
- 序列化表格
-
-
构造器概要
构造器 构造器 说明 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)
-
方法概要
所有方法 实例方法 具体方法 修饰符和类型 方法 说明 org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT>createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext)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)org.apache.flink.api.connector.source.BoundednessgetBoundedness()org.apache.flink.core.io.SimpleVersionedSerializer<EnumStateT>getEnumeratorCheckpointSerializer()org.apache.flink.api.common.typeinfo.TypeInformation<org.apache.seatunnel.api.table.type.SeaTunnelRow>getProducedType()org.apache.flink.core.io.SimpleVersionedSerializer<SplitWrapper<SplitT>>getSplitSerializer()org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT>restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext, EnumStateT checkpoint)
-
-
-
构造器详细资料
-
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
-
createEnumerator
public org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT> createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext) throws Exception
- 指定者:
createEnumerator在接口中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
-
restoreEnumerator
public org.apache.flink.api.connector.source.SplitEnumerator<SplitWrapper<SplitT>,EnumStateT> restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<SplitWrapper<SplitT>> enumContext, EnumStateT checkpoint) throws Exception
- 指定者:
restoreEnumerator在接口中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>
-
-