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

如何在运行时重新配置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

  1. 改造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;
    }
}
  1. 在配置类中注册自定义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:

  1. 改造配置类,通过上下文管理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等保持不变)
}
  1. 在刷新接口中调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 22:45:01