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

