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

已弃用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中定义的Consumer Bean名称一致
  • receiveSse-in-0的绑定配置正确,和receive共用同一个队列和消费组,确保能获取到AccumulatorMessage

3. 核心工作机制

Sinks.Many的作用

Sinks是Reactor官方替代旧版Processor的组件,用于手动向响应式流推送数据:

  • many():表示支持多个订阅者
  • unicast():针对单个订阅者优化(SSE每个客户端是独立长连接,属于单订阅场景)
  • onBackpressureBuffer():当客户端处理速度慢于消息推送速度时,自动缓存消息,避免丢失

整体数据流

  1. RabbitMQ队列中的AccumulatorMessage被receiveSse Consumer接收
  2. 消息转换为String后,通过sink.tryEmitNext推送到Sinks组件
  3. SSE客户端访问/sse接口时,getSse返回sink.asFlux(),建立长连接
  4. Sinks有新消息时,Flux会自动以SSE格式(text/event-stream)推送给客户端
  5. 客户端实时接收消息并处理

4. 关键注意事项

  • 消息类型匹配:原代码中Consumer<Flux<String>>是错误的,Spring Cloud Stream的Consumer Bean应直接接收单个消息对象(AccumulatorMessage),框架会自动处理批量消息的分发
  • 非阻塞推送:优先使用tryEmitNext而非emitNext,前者是非阻塞的,不会因为客户端断开连接导致线程阻塞
  • 多订阅场景:如果需要多个客户端共享同一份SSE流,改用Sinks.many().multicast().onBackpressureBuffer(),并注意订阅者的生命周期管理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 11:50:34