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

Reactor-Netty能否扩展连接相关指标?如何实现Postman式响应时间拆分?

如何基于Reactor-Netty HttpClient实现类似Postman的响应时间拆分指标

我们的Spring Reactive项目已经通过默认HTTP Client指标监控上游服务情况,但现在需要实现类似Postman的响应时间拆分功能——即拆分DNS解析、TCP连接、TLS握手、请求发送、等待响应、响应接收等核心阶段的耗时。请问能否通过当前的Reactor-Netty HttpClient直接实现,或者借助第三方库达成需求?

当前WebClient配置代码

HttpClient client = HttpClient.create(ConnectionProvider.builder("fixed")
            .maxConnections(700)
            .maxIdleTime(Duration.ofMillis(config.getPoolMaxIdleTime()))
            .maxLifeTime(Duration.ofMillis(config.getPoolMaxLifeTime()))
            .metrics(true) // 此处开启默认指标
            .build())
    .option(ChannelOption.SO_TIMEOUT, config.getSocketTimeout())
    .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, config.getConnectionTimeout())
    .doOnConnected(connection -> connection
            .addHandler(new ReadTimeoutHandler(config.getReadTimeout(), TimeUnit.MILLISECONDS))
            .addHandler(new WriteTimeoutHandler(config.getWriteTimeout(), TimeUnit.MILLISECONDS))
    );

Postman响应时间拆分说明

Postman会将HTTP请求的全生命周期拆分为以下关键阶段:

  • DNS解析:从发起请求到完成域名解析的耗时
  • TCP连接:建立TCP连接的耗时
  • TLS握手(HTTPS请求):完成TLS安全握手的耗时
  • 请求发送:从开始发送请求到请求完全发送完毕的耗时
  • 等待响应:请求发送完成后到收到第一个响应字节的耗时
  • 响应接收:从收到第一个响应字节到接收完所有响应数据的耗时

解决方案

1. Reactor-Netty原生实现

Reactor-Netty的HttpClient支持通过**自定义HttpClientMetricsRecorder**扩展指标收集逻辑,从而实现各阶段耗时的拆分。

步骤1:自定义指标收集器

实现HttpClientMetricsRecorder接口,在各生命周期节点打点记录耗时:

import io.netty.channel.Channel;
import io.netty.handler.codec.http.HttpRequest;
import io.netty.handler.codec.http.HttpResponse;
import reactor.netty.http.client.HttpClientMetricsRecorder;

import java.net.SocketAddress;
import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

public class DetailedClientMetricsRecorder implements HttpClientMetricsRecorder {

    private final Map<String, StageTimings> requestTimings = new ConcurrentHashMap<>();

    @Override
    public void onConnectStart(SocketAddress remoteAddress) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = new StageTimings();
        timings.connectStart = System.nanoTime();
        requestTimings.put(key, timings);
    }

    @Override
    public void onConnectSuccess(SocketAddress remoteAddress, Duration duration, Channel channel) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            // 记录TCP连接耗时
            Metrics.timer("http.client.tcp_connect.duration").record(duration);
        }
    }

    @Override
    public void onSslHandshakeStart(SocketAddress remoteAddress) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            timings.sslStart = System.nanoTime();
        }
    }

    @Override
    public void onSslHandshakeSuccess(SocketAddress remoteAddress, Duration duration) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            // 记录TLS握手耗时
            Metrics.timer("http.client.tls_handshake.duration").record(duration);
        }
    }

    @Override
    public void onRequestStart(HttpRequest request, SocketAddress remoteAddress) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            timings.requestStart = System.nanoTime();
        }
    }

    @Override
    public void onRequestSuccess(HttpRequest request, SocketAddress remoteAddress, Duration duration) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            // 记录请求发送耗时
            Metrics.timer("http.client.request_send.duration").record(duration);
            timings.requestEnd = System.nanoTime();
        }
    }

    @Override
    public void onResponseStart(HttpResponse response, SocketAddress remoteAddress) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            // 计算并记录等待响应耗时(请求发送完成到收到第一个响应字节的时间)
            Duration waitDuration = Duration.ofNanos(System.nanoTime() - timings.requestEnd);
            Metrics.timer("http.client.response_wait.duration").record(waitDuration);
            timings.responseStart = System.nanoTime();
        }
    }

    @Override
    public void onResponseSuccess(HttpResponse response, SocketAddress remoteAddress, Duration duration) {
        String key = genRequestKey(remoteAddress);
        StageTimings timings = requestTimings.get(key);
        if (timings != null) {
            // 记录响应接收耗时
            Metrics.timer("http.client.response_receive.duration").record(duration);
            // 清理缓存,避免内存泄漏
            requestTimings.remove(key);
        }
    }

    // 按需实现失败场景的指标记录
    @Override public void onConnectFailed(SocketAddress remoteAddress, Duration duration, Throwable throwable) {}
    @Override public void onSslHandshakeFailed(SocketAddress remoteAddress, Duration duration, Throwable throwable) {}
    @Override public void onRequestFailed(HttpRequest request, SocketAddress remoteAddress, Duration duration, Throwable throwable) {}
    @Override public void onResponseFailed(HttpResponse response, SocketAddress remoteAddress, Duration duration, Throwable throwable) {}
    @Override public void onDataReceived(SocketAddress remoteAddress, long bytes) {}
    @Override public void onDataSent(SocketAddress remoteAddress, long bytes) {}
    @Override public void onReset() {}

    // 生成请求唯一标识(可结合请求ID等优化,避免并发冲突)
    private String genRequestKey(SocketAddress remoteAddress) {
        return remoteAddress.toString() + "-" + Thread.currentThread().getId();
    }

    // 存储各阶段时间点的内部类
    private static class StageTimings {
        long connectStart;
        long sslStart;
        long requestStart;
        long requestEnd;
        long responseStart;
    }
}

步骤2:配置HttpClient使用自定义收集器

修改原代码,将默认的metrics配置替换为自定义实现:

HttpClient client = HttpClient.create(ConnectionProvider.builder("fixed")
            .maxConnections(700)
            .maxIdleTime(Duration.ofMillis(config.getPoolMaxIdleTime()))
            .maxLifeTime(Duration.ofMillis(config.getPoolMaxLifeTime()))
            .build())
    .option(ChannelOption.SO_TIMEOUT, config.getSocketTimeout())
    .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, config.getConnectionTimeout())
    .doOnConnected(connection -> connection
            .addHandler(new ReadTimeoutHandler(config.getReadTimeout(), TimeUnit.MILLISECONDS))
            .addHandler(new WriteTimeoutHandler(config.getWriteTimeout(), TimeUnit.MILLISECONDS))
    )
    // 启用自定义指标收集器
    .metrics(true, DetailedClientMetricsRecorder::new);

2. 第三方库辅助方案

如果需要更开箱即用的实现,可以引入Micrometer Reactor Netty扩展依赖(io.micrometer:micrometer-reactor-netty),它基于Reactor-Netty的MetricsRecorder接口封装了更丰富的默认指标,同时支持自定义扩展各阶段的耗时统计。结合Spring Boot Actuator可以快速将这些指标暴露到Prometheus等监控系统中。


内容的提问来源于stack exchange,提问作者Md. Hasan Basri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 07:40:35