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

多副本场景下Spring Integration SFTP文件摄取的正确配置

多实例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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 04:09:55