You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()中被硬编码,起初我以为无法替换它;且官方文档提到扫描目录时会参考文件最后修改时间,但该逻辑并未在拆分枚举器中生效。

我的疑问

  1. 如何让过滤逻辑结合fileModificationTime字段判断?是否有现成的实现方案?
  2. 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. 优化存储结构:
    如方案1所示,改用Map<String, Long>存储,同一个路径仅保留最新修改时间,避免无效的重复路径占用空间;
  2. 添加定时清理逻辑:
    在自定义枚举器中添加定时任务,比如移除超过指定时长(如7天)未更新的路径记录,减少内存占用;
  3. 结合读后删除逻辑清理:
    如果开启了读后删除,可以在确认文件已被成功处理并删除后,从pathsAlreadyProcessed中移除对应的路径,但需要注意并发安全,确保只有文件处理完成后才执行移除操作。

备注:内容来源于stack exchange,提问作者Dmitry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 09:34:34