【发布时间】:2019-11-14 08:57:59
【问题描述】:
Spring + Flux 的奇怪行为。我有 Python 服务器代码(使用 Flask,但这并不重要,将其视为伪代码),它是流式响应:
def generate():
for row in range(0,10):
time.sleep(1)
yield json.dumps({"count": row}) + '\n'
return Response(generate(), mimetype='application/json')
这样,我模拟处理列表中的一些任务,并在它们准备好后立即向我发送结果,而不是等待所有事情都完成,主要是为了避免将所有内容先保存在服务器的内存中,然后再保存在客户。现在我想用 Spring WebClient 来使用它:
Flux<Count> alerts = webClient
.post()
.uri("/testStream")
.accept(MediaType.APPLICATION_JSON)
.retrieve()
.bodyToFlux( Count.class )
.log();
alerts.subscribe(a -> log.debug("Received count: " + a.count));
Mono<Void> mono = Mono.when(alerts);
mono.block();
log.debug("All done in method");
这是我在日志中得到的内容:
2019-07-03 18:45:08.330 DEBUG 16256 --- [ctor-http-nio-4] c.k.c.restapi.rest.Controller : Received count: 8
2019-07-03 18:45:09.323 INFO 16256 --- [ctor-http-nio-2] reactor.Flux.MonoFlatMapMany.4 : onNext(com.ksftech.chainfacts.restapi.rest.Controller$Count@55d09f83)
2019-07-03 18:45:09.324 INFO 16256 --- [ctor-http-nio-2] reactor.Flux.MonoFlatMapMany.4 : onComplete()
2019-07-03 18:45:09.325 DEBUG 16256 --- [io-28088-exec-4] c.k.c.restapi.rest.Controller : All done in method
2019-07-03 18:45:09.331 INFO 16256 --- [ctor-http-nio-4] reactor.Flux.MonoFlatMapMany.4 : onNext(com.ksftech.chainfacts.restapi.rest.Controller$Count@da447dd)
2019-07-03 18:45:09.332 DEBUG 16256 --- [ctor-http-nio-4] c.k.c.restapi.rest.Controller : Received count: 9
2019-07-03 18:45:09.333 INFO 16256 --- [ctor-http-nio-4] reactor.Flux.MonoFlatMapMany.4 : onComplete()
注意在 mono.block 返回后订阅如何处理最后一个对象。我知道 Reactor 是异步的,一旦它看不到更多对象,它就会释放 Mono 并并行调用我的代码订阅。那么调度器就看什么先运行了。
我想出了一个相当丑陋的组合,即使用 completeConsumer 订阅,并使用旧的等待/通知。然后它工作正常。但是有没有更优雅的方法来确保我的方法等到 Flux 的所有元素都被处理完?
【问题讨论】:
标签: spring webclient flux reactor