Spring Reactor:如何让Flux输出最后一个非Processing状态的订单状态?
问题:Flux轮询无法输出最终非Processing状态项
我开发了一个接口,通过调用Spring服务获取订单状态,该服务返回的枚举值包括Processing、Completed和Error。接口实现如下:
@GetMapping(value = ["/{orderId}/state"], produces = [MediaType.TEXT_EVENT_STREAM_VALUE]) fun orderState(@PathVariable("orderId") orderId: Long): Flux<OrderStateDetails> { return Flux.interval(Duration.ofSeconds(1)) .take(30) .map { orderService.orderState(orderId) } .takeWhile { orderStateDetails -> orderStateDetails.state == OrderState.Procession } }
我用Flux实现简易轮询,期望逻辑是:直到订单状态不再是Processing或轮询时长达到30秒就停止,返回序列类似:
- Processing
- Processing
- Processing
- Completed
但takeWhile会过滤掉最后一个Completed项,导致无法输出。试过添加.also { orderService.orderState(orderId) }和.concatWith(Flux.just(orderService.orderState(orderId))),但会产生额外的不必要调用,还是没法输出这个状态项。想问怎么让Flux输出这个最后状态?
解决方案
核心问题
takeWhile的逻辑是:仅当元素满足条件时才发射,一旦遇到不满足条件的元素,就直接丢弃该元素并终止流,所以最后那个非Processing的状态不会被输出。另外代码中存在拼写错误:OrderState.Procession应为OrderState.Processing。
修复代码
将takeWhile替换为takeUntil,同时修正拼写错误:
@GetMapping(value = ["/{orderId}/state"], produces = [MediaType.TEXT_EVENT_STREAM_VALUE]) fun orderState(@PathVariable("orderId") orderId: Long): Flux<OrderStateDetails> { return Flux.interval(Duration.ofSeconds(1)) .take(30) .map { orderService.orderState(orderId) } .takeUntil { orderStateDetails -> orderStateDetails.state != OrderState.Processing } }
逻辑说明
takeUntil { it.state != Processing }:持续发射元素,直到遇到第一个非Processing的状态(比如Completed/Error),并且会包含这个状态元素,之后终止流。- 保留
take(30):确保轮询最多执行30次,避免无限轮询。
这样就能完美满足需求:输出所有中间的Processing状态,同时输出最终的非Processing状态,且不会产生额外的服务调用。
内容的提问来源于stack exchange,提问作者Jonas Pedersen
相关产品推荐
相关产品推荐

