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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 03:25:33