已弃用EmitterProcessor,Spring Cloud Stream+RabbitMQ如何用Sinks.Many替换?
问题:替换Spring Cloud Stream中弃用的EmitterProcessor为Sinks.Many实现SSE
现有代码(使用已弃用的EmitterProcessor)
我当前用EmitterProcessor实现SSE接口,但该类已被弃用,代码如下:
@RestController @Slf4j public class ResourceRest { private final EmitterProcessor<String> streamProcessor = EmitterProcessor.create(); @GetMapping(value = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> getSee() { return this.streamProcessor; } @Bean public Consumer<Flux<String>> receiveSse() { return recordFlux -> recordFlux .doOnNext(this.streamProcessor::onNext) .doOnNext(value -> log.info("*" +value)) .subscribe(); } }
已知替换方案
了解到可以用Sinks.Many替代:
Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();
项目背景
基于Spring Cloud Stream + RabbitMQ的数据流:
- producer-service发布数据到RabbitMQ
- processor-service处理数据后封装为
AccumulatorMessage,发布到指定队列 - consumer-service消费该队列的消息
当前consumer-service配置
配置类
@Configuration public class ConsumerConfiguration { @Value("${spring.cloud.stream.function.definition}") private String definition; @Bean public Consumer<AccumulatorMessage> receive() { return payload -> { log.info(payload.toString()); }; } }
application.yml
spring: rabbitmq: host: localhost port: 5673 username: guest password: guest application: name: consumer-service cloud: stream: function: definition: receive; receiveSse bindings: receive-in-0: destination: accumulatorMessageChannel group: accumulatorMessageChannelGroup receiveSse-in-0: destination: accumulatorMessageChannel group: accumulatorMessageChannelGroup server: port: 8030
疑问
已将receiveSse()添加到绑定配置,但不清楚如何将Sinks.Many方案应用到ResourceRest类中,也不了解其具体工作机制,求解答。
解决方案与原理说明
1. 改造ResourceRest类(替换为Sinks.Many)
直接替换EmitterProcessor为Sinks.Many,并调整消息推送逻辑,改造后的代码如下:
@RestController @Slf4j public class ResourceRest { // 初始化Sinks.Many:unicast适用于单客户端订阅SSE的场景,缓存积压消息 private final Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer(); @GetMapping(value = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> getSse() { // 将Sink转换为Flux,供SSE客户端订阅长连接 return sink.asFlux(); } @Bean public Consumer<AccumulatorMessage> receiveSse() { // 直接接收队列中的AccumulatorMessage,Spring Cloud Stream会自动映射消息 return message -> { // 根据需求转换消息为String(比如提取字段或直接toString) String sseContent = message.toString(); log.info("*" + sseContent); // 非阻塞推送消息到Sink,处理推送失败场景 Sinks.EmitResult emitResult = sink.tryEmitNext(sseContent); if (!emitResult.isSuccess()) { log.warn("推送SSE消息失败: {}", emitResult.getFailureReason()); } }; } }
2. 配置调整说明
你的application.yml配置无需额外修改,注意两点即可:
function.definition中的receiveSse必须和ResourceRest中定义的ConsumerBean名称一致receiveSse-in-0的绑定配置正确,和receive共用同一个队列和消费组,确保能获取到AccumulatorMessage
3. 核心工作机制
Sinks.Many的作用
Sinks是Reactor官方替代旧版Processor的组件,用于手动向响应式流推送数据:
many():表示支持多个订阅者unicast():针对单个订阅者优化(SSE每个客户端是独立长连接,属于单订阅场景)onBackpressureBuffer():当客户端处理速度慢于消息推送速度时,自动缓存消息,避免丢失
整体数据流
- RabbitMQ队列中的
AccumulatorMessage被receiveSseConsumer接收 - 消息转换为String后,通过
sink.tryEmitNext推送到Sinks组件 - SSE客户端访问
/sse接口时,getSse返回sink.asFlux(),建立长连接 - Sinks有新消息时,Flux会自动以SSE格式(text/event-stream)推送给客户端
- 客户端实时接收消息并处理
4. 关键注意事项
- 消息类型匹配:原代码中
Consumer<Flux<String>>是错误的,Spring Cloud Stream的ConsumerBean应直接接收单个消息对象(AccumulatorMessage),框架会自动处理批量消息的分发 - 非阻塞推送:优先使用
tryEmitNext而非emitNext,前者是非阻塞的,不会因为客户端断开连接导致线程阻塞 - 多订阅场景:如果需要多个客户端共享同一份SSE流,改用
Sinks.many().multicast().onBackpressureBuffer(),并注意订阅者的生命周期管理
内容的提问来源于stack exchange,提问作者skyho
相关产品推荐
相关产品推荐

