类 FlinkSink<InputT,CommT,WriterStateT,GlobalCommT>
- java.lang.Object
-
- org.apache.seatunnel.translation.flink.sink.FlinkSink<InputT,CommT,WriterStateT,GlobalCommT>
-
- 类型参数:
InputT- The generic type of input dataCommT- The generic type of commit messageWriterStateT- The generic type of writer stateGlobalCommT- The generic type of global commit message
- 所有已实现的接口:
Serializable,org.apache.flink.api.connector.sink.Sink<InputT,CommitWrapper<CommT>,FlinkWriterState<WriterStateT>,GlobalCommT>
public class FlinkSink<InputT,CommT,WriterStateT,GlobalCommT> extends Object implements org.apache.flink.api.connector.sink.Sink<InputT,CommitWrapper<CommT>,FlinkWriterState<WriterStateT>,GlobalCommT>
The sink implementation ofSink, the entrypoint of flink sink translation- 另请参阅:
- 序列化表格
-
-
构造器概要
构造器 构造器 说明 FlinkSink(org.apache.seatunnel.api.sink.SeaTunnelSink<org.apache.seatunnel.api.table.type.SeaTunnelRow,WriterStateT,CommT,GlobalCommT> sink, org.apache.seatunnel.api.table.catalog.CatalogTable catalogTable)
-
方法概要
所有方法 实例方法 具体方法 修饰符和类型 方法 说明 Optional<org.apache.flink.api.connector.sink.Committer<CommitWrapper<CommT>>>createCommitter()Optional<org.apache.flink.api.connector.sink.GlobalCommitter<CommitWrapper<CommT>,GlobalCommT>>createGlobalCommitter()org.apache.flink.api.connector.sink.SinkWriter<InputT,CommitWrapper<CommT>,FlinkWriterState<WriterStateT>>createWriter(org.apache.flink.api.connector.sink.Sink.InitContext context, List<FlinkWriterState<WriterStateT>> states)Optional<org.apache.flink.core.io.SimpleVersionedSerializer<CommitWrapper<CommT>>>getCommittableSerializer()Optional<org.apache.flink.core.io.SimpleVersionedSerializer<GlobalCommT>>getGlobalCommittableSerializer()Optional<org.apache.flink.core.io.SimpleVersionedSerializer<FlinkWriterState<WriterStateT>>>getWriterStateSerializer()
-
-
-
构造器详细资料
-
FlinkSink
public FlinkSink(org.apache.seatunnel.api.sink.SeaTunnelSink<org.apache.seatunnel.api.table.type.SeaTunnelRow,WriterStateT,CommT,GlobalCommT> sink, org.apache.seatunnel.api.table.catalog.CatalogTable catalogTable)
-
-
方法详细资料
-
createWriter
public org.apache.flink.api.connector.sink.SinkWriter<InputT,CommitWrapper<CommT>,FlinkWriterState<WriterStateT>> createWriter(org.apache.flink.api.connector.sink.Sink.InitContext context, List<FlinkWriterState<WriterStateT>> states) throws IOException
- 指定者:
createWriter在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>- 抛出:
IOException
-
createCommitter
public Optional<org.apache.flink.api.connector.sink.Committer<CommitWrapper<CommT>>> createCommitter() throws IOException
- 指定者:
createCommitter在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>- 抛出:
IOException
-
createGlobalCommitter
public Optional<org.apache.flink.api.connector.sink.GlobalCommitter<CommitWrapper<CommT>,GlobalCommT>> createGlobalCommitter() throws IOException
- 指定者:
createGlobalCommitter在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>- 抛出:
IOException
-
getCommittableSerializer
public Optional<org.apache.flink.core.io.SimpleVersionedSerializer<CommitWrapper<CommT>>> getCommittableSerializer()
- 指定者:
getCommittableSerializer在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>
-
getGlobalCommittableSerializer
public Optional<org.apache.flink.core.io.SimpleVersionedSerializer<GlobalCommT>> getGlobalCommittableSerializer()
- 指定者:
getGlobalCommittableSerializer在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>
-
getWriterStateSerializer
public Optional<org.apache.flink.core.io.SimpleVersionedSerializer<FlinkWriterState<WriterStateT>>> getWriterStateSerializer()
- 指定者:
getWriterStateSerializer在接口中org.apache.flink.api.connector.sink.Sink<InputT,CommT,WriterStateT,GlobalCommT>
-
-