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

如何在Spring Integration中结合Webflux消费SSE并转换为独立Message实例

实现方案

我们可以通过继承MessageProducerSupport来实现符合Spring Integration生命周期规范的SSE消费组件,自动处理启停、资源释放、消息发送逻辑。

核心组件实现

import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.endpoint.MessageProducerSupport;
import org.springframework.web.reactive.function.client.WebClient;
import org.springframework.http.codec.ServerSentEvent;
import org.springframework.core.ParameterizedTypeReference;
import reactor.core.Disposable;
import reactor.util.retry.Retry;
import java.time.Duration;

@Component
public class SseMessageProducer extends MessageProducerSupport {

    private final WebClient webClient;
    private final String sseBaseUrl = "http://myhost:8080/sse";
    private final String sseEndpoint = "/stream-sse";
    private Disposable subscription;

    // 支持注入自定义WebClient.Builder适配签名校验、代理等特殊场景
    public SseMessageProducer(WebClient.Builder webClientBuilder) {
        this.webClient = webClientBuilder.baseUrl(sseBaseUrl).build();
        // 可直接在此处指定默认输出通道,也可后续在配置类中动态绑定
        // setOutputChannelName("sseEventChannel");
    }

    @Override
    protected void doStart() {
        super.doStart();
        ParameterizedTypeReference<ServerSentEvent<String>> type = new ParameterizedTypeReference<ServerSentEvent<String>>() {};

        this.subscription = webClient.get()
                .uri(sseEndpoint)
                .retrieve()
                .bodyToFlux(type)
                // 可选配置:连接断开自动重试,最多重试5次,每次间隔2秒,可根据业务需求调整
                .retryWhen(Retry.fixedDelay(5, Duration.ofSeconds(2))
                        .doBeforeRetry(retrySignal -> logger.warn("SSE连接断开,第{}次重试", retrySignal.totalRetries() + 1)))
                .subscribe(
                        event -> {
                            // 将SSE原生属性全部存入消息头,payload放业务数据
                            Message<String> message = MessageBuilder.withPayload(event.data())
                                    .setHeader("sseId", event.id())
                                    .setHeader("sseEventType", event.event())
                                    .setHeader("sseRetry", event.retry())
                                    .setHeader("sseComments", event.comment())
                                    .build();
                            // 调用父类方法发送消息到绑定的输出通道
                            sendMessage(message);
                        },
                        error -> logger.error("SSE流消费异常", error),
                        () -> logger.info("SSE流已正常结束")
                );
    }

    @Override
    protected void doStop() {
        // 组件停止时主动取消订阅,释放连接资源
        if (this.subscription != null && !this.subscription.isDisposed()) {
            this.subscription.dispose();
        }
        super.doStop();
    }
}

组件配置与使用

  • 通道绑定:可通过setOutputChannel/setOutputChannelName方法指定消息发送的目标通道,也可通过配置类完成绑定和后续流转逻辑,示例如下:
@Configuration
public class SseIntegrationConfig {

    @Bean
    public MessageChannel sseEventChannel() {
        return new DirectChannel();
    }

    @Bean
    public SseMessageProducer sseMessageProducer(WebClient.Builder webClientBuilder) {
        SseMessageProducer producer = new SseMessageProducer(webClientBuilder);
        producer.setOutputChannel(sseEventChannel());
        return producer;
    }

    // 示例:SSE消息处理流程
    @Bean
    public IntegrationFlow sseProcessingFlow() {
        return IntegrationFlow.from("sseEventChannel")
                .handle(message -> {
                    String payload = (String) message.getPayload();
                    // 编写自定义业务处理逻辑
                    System.out.println("收到SSE业务数据:" + payload);
                })
                .get();
    }
}
  • 生命周期管理:组件会跟随Spring上下文自动启停,也可手动调用start()/stop()方法控制SSE连接的建立与断开,不会出现资源泄漏问题。
  • 类型适配:如果SSE返回结构化JSON数据,只需修改ParameterizedTypeReference的泛型类型即可自动反序列化为对应实体类,无需额外解析逻辑。

内容的提问来源于stack exchange,提问作者Kamil Roman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 16:36:04