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

Spring Integration FileReadingMessageSource停止拾取文件问题排查

问题排查与解决方案

针对你在AWS EKS容器中运行Spring Integration 2.7.3 Job时出现的FileReadingMessageSource停止拾取文件、重启容器后恢复的问题,结合你的配置,主要问题点和修复方案如下:

1. 内存型文件过滤器的局限性

你使用的AcceptOnceFileListFilter是内存级别的过滤器,所有已处理的文件路径会被保存在进程内存中:

  • 当Job长时间运行,处理的文件数量持续增加,缓存会无限膨胀,可能导致内存占用过高、扫描性能急剧下降,最终表现为停止拾取新文件。
  • 容器重启后内存缓存清空,之前的文件会被重新识别并处理,这完全符合你描述的现象。

修复方案:
替换为持久化版本的PersistentAcceptOnceFileListFilter,它会将已处理文件的记录持久化到本地文件或外部存储中,既避免内存溢出,也能在容器重启后保留处理记录,防止重复处理:

@Bean
public DirectoryScanner directoryScanner() {
    // 使用基于文件的元数据存储,确保重启后记录不丢失
    MetadataStore metadataStore = new PropertiesPersistingMetadataStore();
    ((PropertiesPersistingMetadataStore) metadataStore).setBaseDirectory("/path/to/persist"); // 指定持久化目录

    FileListFilter<File> filter = new CompositeFileListFilter<>(
            Arrays.asList(new PersistentAcceptOnceFileListFilter(metadataStore, "file-reader-"), 
                          new LastModifiedFileListFilter()));

    var scanner = new RecursiveDirectoryScanner();
    scanner.setFilter(filter);
    return scanner;
}

2. 伪事务配置的潜在问题

你的Poller配置了pseudoTransactionManager(伪事务管理器)和transactionSynchronizationFactory:
伪事务管理器不会真正执行事务的提交/回滚操作,如果事务同步工厂包含与过滤器状态更新相关的逻辑(比如标记文件为已处理),可能导致状态更新不生效,或者出现异常后状态不一致,最终引发扫描停滞。

修复方案:

  • 如果你的流程不需要事务支持,直接移除transactional(pseudoTransactionManager)和transactionSynchronizationFactory配置:
@Bean
public PollerSpec pollerSpec() {
    return Pollers.fixedDelay(Duration.ofMinutes(5))
                .maxMessagesPerPoll(100);
}
  • 如果确实需要事务,改用真实的事务管理器(比如基于JDBC的DataSourceTransactionManager),确保事务同步逻辑能正确执行。

3. 不规范的Bean注入方式

在startFlow的Bean定义中,你直接调用fileReadingMessageSource(null)来获取实例,这种写法可能导致创建多个FileReadingMessageSource实例,引发上下文不一致或缓存混乱。

修复方案:
通过依赖注入直接获取容器中的FileReadingMessageSource Bean:

@Bean
public IntegrationFlow startFlow(FileReadingMessageSource fileReadingMessageSource) {
    return IntegrationFlows.from(fileReadingMessageSource, p -> p.poller(pollerSpec()))
                        .enrichHeaders(Collections.singletonMap(ERROR_CHANNEL, appErrorChannel))
                        .log(DEBUG, message -> "File Picked: " + message.getHeaders().get(FILENAME))
                        // 后续逻辑
                        ;
}

4. 目录扫描的异常处理缺失

RecursiveDirectoryScanner在扫描目录时如果遇到IO异常(比如EKS存储卷临时故障、权限变化、文件被意外删除等),默认情况下可能导致扫描中断,甚至让后续的轮询任务停止执行。

修复方案:
给扫描器添加异常处理器,确保扫描异常不会影响后续轮询:

@Bean
public DirectoryScanner directoryScanner() {
    // ... 过滤器配置

    var scanner = new RecursiveDirectoryScanner();
    scanner.setFilter(filter);
    scanner.setErrorHandler(throwable -> {
        // 记录异常日志,根据需求添加告警逻辑
        log.error("Directory scan failed, will retry in next poll", throwable);
    });
    return scanner;
}

5. 存储卷的一致性问题(EKS环境特有)

EKS中使用的存储卷(比如EFS)可能存在元数据缓存延迟,导致FileReadingMessageSource无法及时感知到新文件。

修复方案:
设置扫描器的目录元数据刷新间隔,强制定期刷新目录信息:

var scanner = new RecursiveDirectoryScanner();
scanner.setFilter(filter);
scanner.setRefreshInterval(Duration.ofMinutes(1)); // 每隔1分钟刷新一次目录元数据

内容的提问来源于stack exchange,提问作者Rayyan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 01:15:40