Spring Integration递归读取远程SFTP子目录文件问题求助
Spring Integration SFTP 递归读取文件问题及优化需求
初始需求与代码
已基于Spring Integration的入站适配器实现从远程SFTP服务器单个目录获取文件的工作流程,现在需要实现递归读取远程父目录下所有子目录中的文件。初始代码如下:
@Bean public SessionFactory<SftpClient.DirEntry> sftpSessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); factory.setHost("localhost"); factory.setPort(port); factory.setUser("foo"); factory.setPassword("foo"); factory.setAllowUnknownKeys(true); factory.setTestSession(true); return new CachingSessionFactory<>(factory); } @Bean public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() { SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory()); fileSynchronizer.setDeleteRemoteFiles(false); fileSynchronizer.setRemoteDirectory("foo"); fileSynchronizer.setFilter(new SftpSimplePatternFileListFilter("*.xml")); return fileSynchronizer; } @Bean @InboundChannelAdapter(channel = "sftpChannel", poller = @Poller(fixedDelay = "5000")) public MessageSource<File> sftpMessageSource() { SftpInboundFileSynchronizingMessageSource source = new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer()); source.setLocalDirectory(new File("sftp-inbound")); source.setAutoCreateLocalDirectory(true); source.setLocalFilter(new AcceptOnceFileListFilter<File>()); source.setMaxFetchSize(1); return source; } @Bean @ServiceActivator(inputChannel = "sftpChannel") public MessageHandler handler() { return new MessageHandler() { @Override public void handleMessage(Message<?> message) throws MessagingException { System.out.println(message.getPayload()); } }; }
第一次修改后的问题
感谢及时回复。已根据建议代码修改,应用可启动,但文件未被拾取。修改后的代码片段如下:
CompositeFileListFilter<LsEntry> compositeFileListFilter = new CompositeFileListFilter<>(); SftpPersistentAcceptOnceFileListFilter fileListFilter = new SftpPersistentAcceptOnceFileListFilter( (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE"); if (Constants.APP1.equals(appName) || Constants.APP2.equals(appName)) { SftpRegexPatternFileListFilter regexPatternFileListFilter = new SftpRegexPatternFileListFilter(Pattern.compile("^IL.*")); compositeFileListFilter.addFilter(regexPatternFileListFilter); } compositeFileListFilter.addFilter(fileListFilter); return IntegrationFlows.fromSupplier( () -> sftpEnvironment.getSftpGLSIncomingDir(), // remote dir e -> e.autoStartup(true).poller(pollerMetada())) .handle( Sftp.outboundGateway(sftpSessionFactory(), Command.MGET, "payload") .options(Option.RECURSIVE) .filter(compositeFileListFilter) .fileExistsMode(FileExistsMode.IGNORE) .localDirectoryExpression("'/tmp/' + #remoteDirectory")) // re-create tree locally .split() .log() .get();
新实现的问题与待解决需求
改用新代码实现后,出现文件仅被部分处理的问题,例如10个文件中仅能处理5-6个,无法定位问题根源。此外还有两个待解决需求:
- 当前实现可读取远程子目录文件并存储到本地,但希望无需存储到本地,直接在
sftpChannel中处理这些文件; - 希望基于数据库实现去重机制,避免重复处理文件。
完整代码如下:
public class SFTPPollerService { @Bean public SessionFactory<LsEntry> sftpSessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); //code return factory; } //OLD code // @Bean // public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() { // SftpInboundFileSynchronizer fileSynchronizer = // new SftpInboundFileSynchronizer(sftpSessionFactory()); // fileSynchronizer.setDeleteRemoteFiles(sftpEnvironment.isDeleteRemoteFiles()); // fileSynchronizer.setRemoteDirectory(sftpEnvironment.getSftpGLSIncomingDir()); // fileSynchronizer.setPreserveTimestamp(true); // CompositeFileListFilter<LsEntry> compositeFileListFilter = new // CompositeFileListFilter<>(); // SftpPersistentAcceptOnceFileListFilter fileListFilter = // new SftpPersistentAcceptOnceFileListFilter( // (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE"); // if (Constants.app2.equals(appName) // || Constants.app1.equals(appName)) { // SftpRegexPatternFileListFilter regexPatternFileListFilter = // new SftpRegexPatternFileListFilter(Pattern.compile("*.txt")); // compositeFileListFilter.addFilter(regexPatternFileListFilter); // } // compositeFileListFilter.addFilter(fileListFilter); // fileSynchronizer.setFilter(compositeFileListFilter); // return fileSynchronizer; // } // // @Bean // @InboundChannelAdapter(channel = "sftpChannel", poller = @Poller("pollerMetada")) // public MessageSource<File> sftpMessageSource() { // SftpInboundFileSynchronizingMessageSource source = // new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer()); // source.setLocalDirectory(new File(sftpEnvironment.getSftpLocalDir())); // // source.setAutoCreateLocalDirectory(true); // // try { // source.setLocalFilter( // (FileSystemPersistentAcceptOnceFileListFilter) // context.getBean("filelistFilter")); // } catch (Exception e) { // LOG.error( // "Exception caught while setting local filter on // SftpInboundFileSynchronizingMessageSource", // e); // } // source.setMaxFetchSize(sftpEnvironment.getMaxFetchFileSize()); // // return source; // } //new Code @Bean public IntegrationFlow sftpInboundFlow() { CompositeFileListFilter<LsEntry> compositeFileListFilter = new CompositeFileListFilter<>(); SftpPersistentAcceptOnceFileListFilter fileListFilter = new SftpPersistentAcceptOnceFileListFilter( (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE"); if (Constants.app2.equals(appName) || Constants.app1.equals(appName)) { SftpRegexPatternFileListFilter regexPatternFileListFilter = new SftpRegexPatternFileListFilter(Pattern.compile("(subDir | *.txt)")); compositeFileListFilter.addFilter(regexPatternFileListFilter); } fileListFilter.setForRecursion(true); FileSystemPersistentAcceptOnceFileListFilter fileSystemPersistentAcceptOnceFileListFilter = (FileSystemPersistentAcceptOnceFileListFilter) context.getBean( "filelistFilter"); compositeFileListFilter.addFilter(fileListFilter); // IntegrationFlow ir = // IntegrationFlows.from( // Sftp.inboundAdapter(sftpSessionFactory()) // .preserveTimestamp(true) // .remoteDirectory(sftpEnvironment.getSftpGLSIncomingDir()) // .deleteRemoteFiles(sftpEnvironment.isDeleteRemoteFiles()) // .filter(compositeFileListFilter) // .autoCreateLocalDirectory(true) // .localDirectory(new File(sftpEnvironment.getSftpLocalDir())), // e -> e.autoStartup(true).poller(pollerMetada())) // .handle(handler()) // .get(); return IntegrationFlows.fromSupplier( () -> sftpEnvironment.getSftpGLSIncomingDir(), // remote dir e -> e.autoStartup(true).poller(pollerMetada())) .handle( Sftp.outboundGateway(sftpSessionFactory(), Command.MGET, "payload") .options(Option.RECURSIVE) .fileExistsMode(FileExistsMode.IGNORE) .regexFileNameFilter("(dsv[0-9]|.*.xml)") // .filter(compositeFileListFilter) .localDirectoryExpression("'user/localDir/test/'")) // .handle(handler()) // .patternFileNameFilter(".*\\.xml")) // re-create tree locally .split() .channel("sftpChannel") // .handle(handler()) .log() .get(); } @Bean public PollerMetadata pollerMetada() { PollerMetadata pm = new PollerMetadata(); ExpressionEvaluatingTransactionSynchronizationProcessor processor = new ExpressionEvaluatingTransactionSynchronizationProcessor(); ExpressionParser parser = new SpelExpressionParser(); Expression exp = parser.parseExpression("payload.delete()"); processor.setAfterRollbackExpression(exp); TransactionSynchronizationFactory tsf = new DefaultTransactionSynchronizationFactory(processor); pm.setTransactionSynchronizationFactory(tsf); List<Advice> advices = new ArrayList<>(); advices.add(compoundTriggerAdvice()); pm.setAdviceChain(advices); pm.setTrigger(compoundTrigger()); pm.setMaxMessagesPerPoll(sftpEnvironment.getMaxMessagesPerPoll()); return pm; } @Bean public CronTrigger cronTrigger() { if (LOG.isDebugEnabled()) { return new CronTrigger(sftpEnvironment.getPollerCronExpressionWhenDebugModeIsEnabled()); } else { return new CronTrigger(sftpEnvironment.getPollerCronExpression()); } } @Bean public PeriodicTrigger periodicTrigger() { return new PeriodicTrigger(sftpEnvironment.getPeriodicTriggerInMillis()); } @Bean public CompoundTrigger compoundTrigger() { return new CompoundTrigger(cronTrigger()); } @Bean public CompoundTriggerAdvice compoundTriggerAdvice() { return new CompoundTriggerAdvice(compoundTrigger(), periodicTrigger()); } @Bean public FileSystemPersistentAcceptOnceFileListFilter filelistFilter(MetadataStore datastore) { return new FileSystemPersistentAcceptOnceFileListFilter((JdbcMetadataStore) datastore, "INT"); } @Bean public PlatformTransactionManager transactionManager() { return new org.springframework.integration.transaction.PseudoTransactionManager(); } @Bean DataSource dataSource() throws SQLException { OracleDataSource dataSource = new OracleDataSource(); dataSource.setUser(databaseProperties.getOracleUsername()); dataSource.setPassword(databaseProperties.getOraclePassword()); dataSource.setURL(databaseProperties.getOracleUrl()); dataSource.setImplicitCachingEnabled(true); dataSource.setFastConnectionFailoverEnabled(true); return dataSource; } /** * Creates a {@link JdbcMetadataStore} for the de-duplication logic. * * <p>This method uses the "REGION" column of the metadatastore table to differentiate between * multiple apps. The value of the "REGION" column is set equal to the app-name. * * @return a JDBC metadata store * @throws SQLException in case an exception occurs during connection to SQL database */ @Bean public MetadataStore metadataStore() throws SQLException { JdbcMetadataStore jdbcMetadataStore = new JdbcMetadataStore(dataSource()); if (!Constants.app2.equals(appName)) { jdbcMetadataStore.setRegion(appName); } return jdbcMetadataStore; } @Bean @ServiceActivator(inputChannel = "sftpChannel") public MessageHandler handler() { return message -> { File file = (File) message.getPayload(); FileDto fileDto = new FileDto(file); fileHandler.handle(fileDto); LOG.info("controller is here "); try { if (sftpEnvironment.isDeleteLocalFiles()) { Files.deleteIfExists(Paths.get(file.toString())); } } catch (IOException e) { // TODO retry/report/handle gracefully LOG.error(String.format("MessageHandler had error message=%s", message), e); } }; } }
内容的提问来源于stack exchange,提问作者arpit singh
相关产品推荐
相关产品推荐

