如何使用WebFlux WebClient结合Sinks.Many发布Flux对象?
问题描述
我有如下Flux对象:
MyObject o1 = new MyObject("publish", "one"); MyObject o2 = new MyObject("publish", "two"); MyObject o3 = new MyObject("publish", "three"); Flux<Object> myFlux = Flux.just(o1, o2, o3);
我希望通过默认WebClient将该Flux发送至某个端点,代码如下:
public Flux<String> publishFlux(Flux<Object> events) { return eventClient.post() .uri(env.getProperty(BASE)) .body(events, Object.class) .exchangeToFlux(resp -> resp.bodyToFlux(String.class)); }
我尝试使用Sinks.Many()来发布Flux对象,但不清楚正确的实现方式,目前的尝试代码如下:
Many<Object> sink = Sinks.many().multicast().directBestEffort(); sink.tryEmitNext(myFlux); publishFlux(sink.asFlux()).subscribe();
或者:
Many<Object> sink = Sinks.many().multicast().directBestEffort(); value.subscribe(x -> sink.tryEmitNext(x)); publishFlux(sink.asFlux()).subscribe();
请问我哪里操作有误?
错误分析与正确实现
你的两处问题
- 第一个写法的核心错误:
sink.tryEmitNext(myFlux)是把整个Flux<Object>对象当成单个元素往Sink里塞,而不是把Flux里的每个MyObject实例发出去。tryEmitNext的参数是单个元素,不是Flux,这会导致WebClient最终发送的是Flux对象本身,完全不是你要的业务数据。 - 第二个写法的隐患:直接用
subscribe遍历原Flux转发元素,逻辑上能拿到每个MyObject,但没处理原Flux的完成/错误信号——Sink收不到完成信号的话,WebClient的请求会一直挂着,等不到结束;另外也没处理背压,元素多了可能会溢出。
正确的实现方式
方式一:完整转发所有信号(推荐)
要把原Flux里的元素、完成、错误信号都转发给Sink,确保WebClient能正常结束请求,同时处理发送失败的情况:
// 创建多播Sink,若场景为单订阅者,换成unicast更高效 Sinks.Many<Object> sink = Sinks.many().multicast().directBestEffort(); // 绑定原Flux和Sink,转发所有信号 myFlux.doOnNext(item -> { // 发送单个元素,处理发送结果 Sinks.EmitResult result = sink.tryEmitNext(item); if (result.isFailure()) { // 可添加日志或重试逻辑 System.err.println("元素发送失败: " + result); } }) .doOnComplete(() -> sink.tryEmitComplete()) // 原Flux完成时通知Sink .doOnError(err -> sink.tryEmitError(err)) // 原Flux出错时通知Sink .subscribe(); // 启动订阅,触发元素流动 // 用Sink生成的Flux调用发布方法 publishFlux(sink.asFlux()).subscribe( resp -> System.out.println("收到响应: " + resp), err -> System.err.println("请求失败: " + err), () -> System.out.println("请求全部完成") );
方式二:简化版转发
如果不需要处理每个元素的发送结果,也可以直接用Sink的方法作为订阅回调,更简洁:
Sinks.Many<Object> sink = Sinks.many().multicast().directBestEffort(); // 直接把Sink的方法传给subscribe,自动转发元素、错误、完成信号 myFlux.subscribe(sink::tryEmitNext, sink::tryEmitError, sink::tryEmitComplete); publishFlux(sink.asFlux()).subscribe();
额外提示
- 如果你的场景是只有一个订阅者(比如只发一次WebClient请求),用
Sinks.many().unicast().onBackpressureBuffer()比multicast更高效,因为不需要处理多订阅者的逻辑 - 别忽略
tryEmitXXX的返回结果,生产环境里至少要打日志,避免静默失败 - 一定要转发完成/错误信号,否则WebClient会一直等待,导致连接泄漏
内容的提问来源于stack exchange,提问作者rupweb
相关产品推荐
相关产品推荐

