用 Gatling 实现 WebSocket 长连接压测,并统计消息延迟分布
解读
在国内互联网与金融级系统面试中,WebSocket 已广泛用于行情推送、IM、实时风控等场景。面试官抛出“用 Gatling 做长连接压测并统计消息延迟分布”,核心想验证四点:
- 是否真正理解 WebSocket 全双工特性与长连接生命周期;
- 能否用 Gatling 的 WebSocket DSL 正确建立连接、维持心跳、异步收发;
- 如何把“客户端发出时间—收到服务端回包时间”这一关键指标注入 Gatling 的 Session,并自定义延迟直方图;
- 是否具备把压测结果与 SLA(如 P99 ≤ 200 ms)对齐、最终定位瓶颈并给出调优建议的能力。
回答时切忌只贴脚本,要体现“指标定义→脚本设计→数据采集→报告解读→优化闭环”的完整思路,才能与国内一线厂“结果可落地”的面试标准匹配。
知识点
- WebSocket 协议帧格式、掩码规则及 1000/1001/1006 等常见关闭码含义;
- Gatling WebSocket DSL:exec(ws("Connect").connect(...))、await、check、sendText、sendBinary、reconciliate、close;
- Gatling 会话变量与 transform 回调:利用 Session 注入时间戳,实现“请求级”自定义指标;
- Gatling Graphite/InfluxDB 插件 + Grafana 模板,构建“延迟直方图”与“连接数”双轴视图;
- 高并发下内核参数优化:net.core.somaxconn、net.ipv4.tcp_tw_reuse、ulimit nofile;
- 服务端背压场景:Netty 高低水位、backlog、Reactive Streams 限流;
- SLA 量化:P50/P90/P99/P999 分位、吞吐 QPS、错误率 <0.1%、GC 停顿 <50 ms;
- 稳定性陷阱:心跳超时导致连接漂移、NAT 设备老化、TCP 重传;
- 结果归因方法论:CPU 火焰图 → 锁竞争 → 日志异步队列 → 网卡软中断不均衡;
- 交付物规范:压测方案、脚本、监控大盘截图、瓶颈分析、调优 patch、回归报告。
答案
以下示例基于 Gatling 3.9、Scala 2.13,模拟“客户端发送订阅指令,服务端每 1 秒推送一条行情,客户端计算端到端延迟”的核心场景。脚本可直接在 IDEA + Maven 工程运行,也可集成到 Jenkins + InfluxDB 流水线。
-
指标定义
延迟 = 客户端收到推送时间戳 – 客户端本地发送订阅时间戳;
采样维度:连接数 5 k、阶梯递增 500/s,持续 10 min;
SLA:P99 ≤ 200 ms、错误率 <0.1%、CPU ≤ 60%。 -
Maven 依赖
<dependency>
<groupId>io.gatling.highcharts</groupId>
<artifactId>gatling-charts-highcharts</artifactId>
<version>3.9.5</version>
</dependency>
- 核心脚本
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)
)
}
-
延迟分布采集
脚本通过statsEngine.logResponse把每条推送当成一次虚拟“请求”写入,Gatling 原生报告即可输出直方图;若需秒级实时,可打开gatling.conf中 graphite 开关,指向内部 InfluxDB,Grafana 模板变量$latency即可展示 P50/P99 曲线。 -
常见坑与调优
- 心跳:服务端 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,防止长暂停造成延迟毛刺。
- 心跳:服务端 30 s 无数据会 ping,脚本需
-
交付示例
最终报告包含:- Grafana 截图:P99 曲线 180 ms,峰值 QPS 52 k;
- 火焰图:Netty EventLoop 占 38 %,发现 JSON 序列化 hotspot;
- 调优 patch:把 fastjson 换成 jackson-afterburner,CPU 降 8 %,P99 降到 120 ms;
- 回归结论:满足 SLA,可上线。
拓展思考
-
如果服务端支持二进制 Protobuf,如何把解码耗时也纳入延迟统计?
思路:在.check后加transform步骤,把解码开始/结束时间注入 Session,再用logResponse写两条记录,一条网络延迟,一条解码延迟,最终在 Grafana 用latency{layer="net"}与latency{layer="decode"}分层展示。 -
当连接数达到 20 k 时,单机内存 8 G 出现 OOM,如何横向扩展?
国内大厂标准做法:- 采用 Gatling FrontLine(企业版)或开源 gatling-k8s-operator,把 Pod 水平扩到 5 个,每个 4 C8 G;
- 使用 Redis 共享唯一 userId 池,防止不同 Pod 重复登录;
- 通过 Kubernetes HPA 按 CPU 65 % 自动扩容,压测过程无需人工干预。
-
如何验证“消息不丢不重”?
在推送体里加全局递增 sequence,脚本端用 Scala mutable.TreeSet 做去重与断号检测,最终生成“丢号率”、“重号率”指标,若 >0 即触发断言失败,直接阻断发布。 -
面试追问:如果 P99 突然恶化到 2 s,但 CPU、网络、GC 均正常,你会怎么继续定位?
可答:- 抓包看 TCP 重传、窗口满零;
- 打印 Netty 高低水位日志,确认是否写缓冲区阻塞;
- 用
ss -ntp观察 Recv-Q、Send-Q 堆积; - 检查云厂商 SLB 限流策略,是否触发 5 Gbps 带宽封顶;
- 若以上均正常,考虑线程调度延迟,用
perf sched latency看是否被其他 cgroup 抢占。
掌握以上深度,可在性能岗位面试中直接对标 P7/P8 要求,体现“不仅能跑脚本,更能定位瓶颈并推动修复”的完整闭环能力。