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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 12:01:31