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

