WebFlux Sinks多服务流问题:消费者无法接收最新事件
问题分析与修复方案
你的代码存在几个核心问题,导致消费者无法接收到事件,以下是逐一拆解和修复:
1. Sink类型不符合需求
你使用的multicast().onBackpressureBuffer()是实时多播Sink,它不会保留历史事件——只有在事件发送时已经存在的订阅者才能收到消息。而你的消费者是定时发起请求订阅,大概率是在事件发送之后才建立订阅,自然收不到任何内容。
如果要让新订阅者能立即拿到最新的事件,应该改用重放(replay)类型的Sink:
@Bean public Sinks.Many<Event> sink() { // 重放最新的1个事件,新订阅者订阅时直接获取最近的事件 return Sinks.many().replay().latest(); }
2. 消费者的订阅方式错误
每3秒创建一个新的WebClient请求去获取单个事件,这种方式会导致:
- 大量挂起的HTTP请求(无事件时Mono会一直等待,直到超时)
- 只有刚好在事件发送时处于等待状态的请求能收到消息,其余请求要么挂死要么失败
正确的做法是让消费者持续订阅整个事件流,而不是定时发起单次请求:
@EnableScheduling @RequiredArgsConstructor @SpringBootApplication public class ReactiveEventsSubscriberApplication { private final WebClient webClient; public static void main(String[] args) { SpringApplication.run(ReactiveEventsSubscriberApplication.class, args); } // 应用启动时立即开始持续订阅事件流 @PostConstruct public void subscribeToEvents() { webClient.get() .retrieve() .bodyToFlux(Event.class) // 订阅整个事件流,而非单个事件 .subscribe( event -> System.out.println("Received event: " + event), error -> System.err.println("Subscription error: " + error.getMessage()) ); } }
同时需要修改生产者的GetMapping,返回完整的事件流而非单个事件:
@GetMapping() public Flux<Event> getEvents() { return sink.asFlux(); }
3. 生产者参数接收可能存在隐性错误
用@RequestParam Event event接收枚举时,必须确保请求参数名是event,且参数值与枚举的字符串表示完全匹配(比如枚举值是USER_CREATED,请求需传?event=USER_CREATED)。如果参数不匹配,Spring会抛出转换异常,事件根本不会被发送到Sink中。
可以给生产者的PostMapping添加异常捕获和返回值,方便排查问题:
@PostMapping() public ResponseEntity<?> publish(@RequestParam Event event) { try { Sinks.EmitResult emitResult = sink.tryEmitNext(event); if (emitResult.isFailure()) { return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR) .body("Failed to publish event: " + emitResult); } System.out.println("Published event: " + event); return ResponseEntity.ok("Event published: " + event); } catch (Exception e) { System.err.println("Error receiving event: " + e.getMessage()); return ResponseEntity.badRequest().body("Invalid event value"); } }
额外检查点
- 确认消费者的WebClient配置了正确的生产者地址(比如
baseUrl = "http://localhost:8080") - 生产者和消费者的
Event枚举定义必须完全一致(包名、枚举值都要匹配,否则反序列化失败) - 检查应用日志,确认是否有参数转换、网络连接等隐性异常
内容的提问来源于stack exchange,提问作者Smetana Po Aktsii
相关产品推荐
相关产品推荐

