Flink 1.18.1持续文件源重复读取同名文件及路径已处理列表内存问题咨询
我遇到了这样一个场景:外部系统每小时生成一个同名文件,我使用Flink 1.18.1的FileSource读取该目录并开启「读后删除」配置,但第一次处理完成后,后续生成的同名文件无法被识别和读取。
排查后发现问题出在ContinuousFileSplitEnumerator的processDiscoveredSplits方法中:默认逻辑仅通过pathsAlreadyProcessed集合判断文件路径是否已被处理,完全忽略了FileSourceSplit中已存在的fileModificationTime字段,导致同名文件即使更新也会被过滤掉。
默认过滤逻辑代码如下:
private void processDiscoveredSplits(Collection<FileSourceSplit> splits, Throwable error) { if (error != null) { LOG.error("Failed to enumerate files", error); return; } final Collection<FileSourceSplit> newSplits = splits.stream() .filter((split) -> pathsAlreadyProcessed.add(split.path())) .collect(Collectors.toList()); splitAssigner.addSplits(newSplits); assignSplits(); }
同时ContinuousFileSplitEnumerator在AbstractFileSource的createSplitEnumerator()中被硬编码,起初我以为无法替换它;且官方文档提到扫描目录时会参考文件最后修改时间,但该逻辑并未在拆分枚举器中生效。
我的疑问
- 如何让过滤逻辑结合
fileModificationTime字段判断?是否有现成的实现方案? pathsAlreadyProcessed会持续累加文件路径,它何时会被清空?这会导致检查点内存占用越来越大,该如何处理?
问题解答
1. 结合文件修改时间实现过滤的方案
针对Flink 1.18.1版本,有两种可行的解决思路:
方案1:自定义SplitEnumerator并注入(推荐)
虽然默认枚举器被硬编码在AbstractFileSource中,但Flink的DiscoverySettings提供了扩展点,可以通过自定义EnumeratorProvider替换默认实现:
- 步骤1:自定义枚举器逻辑
复制默认ContinuousFileSplitEnumerator的代码,将pathsAlreadyProcessed从Set<String>改为Map<String, Long>(key为文件路径,value为最后处理的修改时间),调整过滤逻辑:private final Map<String, Long> pathsAlreadyProcessed = new HashMap<>(); private void processDiscoveredSplits(Collection<FileSourceSplit> splits, Throwable error) { if (error != null) { LOG.error("Failed to enumerate files", error); return; } final Collection<FileSourceSplit> newSplits = splits.stream() .filter(split -> { Long lastProcessedTime = pathsAlreadyProcessed.get(split.path()); // 路径未处理过,或当前文件版本更新则视为新Split if (lastProcessedTime == null || split.fileModificationTime() > lastProcessedTime) { pathsAlreadyProcessed.put(split.path(), split.fileModificationTime()); return true; } return false; }) .collect(Collectors.toList()); splitAssigner.addSplits(newSplits); assignSplits(); } - 步骤2:注入自定义枚举器
在构建FileSource时,通过DiscoverySettings注入自定义的枚举器提供者:FileSource<String> source = FileSource .forRecordStreamFormat(new TextLineInputFormat(), new Path("/your/target/dir")) .monitorContinuously(Duration.ofHours(1)) // 匹配外部系统的文件生成间隔 .discoverySettings(settings -> settings.setEnumeratorProvider((context, assigner, discoverySettings, checkpoint, processedPaths) -> new CustomContinuousFileSplitEnumerator( context, assigner, discoverySettings, checkpoint, processedPaths.isEmpty() ? new HashMap<>() : processedPaths ) ) ) .build();
方案2:升级到Flink 1.19+版本
Flink 1.19及后续版本已经修复了这个问题,ContinuousFileSplitEnumerator的默认逻辑会同时结合文件路径和修改时间判断是否为新文件,升级后无需自定义代码即可解决同名文件重复读取的问题。
2. pathsAlreadyProcessed的内存占用与清理问题
默认的ContinuousFileSplitEnumerator中,pathsAlreadyProcessed是一个HashSet<String>,会被完整写入检查点,且不会自动清空:
- 对于同名文件场景:默认逻辑只会存储一次路径,内存占用不会无限增长;
- 对于大量不同路径的文件:如果新路径持续新增,该集合会不断膨胀,导致检查点体积增大、JobManager内存占用上升。
解决建议:
- 优化存储结构:
如方案1所示,改用Map<String, Long>存储,同一个路径仅保留最新修改时间,避免无效的重复路径占用空间; - 添加定时清理逻辑:
在自定义枚举器中添加定时任务,比如移除超过指定时长(如7天)未更新的路径记录,减少内存占用; - 结合读后删除逻辑清理:
如果开启了读后删除,可以在确认文件已被成功处理并删除后,从pathsAlreadyProcessed中移除对应的路径,但需要注意并发安全,确保只有文件处理完成后才执行移除操作。
备注:内容来源于stack exchange,提问作者Dmitry

