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

递归调用中非阻塞处理Mono:无需阻塞判断下游响应是否为空

问题分析与解决方案

为什么只打印"YY",没打印"ZZ""XX"?

因为Mono是冷发布者,只有当被subscribe()(订阅)时才会触发内部逻辑执行。你现在只是调用了flatMap操作符,但没有订阅这个Mono,所以flatMap里的代码完全不会运行,只会执行listOfMonos.last块里的println("YY")。

是否需要阻塞?完全不需要

阻塞会破坏Reactor的非阻塞特性,导致服务性能下降。用响应式编程的方式就能实现需求,全程异步非阻塞。

非阻塞实现方案

原来的同步递归逻辑有问题:你在Mono还没完成的时候就直接调用了rec,导致startingCoordinate根本没机会被响应结果更新。正确的做法是把递归逻辑嵌入到响应式调用链中,等当前Mono完成后再触发下一次递归。

修改后的代码示例

import reactor.core.publisher.Mono
import java.util.concurrent.ConcurrentHashMap

private fun rec(
    storeId: Int,
    visited: MutableSet<CartesianCoordinate>,
    previousCartesianCoordinate: CartesianCoordinate?,
    destinations: List<PathingCartesianPosition>
): Mono<Unit> {
    // 终止条件:所有目的地都已访问
    if (visited.size == destinations.size) {
        return Mono.just(Unit)
    }

    val currentCoordinate = findClosestCoordinate(destinations, previousCartesianCoordinate, visited)
    currentCoordinate.first?.let { visited.add(it) }

    // 构建查询参数(处理null情况避免NPE)
    val query = previousCartesianCoordinate?.let { prev ->
        currentCoordinate.first?.let { curr ->
            PathServiceQuery(prev, curr)
        }
    }
    val queryString = query?.let { ObjectMapper().writeValueAsString(it) } ?: ""

    // 调用下游服务,处理响应后触发递归
    return webClient.getPath(storeId, queryString, currentCoordinate.second)
        .cast(PathingCartesianResponse::class.java)
        .flatMap { response ->
            println("ZZ")
            var nextCoordinate = previousCartesianCoordinate
            if (!response.coordinates.isNullOrEmpty()) {
                println("XX")
                nextCoordinate = currentCoordinate.first
            }
            // 当前响应处理完成后,触发下一次递归
            rec(storeId, visited, nextCoordinate, destinations)
        }
}

关键说明

  1. 函数返回Mono:代表当前异步步骤完成,整个递归流程是一个链式的异步操作
  2. 线程安全集合:用ConcurrentHashMap.newKeySet()(或CopyOnWriteArraySet)代替原来的HashSet,避免响应式环境下多线程操作的并发问题
  3. 响应式递归触发:只有当前下游服务调用完成(flatMap执行)后,才会调用下一次rec,保证nextCoordinate是被响应结果更新后的值
  4. 触发整个流程:在业务代码中订阅这个递归函数返回的Mono即可启动整个流程,比如:
    val visited = ConcurrentHashMap.newKeySet<CartesianCoordinate>()
    rec(storeId, visited, initialCoordinate, destinations)
        .subscribe({
            // 整个递归流程完成后的回调
        }, { error ->
            // 异常处理
        })
    

原代码其他问题修正

  • 避免在响应式链外修改可变变量(比如原来的startingCoordinate),响应式编程更推荐通过调用链传递状态
  • 处理query为null的情况,避免ObjectMapper序列化null时抛出异常

内容的提问来源于stack exchange,提问作者anonymous

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 04:07:10