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

如何使用Micrometer Timer记录返回Mono/Flux的异步方法执行时长

使用Micrometer记录异步方法(如Kafka ReplyingTemplate)的实际耗时

你遇到的问题很典型——同步场景下的Timer.record(Supplier)没法直接适配异步/反应式操作,因为它只会记录Supplier创建的瞬间,而不会等待异步任务完成。针对你用的Reactor(Mono)和Kafka ReplyingTemplate场景,我推荐这几种靠谱的实现方式:

1. 直接用Micrometer针对反应式类型的record方法

Micrometer专门为Reactor的Mono/Flux提供了重载的record方法,它会自动订阅流并在任务完成(成功或失败)时记录从订阅到结束的耗时,完美匹配sendAndReceive的实际执行周期。

修改你的代码如下:

public Mono<String> sendRequest(Mono<String> request) {
    return request
        .map(r -> new ProducerRecord<String, String>(requestsTopic, r))
        .map(pr -> {
            pr.headers()
                    .add(new RecordHeader(KafkaHeaders.REPLY_TOPIC,
                            "reply-topic".getBytes()));
            return pr;
        })
        // 用record包裹sendAndReceive返回的Mono,自动记录执行耗时
        .map(pr -> responseGenerationTimer.record(replyingKafkaTemplate.sendAndReceive(pr)))
        // 后续的map/filter等操作不受影响
        ... 
}

2. 手动计时(适合需要自定义逻辑的场景)

如果你需要对成功/失败场景分别计时,或者要添加额外的业务逻辑,可以用Stopwatch结合doOnSuccess/doOnError来手动控制:

public Mono<String> sendRequest(Mono<String> request) {
    return request
        .map(r -> new ProducerRecord<String, String>(requestsTopic, r))
        .map(pr -> {
            pr.headers()
                    .add(new RecordHeader(KafkaHeaders.REPLY_TOPIC,
                            "reply-topic".getBytes()));
            return pr;
        })
        .flatMap(pr -> Mono.defer(() -> {
            // 创建并启动计时器
            Stopwatch stopwatch = Stopwatch.createStarted();
            return replyingKafkaTemplate.sendAndReceive(pr)
                // 成功完成时记录耗时
                .doOnSuccess(response -> {
                    stopwatch.stop();
                    responseGenerationTimer.record(stopwatch.elapsed());
                    // 可选:单独记录成功场景的计时器
                    // successTimer.record(stopwatch.elapsed());
                })
                // 出错时也记录耗时(避免只统计成功场景)
                .doOnError(error -> {
                    stopwatch.stop();
                    responseGenerationTimer.record(stopwatch.elapsed());
                    // 可选:单独记录错误场景的计时器
                    // errorTimer.record(stopwatch.elapsed());
                });
        }))
        ...
}

3. 用Reactor原生的Metrics集成

Reactor自带了和Micrometer的集成,可以直接通过metrics()方法自动生成计时器,不需要手动创建Timer实例:

public Mono<String> sendRequest(Mono<String> request) {
    return request
        .map(r -> new ProducerRecord<String, String>(requestsTopic, r))
        .map(pr -> {
            pr.headers()
                    .add(new RecordHeader(KafkaHeaders.REPLY_TOPIC,
                            "reply-topic".getBytes()));
            return pr;
        })
        .map(pr -> replyingKafkaTemplate.sendAndReceive(pr)
            // 设置计时器名称和标签,方便监控区分不同场景
            .name("kafka.send_and_receive.duration")
            .tag("request_topic", requestsTopic)
            .tag("reply_topic", "reply-topic")
            // 开启自动metrics记录
            .metrics())
        ...
}

几个关键注意点:

  • 标签要精准:给计时器添加主题、状态(成功/失败)等标签,后续监控时能快速定位问题。
  • 不要漏订阅:反应式流只有被订阅才会执行,确保你的Mono最终被订阅(比如控制器返回给WebFlux,或者调用subscribe()),否则计时器不会触发。
  • 错误场景也要统计:别只记录成功的耗时,错误场景的耗时同样重要,能帮你排查超时、异常等问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:11:38