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]);
问题现象
- 后端调试日志显示事件已通过Flux传递并映射为
ServerSentEvent,但网络面板中EventSource始终处于挂起状态,客户端从未收到消息 - 客户端与服务端之间存在Spring Boot Cloud Gateway,已排除WebSocket影响(仅需单向通信)
- 尝试过添加重试机制、心跳消息,更换
Sinks.Many替代SubmissionPublisher,均未解决问题 - 偶发事件无法到达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
相关产品推荐
相关产品推荐

