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

Spring SFTP集成:如何实现下载与处理并行,及时触发ServiceActivator

Spring SFTP入站集成多线程处理阻塞问题解决

问题背景

无法在入站Spring SFTP集成中实现多线程处理,期望轮询器获取到文件后立即在单独线程触发ServiceActivator,但当前代码中ServiceActivator需等待messageSource处理完所有文件才执行,当SFTP目录存在2GB以上大文件时,等待时间过长。

初始代码

DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory();
factory.setHost(sftpHost);
factory.setPort(sftpPort);
factory.setAllowUnknownKeys(true);
factory.setUser(sftpUser);
factory.setPassword(sftpPassword);
factory.setAllowUnknownKeys(true);
return factory;
}

@Bean(name="defaultsync")
public SftpInboundFileSynchronizer synchronizer(){
    SftpInboundFileSynchronizer sync = new SftpInboundFileSynchronizer(createNewSFTPSessionFactory());
    sync.setDeleteRemoteFiles(false);
    sync.setRemoteDirectory(sftpDirectory);
    sync.setFilter(new SftpSimplePatternFileListFilter("*.xml"));
    return sync;
}

@Bean(name="sftpMessageSource")
@InboundChannelAdapter(channel="fileuploaded", poller = @Poller(fixedDelay = "3000"))
public MessageSource<File> sftpMessageSource(){
    SftpInboundFileSynchronizingMessageSource source =
            new SftpInboundFileSynchronizingMessageSource(synchronizer());
    source.setLocalDirectory(new File("C:\\Users\\Administrator\\Documents\\myfolder"));
    source.setAutoCreateLocalDirectory(true);
    source.setMaxFetchSize(-1);
    return source;
}


@ServiceActivator(inputChannel = "fileuploaded")
public void handleIncomingFile(File file) throws IOException {
    log.info(String.format("handleIncomingFile BEGIN %s", file.getName()));
    String content = FileUtils.readFileToString(file, "UTF-8");
    log.info(String.format("Content: %s", content));
    if(awsFileService.putFileInAWSS3(file))
    System.out.println("Processed file : "+file.getName());

} 

排查发现的核心问题

  • maxFetchSize设为负数,导致必须处理完输入通道中所有文件才会触发ServiceActivator
  • 缺少防重复处理过滤器,同一文件会被重复轮询处理
  • 当文件数量多、体积大时,SimpleMetadataStore方案性能不足,需改用基于JDBC的持久化存储过滤器

期间出现的错误

Caused by: org.springframework.jdbc.UncategorizedSQLException: PreparedStatementCallback; uncategorized SQLException for SQL [INSERT INTO INT_METADATA_STORE(METADATA_KEY, METADATA_VALUE, REGION) SELECT ?, ?, ? FROM INT_METADATA_STORE WHERE METADATA_KEY=? AND REGION=? HAVING COUNT(*)=0]; SQL state [S0002]; error code [208]; Invalid object name 'INT_METADATA_STORE'

最终解决方案

1. 配置默认Spring数据源

确保项目中已配置好可用的Spring DataSource(通过application.properties/yaml配置数据库连接信息)。

2. 修改SFTP同步器与过滤器代码

@Autowired
DataSource dataSource;

@Bean(name="defaultsync")
public SftpInboundFileSynchronizer synchronizer(){
    SftpInboundFileSynchronizer sync = new SftpInboundFileSynchronizer(createNewSFTPSessionFactory());
    sync.setDeleteRemoteFiles(false);
    sync.setRemoteDirectory(sftpDirectory);
    // sync.setFilter(new SftpSimplePatternFileListFilter("*.xml"));
    sync.setFilter(getFilter());
    return sync;
}

public ChainFileListFilter<DirEntry> getFilter()
{
    ChainFileListFilter<DirEntry> chainFilter = new ChainFileListFilter<>();
    chainFilter.addFilter(new SftpRegexPatternFileListFilter(".*\\.(xml|XML)"));
    chainFilter.addFilter(new SftpPersistentAcceptOnceFileListFilter(getMetadataStore(),"MyPrefix"));
    return chainFilter;
}
@Bean
public ConcurrentMetadataStore getMetadataStore()
{
    return new JdbcMetadataStore(dataSource);
}

3. 初始化数据库表

需使用Spring Integration JDBC模块提供的对应数据库初始化脚本,创建INT_METADATA_STORE表,脚本可从Spring Integration官方代码仓库的JDBC资源目录获取。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 14:47:31