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

Reactor中List<Mono<String>>空结果的重试实现咨询

解决Mono返回空字符串的重试问题

首先得提一句,你现在用block()的方式其实违背了Reactor响应式设计的核心原则——阻塞操作会破坏非阻塞模型的性能优势,我们完全可以用Reactor原生操作符来实现精准的空结果重试逻辑,全程保持响应式流的链式调用。

步骤1:给单个Mono添加空结果重试逻辑

针对每个返回String的Mono,我们可以组合filter和repeatWhenEmpty操作符,专门处理“返回空字符串就重试”的场景:

fun retryOnEmpty(mono: Mono<String>, maxRetries: Int = 3): Mono<String> {
    return mono
        // 过滤掉空字符串,触发后续的重试逻辑
        .filter { it.isNotEmpty() }
        // 当流为空(即空字符串被过滤后无元素发出)时,重试指定次数
        .repeatWhenEmpty(Retry.max(maxRetries))
}

这里repeatWhenEmpty是Reactor专门为“流无元素完成”场景设计的重试操作符,比通用的retryWhen更精准,不需要额外判断错误类型。如果需要避免短时间频繁调用API,还可以换成指数退避重试:Retry.backoff(maxRetries, Duration.ofSeconds(1))。

步骤2:批量处理List<Mono>集合

如果你手里是一个List<Mono<String>>,可以把它转成Flux并行处理所有请求,同时给每个Mono加上上面的重试逻辑:

val apiCallMonos: List<Mono<String>> = // 你的API调用Mono集合
val processedResults: Mono<List<String>> = Flux.fromIterable(apiCallMonos)
    // 给每个Mono绑定空结果重试逻辑
    .flatMap { retryOnEmpty(it) }
    // 收集所有结果到List,等待所有请求完成
    .collectList()

结合你现有代码的优化

看你当前代码里用block()同步获取statistic的值,这会阻塞线程,我们把这段改成响应式链式调用,同时如果statistic对应的Mono返回的Map里存在空字符串需要重试,也可以套用上面的逻辑:

// 把原来的block()替换成响应式处理
val firstAskForItemsRetrieved = firstAskForItems.flatMap { item ->
    val statisticMono = item["statistic"] as Mono<Map<Any, Any>>
    
    // 假设我们要检查statistic里的"targetField"字段是否为空,为空则重试
    statisticMono
        .filter { stat ->
            val targetValue = stat["targetField"] as? String ?: ""
            targetValue.isNotEmpty()
        }
        .repeatWhenEmpty(Retry.max(3))
        .map { statistic ->
            item.apply { this["statistic"] = statistic }
        }
}

// 针对第一次处理后仍存在的空结果,发起二次重试
val secondAskForItems = firstAskForItemsRetrieved
    // 筛选出仍有空结果的项
    .filter { item ->
        val stat = item["statistic"] as Map<Any, Any>
        val targetValue = stat["targetField"] as? String ?: ""
        targetValue.isEmpty()
    }
    // 重新发起API调用并绑定重试逻辑
    .flatMap { item ->
        // 替换成你实际的重新获取statistic的API调用
        fetchStatisticApi()
            .filter { stat ->
                val targetValue = stat["targetField"] as? String ?: ""
                targetValue.isNotEmpty()
            }
            .repeatWhenEmpty(Retry.max(2))
            .map { newStat ->
                item.apply { this["statistic"] = newStat }
            }
    }

关键提醒

  • 除非是应用启动等极少数同步场景,否则永远不要在响应式流中使用block(),它会阻塞线程池,严重降低系统吞吐量。
  • repeatWhenEmpty只会在“流完成但无元素发出”时触发重试,刚好匹配我们过滤空字符串后的场景,逻辑精准且高效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 02:30:17