【问题标题】:Reading file that is being appended in Flink读取 Flink 中追加的文件
【发布时间】:2020-03-22 14:31:27
【问题描述】:

我们有一个旧版应用程序正在将结果作为记录写入一些本地文件。我们希望实时处理这些记录,因此我们计划使用 Flink 作为引擎。我知道我可以使用StreamingExecutionEnvironment#readFile 读取文本文件。似乎我们需要类似于 PROCESS_CONTINUOUSLY 的东西,但是这个标志会导致在每次更改时重新处理整个文件,这不是我们想要的。

当然,我可以编写我的自定义源,以在其状态下保存每个文件的记录数。但我认为这种方法在检查点或其他方面可能存在一些问题——我的理由是,如果这很容易可靠地实现,那么它已经在 Flink 中实现了。

任何提示/建议如何解决这个问题?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    只要您愿意从单个文件(每个源实例)中读取内容,您就可以使用自定义源相当轻松地做到这一点。您将需要使用操作员状态并实施检查点。状态处理和检查点将如下所示:

    public class CheckpointedFileSource implements SourceFunction<Event>, ListCheckpointed<Long> {
        private long eventCnt = 0;
    
        public void run(SourceContext<Event> sourceContext) throws Exception {
            final Object lock = sourceContext.getCheckpointLock();
    
            // skip over previously emitted events
            ...
    
            while (not cancelled) {
                read event from file;
    
                synchronized (lock) {
                    eventCnt++;
                    sourceContext.collectWithTimestamp(event, timestamp);
                }
    
            }
        }
    
        @Override
        public List<Long> snapshotState(long checkpointId, long checkpointTimestamp) throws Exception {
            return Collections.singletonList(eventCnt);
        }
    
        @Override
        public void restoreState(List<Long> state) throws Exception {
            for (Long s : state)
                this.eventCnt = s;
        }
    
    }
    

    有关完整示例,请参阅 Flink 训练练习中使用的 the checkpointed taxi ride data source。您必须对其进行一些调整,因为它旨在读取静态文件,而不是附加到的文件。

    【讨论】:

      猜你喜欢
      • 2017-02-07
      • 2023-04-09
      • 2013-02-16
      • 2011-08-03
      • 1970-01-01
      • 1970-01-01
      • 2022-12-07
      • 1970-01-01
      • 2011-05-05
      相关资源
      最近更新 更多