用 Gatling 实现 WebSocket 长连接压测,并统计消息延迟分布

解读

在国内互联网与金融级系统面试中,WebSocket 已广泛用于行情推送、IM、实时风控等场景。面试官抛出“用 Gatling 做长连接压测并统计消息延迟分布”,核心想验证四点:

  1. 是否真正理解 WebSocket 全双工特性与长连接生命周期;
  2. 能否用 Gatling 的 WebSocket DSL 正确建立连接、维持心跳、异步收发;
  3. 如何把“客户端发出时间—收到服务端回包时间”这一关键指标注入 Gatling 的 Session,并自定义延迟直方图;
  4. 是否具备把压测结果与 SLA(如 P99 ≤ 200 ms)对齐、最终定位瓶颈并给出调优建议的能力。
    回答时切忌只贴脚本,要体现“指标定义→脚本设计→数据采集→报告解读→优化闭环”的完整思路,才能与国内一线厂“结果可落地”的面试标准匹配。

知识点

  1. WebSocket 协议帧格式、掩码规则及 1000/1001/1006 等常见关闭码含义;
  2. Gatling WebSocket DSL:exec(ws("Connect").connect(...))、await、check、sendText、sendBinary、reconciliate、close;
  3. Gatling 会话变量与 transform 回调:利用 Session 注入时间戳,实现“请求级”自定义指标;
  4. Gatling Graphite/InfluxDB 插件 + Grafana 模板,构建“延迟直方图”与“连接数”双轴视图;
  5. 高并发下内核参数优化:net.core.somaxconn、net.ipv4.tcp_tw_reuse、ulimit nofile;
  6. 服务端背压场景:Netty 高低水位、backlog、Reactive Streams 限流;
  7. SLA 量化:P50/P90/P99/P999 分位、吞吐 QPS、错误率 <0.1%、GC 停顿 <50 ms;
  8. 稳定性陷阱:心跳超时导致连接漂移、NAT 设备老化、TCP 重传;
  9. 结果归因方法论:CPU 火焰图 → 锁竞争 → 日志异步队列 → 网卡软中断不均衡;
  10. 交付物规范:压测方案、脚本、监控大盘截图、瓶颈分析、调优 patch、回归报告。

答案

以下示例基于 Gatling 3.9、Scala 2.13,模拟“客户端发送订阅指令,服务端每 1 秒推送一条行情,客户端计算端到端延迟”的核心场景。脚本可直接在 IDEA + Maven 工程运行,也可集成到 Jenkins + InfluxDB 流水线。

  1. 指标定义
    延迟 = 客户端收到推送时间戳 – 客户端本地发送订阅时间戳;
    采样维度:连接数 5 k、阶梯递增 500/s,持续 10 min;
    SLA:P99 ≤ 200 ms、错误率 <0.1%、CPU ≤ 60%。

  2. Maven 依赖

<dependency>
  <groupId>io.gatling.highcharts</groupId>
  <artifactId>gatling-charts-highcharts</artifactId>
  <version>3.9.5</version>
</dependency>
  1. 核心脚本
package computerdatabase

import io.gatling.core.Predef._
import io.gatling.http.Predef._
import io.gatling.websocket.Predef._
import scala.concurrent.duration._

class WebSocketLatency extends Simulation {

  val httpProtocol = http
    .baseUrl("https://gateway.example.com") // 仅用于 cookie 或鉴权
    .wsBaseUrl("wss://gateway.example.com/ws")

  val feeder = Iterator.continually(Map(
    "userId" -> java.util.UUID.randomUUID.toString
  ))

  val subscribe =
    exec(ws("Connect").connect("/quote"))
      .exec(session => session.set("subTime", System.currentTimeMillis()))
      .exec(ws("SubScript")
        .sendText("""{"op":"sub","symbol":"BTCUSDT"}""")
        .await(5 seconds)(
          ws.checkTextMessage("firstPush")
            .check(jsonPath("$.price").exists.saveAs("price"))
            .check(jsonPath("$.ts").saveAs("srvTs"))
        ))
      .exec(session => {
        val recv = System.currentTimeMillis()
        val sent = session("subTime").as[Long]
        val latency = recv - sent
        // 注入自定义指标
        session.set("latency", latency)
      })
      // 持续接收推送并更新延迟
      .during(10 minutes) {
        exec(ws("WaitPush")
          .await(30 seconds)(
            ws.checkTextMessage("push")
              .check(jsonPath("$.ts").saveAs("srvTs"))
          )
          .exec(session => {
            val recv = System.currentTimeMillis()
            val sent = session("srvTs").as[String].toLong
            val latency = recv - sent
            session.set("latency", latency)
          })
          .exec(session => {
            // 把 latency 写入 Gatling 直方图
            statsEngine.logResponse(
              session.scenario,
              session.groups,
              "pushLatency",
              startTimestamp = session("srvTs").as[String].toLong,
              endTimestamp   = System.currentTimeMillis(),
              status = OK,
              // 额外标签
              Map("latency" -> session("latency").as[Long].toString)
            )
            session
          })
      }
      .exec(ws("Close").close)

  val scn = scenario("WS_Quote_Latency")
    .feed(feeder)
    .exec(subscribe)

  setUp(
    scn.inject(
      incrementUsersPerSec(500)
        .times(10)
        .eachLevelLasting(1 minute)
        .separatedByRampsLasting(10 seconds)
        .startingFrom(500)
    ).protocols(httpProtocol)
  ).maxDuration(15 minutes)
   .assertions(
     global.responseTime.percentile(99).lt(200),
     global.failedRequests.percent.lte(0.1)
   )
}
  1. 延迟分布采集
    脚本通过 statsEngine.logResponse 把每条推送当成一次虚拟“请求”写入,Gatling 原生报告即可输出直方图;若需秒级实时,可打开 gatling.conf 中 graphite 开关,指向内部 InfluxDB,Grafana 模板变量 $latency 即可展示 P50/P99 曲线。

  2. 常见坑与调优

    • 心跳:服务端 30 s 无数据会 ping,脚本需 .ping.autoReply 防止被踢;
    • NAT 超时:云厂商 LB 默认 60 s 断链,脚本里加 keepAlive 并设置 connectionHeader = "Upgrade"
    • 文件句柄:压测机 ulimit -n 65535,/etc/security/limits.conf 永久生效;
    • 网络软中断:多队列网卡 + RPS 绑定,避免单 CPU 100%;
    • GC:G1 堆 ≤ 8 G,-XX:MaxGCPauseMillis=100,防止长暂停造成延迟毛刺。
  3. 交付示例
    最终报告包含:

    • Grafana 截图:P99 曲线 180 ms,峰值 QPS 52 k;
    • 火焰图:Netty EventLoop 占 38 %,发现 JSON 序列化 hotspot;
    • 调优 patch:把 fastjson 换成 jackson-afterburner,CPU 降 8 %,P99 降到 120 ms;
    • 回归结论:满足 SLA,可上线。

拓展思考

  1. 如果服务端支持二进制 Protobuf,如何把解码耗时也纳入延迟统计?
    思路:在 .check 后加 transform 步骤,把解码开始/结束时间注入 Session,再用 logResponse 写两条记录,一条网络延迟,一条解码延迟,最终在 Grafana 用 latency{layer="net"}latency{layer="decode"} 分层展示。

  2. 当连接数达到 20 k 时,单机内存 8 G 出现 OOM,如何横向扩展?
    国内大厂标准做法:

    • 采用 Gatling FrontLine(企业版)或开源 gatling-k8s-operator,把 Pod 水平扩到 5 个,每个 4 C8 G;
    • 使用 Redis 共享唯一 userId 池,防止不同 Pod 重复登录;
    • 通过 Kubernetes HPA 按 CPU 65 % 自动扩容,压测过程无需人工干预。
  3. 如何验证“消息不丢不重”?
    在推送体里加全局递增 sequence,脚本端用 Scala mutable.TreeSet 做去重与断号检测,最终生成“丢号率”、“重号率”指标,若 >0 即触发断言失败,直接阻断发布。

  4. 面试追问:如果 P99 突然恶化到 2 s,但 CPU、网络、GC 均正常,你会怎么继续定位?
    可答:

    • 抓包看 TCP 重传、窗口满零;
    • 打印 Netty 高低水位日志,确认是否写缓冲区阻塞;
    • ss -ntp 观察 Recv-Q、Send-Q 堆积;
    • 检查云厂商 SLB 限流策略,是否触发 5 Gbps 带宽封顶;
    • 若以上均正常,考虑线程调度延迟,用 perf sched latency 看是否被其他 cgroup 抢占。

掌握以上深度,可在性能岗位面试中直接对标 P7/P8 要求,体现“不仅能跑脚本,更能定位瓶颈并推动修复”的完整闭环能力。