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

Mono.subscribe()执行后计数变量始终为0,如何统计执行次数并打印日志?

问题根因

你遇到的是响应式编程的典型异步时序问题:subscribe()为非阻塞调用,触发异步逻辑后会立即返回,不会等待Mono流执行完成。你调用两个业务方法后立刻打印计数器时,所有doOnSuccess的自增逻辑还未执行,因此输出始终为初始值0。同时要注意如果计数器未使用线程安全实现,并发更新时也会出现计数不准的问题。

实现方案

推荐优先使用标准响应式非阻塞方案,确实无法修改业务方法时可选用阻塞兼容方案:

方案1:响应式无阻塞实现(推荐)

不要在业务方法内部硬编码subscribe()调用,将组装好的Mono/Flux流向上返回,上层统一等待所有流执行完成后再打印日志:

  • 调整业务方法,删除内部subscribe()调用,返回组装完成的响应式流:
// 修改前的业务方法内部逻辑
Mono.just(业务逻辑)
    .doOnSuccess { /* 自增计数逻辑 */ }
    .subscribe() // 删除该行

// 修改后的业务方法示例
fun sendAllTalktalkMessages(naverId: String): Mono<Void> {
    return Mono.just(业务逻辑)
        .doOnSuccess { totalNumInst.totalNum.incrementAndGet() } // 这里使用AtomicInteger保证线程安全
        .then()
}
  • 上层统一拼接流并等待执行完成后打印:
val totalNumInst = TotalNumObj()

Mono.when(
    devSupportService.sendAllTalktalkMessages(naverId),
    devSupportService.sendAllAutoDepositTalktalkMessages(naverId, totalNum)
)
.doOnSuccess { 
    logger.info("实际执行次数:${totalNumInst.totalNum}")
}
.subscribe() // 整个业务流统一触发订阅

方案2:阻塞兼容实现(仅适配无法修改业务方法的场景)

使用CountDownLatch阻塞等待所有异步任务执行完成:

  • 首先确认两个业务方法触发的Mono总数量,初始化对应计数的CountDownLatch
  • 在每个doOnSuccess回调中追加latch减1逻辑:
.doOnSuccess {
    totalNumInst.totalNum.incrementAndGet()
    latch.countDown()
}
  • 上层调用后等待所有任务完成再打印日志:
val totalNumInst = TotalNumObj()
val latch = CountDownLatch(Mono总数量) // 替换为实际的Mono总数

devSupportService.sendAllTalktalkMessages(naverId)
devSupportService.sendAllAutoDepositTalktalkMessages(naverId, totalNum)

latch.await() // 阻塞等待所有任务执行完成
logger.info("实际执行次数:${totalNumInst.totalNum}")
注意事项
  • 计数器必须使用线程安全实现,比如java.util.concurrent.atomic.AtomicInteger,避免多线程并发更新导致的计数丢失、可见性问题
  • 方案1符合响应式编程规范,不会阻塞线程,性能更好;方案2会阻塞调用线程,仅建议用于兼容旧代码的场景

内容的提问来源于stack exchange,提问作者dev.sunset

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 23:39:03