多副本场景下Spring Integration SFTP文件摄取的正确配置
问题现象
运行多实例Spring Boot应用时,出现SFTP文件已被其他实例删除后,当前实例仍尝试读取的错误,错误栈如下:
org.springframework.messaging.MessagingException at org.springframework.integration.endpoint.AbstractPollingEndpoint.pollForMessage(AbstractPollingEndpoint.java:433) ... Caused by: java.io.UncheckedIOException: IOException when retrieving /to_erp/orders/49230084_COMPLETED_ORDER.json ... Caused by: SFTP error (SSH_FX_NO_SUCH_FILE): The file does not exist. ...
已配置SftpPersistentAcceptOnceFileListFilter配合JdbcMetadataStore实现多实例互斥,但仍出现重复处理问题。
当前配置代码
SFTP订阅者配置
@Profile("!test") @Component public class CreatedOrderSubscriber { private static final String REMOTE_DIRECTORY = "/to_wms/orders/"; private static final String FILE_PREFIX = "created_order_stream_"; private static final int MAX_FETCH_SIZE = 50; private final ConcurrentMetadataStore metadataStore; private final OilOrderCreatedService oilOrderCreatedService; private final SessionFactory<DirEntry> sftpSessionFactory; private final AdviceFactory adviceChainFactory; public CreatedOrderSubscriber( ConcurrentMetadataStore metadataStore, AdviceFactory adviceChainFactory, OilOrderCreatedService oilOrderCreatedService, @Qualifier("wmsSftpSessionFactory") SessionFactory<DirEntry> sftpSessionFactory) { this.metadataStore = metadataStore; this.oilOrderCreatedService = oilOrderCreatedService; this.sftpSessionFactory = sftpSessionFactory; this.adviceChainFactory = adviceChainFactory; } @Bean @ServiceActivator(inputChannel = "createdOrderV8Data", adviceChain = "afterCreatedOrder") public MessageHandler handleCreatedOrder() { var callback = createMessageHandlerCallback(); return MessageHandlerFactory.createHandler(callback, OrderUpdateHeaderDto.class); } @Bean @InboundChannelAdapter(channel = "createdOrderV8Stream", poller = @Poller(fixedDelay = "1000", maxMessagesPerPoll = "10")) public MessageSource<InputStream> createdOrderFtpMessageSource(SftpRemoteFileTemplate createdOrderV8Template) { SftpStreamingMessageSource messageSource = new SftpStreamingMessageSource(createdOrderV8Template); messageSource.setRemoteDirectory(REMOTE_DIRECTORY); messageSource.setFilter(fileListFilter()); messageSource.setMaxFetchSize(MAX_FETCH_SIZE); return messageSource; } @Bean @Transformer(inputChannel = "createdOrderV8Stream", outputChannel = "createdOrderV8Data") public org.springframework.integration.transformer.Transformer createdOrderTransformer() { return new StreamTransformer("UTF-8"); } @Bean public ExpressionEvaluatingRequestHandlerAdvice afterCreatedOrder() { return adviceChainFactory.createAdviceWithRemovalUponSuccess("createdOrderV8Template"); } @Bean public SftpRemoteFileTemplate createdOrderV8Template() { return new SftpRemoteFileTemplate(sftpSessionFactory); } private Function<OrderUpdateHeaderDto, Void> createMessageHandlerCallback() { return dto -> { oilOrderCreatedService.handleMessage(dto); return null; }; } private ChainFileListFilter<DirEntry> fileListFilter() { var sftpPersistentAcceptOnceFilter = new SftpPersistentAcceptOnceFileListFilter(metadataStore, FILE_PREFIX); var nameFilter = new SftpSimplePatternFileListFilter("*_NEW_ORDER.json"); var filters = Set.of(nameFilter, sftpPersistentAcceptOnceFilter); return new ChainFileListFilter<>(filters); } }
元数据存储配置
@Profile("!test") @Configuration public class MetadataStoreConfig { @Bean public ConcurrentMetadataStore metadataStore(DataSource dataSource) { return new JdbcMetadataStore(dataSource); } }
问题原因与修复方案
1. 过滤器链顺序错误
当前过滤器链先执行名称过滤,再执行持久化互斥过滤。多实例并发时,多个实例会同时拿到符合名称规则的文件列表,之后才去检查元数据存储,此时可能因为并发导致多个实例都通过过滤,进而重复处理文件。
修复:调整过滤器顺序,先执行持久化互斥过滤
修改fileListFilter方法,将SftpPersistentAcceptOnceFileListFilter放在过滤器链的最前面:
private ChainFileListFilter<DirEntry> fileListFilter() { var sftpPersistentAcceptOnceFilter = new SftpPersistentAcceptOnceFileListFilter(metadataStore, FILE_PREFIX); var nameFilter = new SftpSimplePatternFileListFilter("*_NEW_ORDER.json"); // 先执行持久化锁检查,再过滤文件名 return new ChainFileListFilter<>(List.of(sftpPersistentAcceptOnceFilter, nameFilter)); }
2. 添加锁过期时间(可选但推荐)
若某个实例在标记文件为已处理后异常崩溃,可能导致文件被永久锁定。可通过设置lockTtl让锁自动过期:
sftpPersistentAcceptOnceFilter.setLockTtl(60000); // 锁1分钟后自动过期
3. 优化轮询与拉取参数
当前maxFetchSize=50、maxMessagesPerPoll=10,一次拉取过多文件会增加并发冲突概率。可适当调小maxFetchSize,或设置setUseWatchService(false)(无需实时监听文件变化时):
SftpStreamingMessageSource messageSource = new SftpStreamingMessageSource(createdOrderV8Template); messageSource.setRemoteDirectory(REMOTE_DIRECTORY); messageSource.setFilter(fileListFilter()); messageSource.setMaxFetchSize(10); // 调小拉取数量 messageSource.setUseWatchService(false); // 关闭watch服务,减少资源消耗 return messageSource;
4. 确保元数据存储的原子性
JdbcMetadataStore默认支持原子性的putIfAbsent操作,这是SftpPersistentAcceptOnceFileListFilter实现互斥的核心。需确保:
- 数据源配置正确,支持事务
- 数据库中
INT_METADATA_STORE表已创建(默认会自动生成,若未生成可手动创建):
CREATE TABLE INT_METADATA_STORE ( METADATA_KEY VARCHAR(255) NOT NULL PRIMARY KEY, METADATA_VALUE VARCHAR(255), REGION VARCHAR(100) NOT NULL DEFAULT '' );
总结
通过调整过滤器链顺序,确保先进行持久化锁检查,结合合理的轮询参数与锁过期配置,即可实现多实例Spring Integration SFTP订阅者的互不干扰,避免重复处理文件导致的错误。
内容的提问来源于stack exchange,提问作者FOR_SCIENCE

