如何在运行时重新配置Spring Integration的SFTP相关Bean?
动态刷新SFTP配置的解决方案
我有如下Spring Integration SFTP配置,该配置从数据库加载SFTP服务器信息,提供了用于上传文件的RemoteFileTemplate,以及轮询多台SFTP服务器拉取文件的IntegrationFlow:
@Configuration @EnableIntegration public class SftpConfiguration { @Autowired private InterfaceRepository interfaceRepo; public record SessionFactoryKey(String host, int port, String user) { } @Bean SessionFactoryLocator<LsEntry> sessionFactoryLocator() { Map<Object, SessionFactory<LsEntry>> factories = interfaceRepo.findAll().stream() .map(x -> new SimpleEntry<>(new SessionFactoryKey(x.getHostname(), x.getPort(), x.getUsername()), sessionFactory(x.getHostname(), x.getPort(), x.getUsername(), x.getPassword()))) .collect(Collectors.toMap(Entry::getKey, Entry::getValue, (a, b) -> a)); return new DefaultSessionFactoryLocator<>(factories); } @Bean RemoteFileTemplate<LsEntry> fileTemplateResolver(DelegatingSessionFactory<LsEntry> delegatingSessionFactory) { return new SftpRemoteFileTemplate(delegatingSessionFactory); } @Bean DelegatingSessionFactory<LsEntry> delegatingSessionFactory(SessionFactoryLocator<LsEntry> sessionFactoryLocator) { return new DelegatingSessionFactory<>(sessionFactoryLocator); } @Bean RotatingServerAdvice advice(DelegatingSessionFactory<LsEntry> delegatingSessionFactory) { List<RotationPolicy.KeyDirectory> keyDirectories = interfaceRepo.findAll().stream() .filter(Interface::isReceivingData) .map(x -> new RotationPolicy.KeyDirectory( new SessionFactoryKey(x.getHostname(), x.getPort(), x.getUsername()), x.getDirectory())) .toList(); return keyDirectories.isEmpty() ? null : new RotatingServerAdvice(delegatingSessionFactory, keyDirectories); } @Bean PropertiesPersistingMetadataStore store() { return new PropertiesPersistingMetadataStore(); } @Bean public IntegrationFlow flow(ObjectProvider<RotatingServerAdvice> adviceProvider, DelegatingSessionFactory<LsEntry> delegatingSessionFactory, PropertiesPersistingMetadataStore store) { RotatingServerAdvice advice = adviceProvider.getIfAvailable(); return advice == null ? null : IntegrationFlows .from(Sftp.inboundAdapter(delegatingSessionFactory) .filter(new SftpPersistentAcceptOnceFileListFilter(store, "rotate_")) .localDirectory(new File("C:\\tmp\\sftp")) .localFilenameExpression("#remoteDirectory + T(java.io.File).separator + #root") .remoteDirectory("."), e -> e.poller(Pollers.fixedDelay(1).advice(advice))) .channel(MessageChannels.queue("files")).get(); } private SessionFactory<LsEntry> sessionFactory(String host, int port, String user, String password) { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); factory.setHost(host); factory.setPort(port); factory.setUser(user); factory.setPassword(password); factory.setAllowUnknownKeys(true); return factory; } }
我希望数据库中的SFTP配置发生变更时,能重新加载相关Bean,但尝试了Spring Cloud的@RefreshScope,由于IntegrationFlow仅支持singleton作用域而失效。请问除了重启应用上下文(重新调用SpringApplication.run)之外,还有其他可行的解决方案吗?
核心思路:避免直接刷新Bean,动态更新内部配置
因为IntegrationFlow和相关会话工厂组件是单例,无法通过Scope机制刷新,所以重点是让这些组件的内部配置可以动态更新,而非重建Bean。
方案1:自定义可刷新的SessionFactoryLocator
- 改造
SessionFactoryLocator实现,支持动态更新SFTP会话工厂映射:
public class RefreshableSessionFactoryLocator extends DefaultSessionFactoryLocator<LsEntry> { private final InterfaceRepository interfaceRepo; public RefreshableSessionFactoryLocator(InterfaceRepository interfaceRepo) { super(Collections.emptyMap()); this.interfaceRepo = interfaceRepo; // 初始化加载配置 refreshFactories(); } // 暴露刷新方法 public void refreshFactories() { Map<Object, SessionFactory<LsEntry>> newFactories = interfaceRepo.findAll().stream() .map(x -> new AbstractMap.SimpleEntry<>(new SessionFactoryKey(x.getHostname(), x.getPort(), x.getUsername()), createSessionFactory(x.getHostname(), x.getPort(), x.getUsername(), x.getPassword()))) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue, (a, b) -> a)); // 更新父类的会话工厂映射 super.setSessionFactories(newFactories); } private SessionFactory<LsEntry> createSessionFactory(String host, int port, String user, String password) { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); factory.setHost(host); factory.setPort(port); factory.setUser(user); factory.setPassword(password); factory.setAllowUnknownKeys(true); return factory; } }
- 在配置类中注册自定义Locator,并添加刷新端点:
@Bean RefreshableSessionFactoryLocator sessionFactoryLocator(InterfaceRepository interfaceRepo) { return new RefreshableSessionFactoryLocator(interfaceRepo); } // 自定义刷新接口(也可通过Spring Boot Actuator实现) @RestController @RequestMapping("/sftp") public class SftpRefreshController { private final RefreshableSessionFactoryLocator sessionFactoryLocator; private final RotatingServerAdvice rotatingServerAdvice; private final InterfaceRepository interfaceRepo; public SftpRefreshController(RefreshableSessionFactoryLocator sessionFactoryLocator, @Autowired(required = false) RotatingServerAdvice rotatingServerAdvice, InterfaceRepository interfaceRepo) { this.sessionFactoryLocator = sessionFactoryLocator; this.rotatingServerAdvice = rotatingServerAdvice; this.interfaceRepo = interfaceRepo; } @PostMapping("/refresh") public void refreshSftpConfig() { // 刷新会话工厂 sessionFactoryLocator.refreshFactories(); // 刷新RotatingServerAdvice的密钥目录 if (rotatingServerAdvice != null) { List<RotationPolicy.KeyDirectory> newKeyDirectories = interfaceRepo.findAll().stream() .filter(Interface::isReceivingData) .map(x -> new RotationPolicy.KeyDirectory( new SessionFactoryKey(x.getHostname(), x.getPort(), x.getUsername()), x.getDirectory())) .toList(); // 通过反射修改内部字段(也可改造RotatingServerAdvice使其支持刷新) try { Field field = RotatingServerAdvice.class.getDeclaredField("keyDirectories"); field.setAccessible(true); field.set(rotatingServerAdvice, newKeyDirectories); } catch (NoSuchFieldException | IllegalAccessException e) { throw new RuntimeException("刷新RotatingServerAdvice失败", e); } } } }
方案2:利用IntegrationFlowContext动态管理Flow
如果需要完全重建IntegrationFlow,可以使用Spring Integration的IntegrationFlowContext实现动态注册/销毁Flow:
- 改造配置类,通过上下文管理Flow:
@Configuration @EnableIntegration public class SftpConfiguration { @Autowired private InterfaceRepository interfaceRepo; @Autowired private IntegrationFlowContext flowContext; @Autowired private DelegatingSessionFactory<LsEntry> delegatingSessionFactory; @Autowired private PropertiesPersistingMetadataStore store; private static final String SFTP_FLOW_ID = "sftp-inbound-flow"; @PostConstruct public void initFlow() { registerOrRefreshFlow(); } // 暴露刷新Flow的方法 public void registerOrRefreshFlow() { // 先销毁旧Flow flowContext.remove(SFTP_FLOW_ID); // 加载新配置创建Flow List<RotationPolicy.KeyDirectory> keyDirectories = interfaceRepo.findAll().stream() .filter(Interface::isReceivingData) .map(x -> new RotationPolicy.KeyDirectory( new SessionFactoryKey(x.getHostname(), x.getPort(), x.getUsername()), x.getDirectory())) .toList(); if (!keyDirectories.isEmpty()) { RotatingServerAdvice advice = new RotatingServerAdvice(delegatingSessionFactory, keyDirectories); IntegrationFlow flow = IntegrationFlows .from(Sftp.inboundAdapter(delegatingSessionFactory) .filter(new SftpPersistentAcceptOnceFileListFilter(store, "rotate_")) .localDirectory(new File("C:\\tmp\\sftp")) .localFilenameExpression("#remoteDirectory + T(java.io.File).separator + #root") .remoteDirectory("."), e -> e.poller(Pollers.fixedDelay(1).advice(advice))) .channel(MessageChannels.queue("files")) .get(); // 注册新Flow flowContext.register(flow, registration -> registration.id(SFTP_FLOW_ID)); } } // 其他Bean配置(SessionFactoryLocator、DelegatingSessionFactory等保持不变) }
- 在刷新接口中调用Flow重建方法:
@RestController @RequestMapping("/sftp") public class SftpRefreshController { private final SftpConfiguration sftpConfiguration; private final RefreshableSessionFactoryLocator sessionFactoryLocator; public SftpRefreshController(SftpConfiguration sftpConfiguration, RefreshableSessionFactoryLocator sessionFactoryLocator) { this.sftpConfiguration = sftpConfiguration; this.sessionFactoryLocator = sessionFactoryLocator; } @PostMapping("/refresh") public void refreshSftpConfig() { // 先刷新会话工厂 sessionFactoryLocator.refreshFactories(); // 重建IntegrationFlow sftpConfiguration.registerOrRefreshFlow(); } }
注意事项
- 方案1更轻量,无需重建Flow,适合配置变更不频繁的场景;方案2更彻底,适合需要完全重置Flow逻辑的场景。
- 触发刷新的方式:可以通过数据库事件监听(如PostgreSQL的LISTEN/NOTIFY)、定时轮询配置版本号,或者外部API调用刷新端点实现。
- 刷新过程需保证线程安全,可添加锁控制避免并发刷新导致的异常。
内容的提问来源于stack exchange,提问作者Guerric P
相关产品推荐
相关产品推荐

