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

Spring Boot中如何将RabbitMQ队列事件发布到响应式ServerSentEvents Flux?

问题:Server Sent Events(SSE)无法正常推送任务进度

业务场景

接收RabbitMQ队列中的任务运行消息,保存到数据库的同时,通过SSE向等待的客户端推送任务运行状态。

后端实现代码

Controller中的SSE接口

@GetMapping("/{id}/progress/stream")
public Flux<ServerSentEvent<JobProgress>> subscribeToProgress(@PathVariable("id") String id) {
  return Flux.merge(JdkFlowAdapter.flowPublisherToFlux(progressPublisher)
    .filter(progress -> progress.id().equals(id))
    .map(progress -> ServerSentEvent.<JobProgress>builder()
      .id(progress.id())
      .event("job-progress")
      .data(progress)
      .retry(Duration.ofSeconds(10))
      .build()
    ), 
    Flux.interval(Duration.ofMillis(500L))
      .map(sequence -> ServerSentEvent.<JobProgress>builder()
        .event("keep-alive")
        .comment("Keeping the connection alive to avoid abrupt closing.")
        .retry(Duration.ofSeconds(1))
        .build()))
  );
}

注:JobProgress是包含状态码和消息的实体类,事件通过java.util.concurrent.SubmissionPublisher progressPublisher发送

发布者Bean配置

@Configuration
public class PublisherConfiguration {
  @Bean
  public SubmissionPublisher<JobProgress> progressSink() {
    return new SubmissionPublisher();
  }
}

事件发布组件

@Component
public class JobEventPublisher {
  @Autowired
  private final SubmissionPublisher<JobProgress> progressPublisher;

  public void dispatchToClient(JobProgress progress) {
    progressPublisher.submit(progress);
  }
}

前端接收代码(TypeScript/React)

useEffect(() => {
  const id = "some-uuid-received-on-post";
  const source = new EventSource(`${baseUrl}/api/v1/job/${id}/progress/stream`, {withCredentials: true});

  // 此方法从未触发
  source.onmessage = console.log;

 return source.close;
}, [baseUrl]);

问题现象

  1. 后端调试日志显示事件已通过Flux传递并映射为ServerSentEvent,但网络面板中EventSource始终处于挂起状态,客户端从未收到消息
  2. 客户端与服务端之间存在Spring Boot Cloud Gateway,已排除WebSocket影响(仅需单向通信)
  3. 尝试过添加重试机制、心跳消息,更换Sinks.Many替代SubmissionPublisher,均未解决问题
  4. 偶发事件无法到达Controller的情况,但暂未关联到SSE问题

解决方案建议

1. 修复SSE接口的响应头配置

Spring WebFlux默认的SSE响应头可能存在缺失,需显式指定Content-Type为text/event-stream,并禁用缓存:

@GetMapping("/{id}/progress/stream")
public Flux<ServerSentEvent<JobProgress>> subscribeToProgress(@PathVariable("id") String id, ServerHttpResponse response) {
  // 设置SSE响应头
  response.getHeaders().setContentType(MediaType.TEXT_EVENT_STREAM);
  response.getHeaders().setCacheControl(CacheControl.noCache().getHeaderValue());
  
  return Flux.merge(JdkFlowAdapter.flowPublisherToFlux(progressPublisher)
    .filter(progress -> progress.id().equals(id))
    .map(progress -> ServerSentEvent.<JobProgress>builder()
      .id(progress.id())
      .event("job-progress")
      .data(progress)
      .build()
    ), 
    Flux.interval(Duration.ofMillis(5000L)) // 调整心跳间隔为5秒,避免频繁请求
      .map(sequence -> ServerSentEvent.<JobProgress>builder()
        .event("keep-alive")
        .comment("Keeping the connection alive.")
        .build()))
    .retryWhen(Retry.fixedDelay(3, Duration.ofSeconds(10))); // 将重试逻辑移到Flux顶层,而非单个事件
}

2. 替换SubmissionPublisher为Sinks.Many(更适配Reactor)

SubmissionPublisher是JDK Flow API实现,与Reactor的Flux适配存在潜在问题,改用Reactor原生的Sinks.Many:

配置类修改:

@Configuration
public class PublisherConfiguration {
  @Bean
  public Sinks.Many<JobProgress> progressSink() {
    // 多播模式,支持多个订阅者,保留最新事件供新订阅者
    return Sinks.many().multicast().onBackpressureBuffer();
  }
}

事件发布组件修改:

@Component
public class JobEventPublisher {
  private final Sinks.Many<JobProgress> progressSink;

  // 构造注入替代@Autowired,避免final字段注入问题
  public JobEventPublisher(Sinks.Many<JobProgress> progressSink) {
    this.progressSink = progressSink;
  }

  public void dispatchToClient(JobProgress progress) {
    // 非阻塞发送,忽略发送失败(可根据业务添加回调)
    progressSink.tryEmitNext(progress);
  }
}

Controller接口修改:

@GetMapping("/{id}/progress/stream")
public Flux<ServerSentEvent<JobProgress>> subscribeToProgress(@PathVariable("id") String id, ServerHttpResponse response) {
  response.getHeaders().setContentType(MediaType.TEXT_EVENT_STREAM);
  response.getHeaders().setCacheControl(CacheControl.noCache().getHeaderValue());
  
  return Flux.merge(
    progressSink.asFlux()
      .filter(progress -> progress.id().equals(id))
      .map(progress -> ServerSentEvent.<JobProgress>builder()
        .id(progress.id())
        .event("job-progress")
        .data(progress)
        .build()
      ),
    Flux.interval(Duration.ofMillis(5000L))
      .map(sequence -> ServerSentEvent.<JobProgress>builder()
        .event("keep-alive")
        .build()
      )
  ).retryWhen(Retry.fixedDelay(3, Duration.ofSeconds(10)));
}

3. 前端调整事件监听方式

后端指定了event字段为job-progress和keep-alive,前端需要对应监听指定事件,而非默认的onmessage(仅监听无event字段的消息):

useEffect(() => {
  const id = "some-uuid-received-on-post";
  const source = new EventSource(`${baseUrl}/api/v1/job/${id}/progress/stream`, {withCredentials: true});

  // 监听任务进度事件
  source.addEventListener("job-progress", (event) => {
    const progress = JSON.parse(event.data);
    console.log("任务进度:", progress);
  });

  // 监听心跳事件(可选)
  source.addEventListener("keep-alive", () => {
    console.log("心跳连接保持");
  });

  // 监听错误
  source.onerror = (error) => {
    console.error("SSE连接错误:", error);
    source.close();
  };

  return () => source.close();
}, [baseUrl]);

4. 检查Spring Cloud Gateway配置

确保Gateway允许SSE的长连接,添加以下配置:

spring:
  cloud:
    gateway:
      httpclient:
        connect-timeout: 30000
        response-timeout: 3600000 # 1小时超时,适配长连接
      default-filters:
        - AddResponseHeader=Cache-Control, no-cache
        - AddResponseHeader=Connection, keep-alive

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 11:58:13