类 FlinkRowCollector

  • 所有已实现的接口:
    org.apache.seatunnel.api.source.Collector<org.apache.seatunnel.api.table.type.SeaTunnelRow>

    public class FlinkRowCollector
    extends Object
    implements org.apache.seatunnel.api.source.Collector<org.apache.seatunnel.api.table.type.SeaTunnelRow>
    The implementation of Collector for flink engine.
    • 构造器详细资料

      • FlinkRowCollector

        public FlinkRowCollector​(org.apache.seatunnel.shade.com.typesafe.config.Config envConfig,
                                 org.apache.seatunnel.api.common.metrics.MetricsContext metricsContext)
    • 方法详细资料

      • collect

        public void collect​(org.apache.seatunnel.api.table.type.SeaTunnelRow record)
        指定者:
        collect 在接口中 org.apache.seatunnel.api.source.Collector<org.apache.seatunnel.api.table.type.SeaTunnelRow>
      • getCheckpointLock

        public Object getCheckpointLock()
        指定者:
        getCheckpointLock 在接口中 org.apache.seatunnel.api.source.Collector<org.apache.seatunnel.api.table.type.SeaTunnelRow>
      • withReaderOutput

        public FlinkRowCollector withReaderOutput​(org.apache.flink.api.connector.source.ReaderOutput<org.apache.seatunnel.api.table.type.SeaTunnelRow> readerOutput)