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
相关产品推荐
相关产品推荐

