【问题标题】: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。您必须对其进行一些调整,因为它旨在读取静态文件,而不是附加到的文件。