ActiveMQ MQTT 连接风暴排查手记
业务背景
映美云打印平台,在线云打印机 4 万+ 台,涵盖云针式打印机、热敏打印机、喷墨电票打印机等多个产品线。打印机设备通过 MQTT 协议接入 ActiveMQ,与上层应用通信。
应用端主要包括:
- 映美云打印 App— 移动端打印入口,下发打印订单
- e 标签— 标签打印业务,产送标签模板与打印任务
- 云驱动— PC 端云打印驱动,模拟本地打印机体验
- 映管家— 设备管理与运维平台,管理打印机状态与配置
消息流向:各应用端生产/订阅打印状态和打印订单状态,通过 ActiveMQ 的 Topic 模式推送到对应的打印机设备。设备上线时订阅自身的print/{设备号}主题,应用端通过该主题下发打印任务,设备回传打印状态。
故障概述
| 项目 | 说明 |
|---|---|
| 服务器 | 4 核 8G |
| 中间件 | Apache ActiveMQ 6.0.1 |
| 协议 | MQTT (mqtt+nio) |
| 在线设备 | 4 万+ 台云打印机 |
| 现象 | 连接数飙升 3-4 万,线程数爆炸到 2 万+,CPU 跑满 350% |
| 结果 | 进程反复崩溃重启,打印任务大面积延迟或失败 |
| 排查周期 | 约一周的持续排查与多轮优化 |
一、问题初现:Timer already cancelled
最初的告警很模糊——客户端报 “Could not accept connection: Timer already cancelled”。监控脚本显示 CPU 不高(当时监控脚本有 bug,CPU 读数始终为 0),但netstat一查吓一跳:
[root@xxx ~]# netstat -antp | grep java | grep 1883 | wc -l 348193 万+ 连接,而当时配置的maximumConnections是默认值。第一反应是"连接数超限了",但深入看监控日志发现更诡异的现象:
- 线程数一度冲到 3088
- 连接数忽高忽低,像潮水一样涨落
- 进程每隔几十分钟就崩一次,监控脚本自动拉起来
第一轮猜测(错的)
一开始怀疑是 Linux 系统参数问题:tcp_max_tw_buckets只有 5000,somaxconn不够,ulimit -n限制等等。调了一轮系统参数,效果不大——问题根本不在系统层,在应用层。
二、定位根因:ios-print 重复 clientId
翻 activemq.log,满屏的Stealing link for clientId ios-print。
Stealing link for clientId ios-print from tcp://59.50.11.69:63080 Stealing link for clientId ios-print from tcp://59.50.11.69:63081 Stealing link for clientId ios-print from tcp://36.109.251.214:23456 ...每秒 7 次。所有 iOS 打印机设备都在用同一个 clientId = “ios-print”连接 MQTT。
MQTT 协议的坑
MQTT 协议规定 clientId 必须唯一。当第二个连接用相同 clientId 连上来时,ActiveMQ 默认行为是“Stealing link”—— 踢掉旧连接,新连接顶上。
正常情况下这没什么,但如果有几千台设备共用一个 clientId,就会发生:
设备A连接 → 注册 ios-print 设备B连接 → 踢掉设备A → 设备A断线重连 → 踢掉设备B 设备C连接 → 踢掉设备A → 设备B重连 → 踢掉设备C ...几千台设备互相踢,形成连接风暴。每次踢连接都会创建新线程、销毁旧线程,叠加closeAsync=false(同步关闭)和maxInactivityDuration=10s(超时断连太快),线程数越积越多,最终拖垮进程。
为什么会有相同 clientId?
客户端代码写死了clientId = "ios-print",每台 iOS 设备都用这个值连。正确做法是每台设备生成唯一 clientId(比如用设备序列号),但客户端发版周期长,服务端必须先扛住。
三、第一次尝试:allowLinkStealing=“false”(失败)
最直觉的想法:既然 Stealing link 是问题,那禁止 Stealing 不就行了?
ActiveMQ 有个参数叫allowLinkStealing="false",放在<transportConnector>上。
结果:启动失败。原因是这个参数是 transport 层的 URI 参数,不是 XML 属性,而且 ActiveMQ 6.x 对这个参数的支持方式变了。直接加在 XML 标签上会解析报错。
教训:ActiveMQ 的 transport 参数要放在 URI 里,不是 XML 属性里。
四、第二次尝试:自定义插件 RejectDuplicateClientIdPlugin
既然配置层面做不到,就写插件。思路很简单:
- 写一个 BrokerPlugin,拦截
addConnection() - 如果 clientId 在黑名单里(如 ios-print)且已有活跃连接,直接拒绝
- 拒绝时同步关闭底层 transport,减少资源泄漏
- 限频日志,不要每次拒绝都打日志
插件核心逻辑
publicvoidaddConnection(ConnectionContextcontext,ConnectionInfoinfo)throwsException{StringclientId=info.getClientId();if(clientId!=null&&blockedClientIds.contains(clientId)){if(!activeBlockedIds.add(clientId)){// 已存在 → 拒绝logRateLimited(clientId);// 限频日志stopTransportSync(context);// 同步关 I/OthrownewIllegalStateException("clientId '"+clientId+"' already connected");}}super.addConnection(context,info);}插件设计的几个细节
为什么不用InvalidClientIDException?
因为这个异常类在 activemq-broker 包里不一定有,用标准的IllegalStateException更稳妥。
为什么要stopTransportSync()同步关?
抛异常后 ActiveMQ 的 catch 块还会调一次stop(),如果 transport 没关,就走完整的同步清理路径,很慢。先关了 I/O,后面的 stop() 就是 no-op。
为什么限频日志?
每秒 13 次拦截,每次 3 行日志(WARN + 堆栈),一天能写满磁盘。改成每 30 秒输出一条汇总:
[RejectDuplicateClientIdPlugin] 拒绝 'ios-print' 重复连接 (最近30秒内 735 次)配套:log4j2 日志抑制
插件自己的日志好控制,但 ActiveMQ 内部在连接被拒绝时还会打一堆 WARN(Failed to add Connection、Stopping、Transport Connection failed),需要在 log4j2.properties 里把相关 logger 调到 ERROR:
logger.transport.name = org.apache.activemq.broker.TransportConnection logger.transport.level = ERROR logger.transportConnector.name = org.apache.activemq.broker.TransportConnector logger.transportConnector.level = ERROR logger.mqttConverter.name = org.apache.activemq.transport.mqtt.MQTTProtocolConverter logger.mqttConverter.level = ERROR第三条是后来加的——批量断连时每个死连接会产生 32 行堆栈(Broken pipe),15 秒能产生 3000+ 行日志。
五、最大的坑:URI 多行书写导致参数全部失效
插件部署了,配置也改了,但线上日志显示maximumConnections=10000、closeAsync=false——都是默认值,配置文件里写的参数一个都没生效。
折腾了很久才发现原因:activemq.xml 中 URI 跨多行书写,换行和缩进空格被原样包含在 URI 字符串中。
<!-- ❌ 错误写法:URI 跨多行,空格被包含进属性值 --><transportConnectorname="mqtt+nio"uri="mqtt+nio://0.0.0.0:1883? maximumConnections=20000&wireFormat.maxInactivityDuration=10000&..."/>ActiveMQ 拿到的实际 URI 是:
mqtt+nio://0.0.0.0:1883?\n maximumConnections=20000&\n wireFormat.maxInactivityDuration=10000&\n ...参数名前面带了一堆空格(日志里显示为%20%20%20%20%20),URI 解析器匹配不上,全部回退默认值。
正确写法:URI 必须在一行内。
<!-- ✅ 正确写法:URI 一行写完 --><transportConnectorname="mqtt+nio"uri="mqtt+nio://0.0.0.0:1883?maximumConnections=33000&transport.useInactivityMonitor=true&wireFormat.maxInactivityDuration=60000&transport.closeIdleConnectionTimeout=60000&transport.keepAliveTime=15000&transport.soKeepAlive=true&transport.tcpNoDelay=true&transport.closeAsync=true&allowLinkStealing=false"/>这个坑浪费了好几天。每次改了配置重启,以为生效了,实际全是默认值在跑。验证参数是否生效的唯一方法:看
Listening for connections日志,对比参数值。
六、更深层的问题:NIO 线程池无限制
插件生效、closeAsync=true、allowLinkStealing=false都配置正确后,连接数稳定在 3 万+,线程数正常 130 左右。但每次重启后 8-10 分钟,必然出现一次线程爆炸:
15:48:47 THR=105 CONN=30138 CPU=26% ← 正常 15:50:28 THR=1822 CONN=30564 CPU=61% ← 突然跳升 16:02:24 THR=15925 CONN=30937 CPU=236% ← 线程爆炸 16:06:16 THR=21442 CONN=32645 CPU=340% ← 濒临崩溃327 个新连接产生了 1685 个新线程(5 线程/连接),严重不正常。
原因
ActiveMQ 的 NIO Selector 线程池没有上限。正常情况下 8 个线程处理 3 万连接没问题,但重连风暴来临时(每分钟 1 万+ 新连接),线程池处理不过来就疯狂创建新线程,从 100 涨到 20000+,CPU 被线程调度耗尽,实际业务处理趋近于 0。
每次重启后所有客户端同时重连,必然触发这个风暴——形成"重启 → 风暴 → 崩溃 → 重启"的死循环。
修复:限制 NIO 线程池
在 setenv 的 JVM 参数中加三行:
-Dorg.apache.activemq.transport.nio.SelectorPool.corePoolSize=8-Dorg.apache.activemq.transport.nio.SelectorPool.maximumPoolSize=32-Dorg.apache.activemq.transport.nio.SelectorPool.workQueueCapacity=2000效果:重连风暴来临时,NIO 线程最多 32 个,超出的连接在 2000 容量的队列里排队。队列满了新连接被maximumConnections拒绝,而不是无限创建线程拖垮进程。
NIO 的设计就是少量线程处理大量连接。3 万连接 / 8 个核心线程 = 每线程 3750 连接,这是 NIO 的正常工作模式。32 的上限对于 4 核 CPU 已经很充裕了。
七、监控脚本的进化
整个过程中监控脚本也迭代了好几版,记录一下踩过的坑:
坑1:CPU 读数为 0
v3 版本用CPU=$(get_cpu)赋值,但$(...)是子 shell,里面修改的全局变量(_PREV_TICKS、_PREV_SEC)传不出来,导致 CPU 计算永远是 0。
修复:直接调用函数不用$(),函数里改全局变量,调用者直接读。
坑2:grep -c || echo 0输出两行
COUNT=$(grep-cpatternfile||echo0)当 grep 匹配 0 次时,grep -c输出0然后 exit 1,触发|| echo 0,结果变量里是0\n0(两行),后面的整数比较直接报错。
修复:改成grep -c pattern file || true,再用safe_num函数处理。
坑3:ss命令缺-a
ss -Ht -n只统计 ESTABLISHED 状态的连接,TIME_WAIT、CLOSE_WAIT 都漏掉了,连接数统计少了一大截。
修复:加-a参数。
坑4:fork 太多,OOM 时自己先挂
v3 版本每 30 秒要 fork 15 次(ss、grep、ps、free 等),高峰期内存 25MB。系统 OOM 时,监控脚本因为频繁 fork 先被 kill,反而起不到监控作用。
修复 v4:能读/proc就不用命令。CPU、RSS、线程数、系统内存全部用 bash 读/proc文件,0 fork。连接数用 awk 读/proc/net/tcp,1 fork。每周期从 15 fork 降到 2 fork,内存从 25MB 降到 1MB。
坑5:重启后不验证插件
监控脚本触发重启后只检查 PID 是否存在,不验证插件是否加载。曾经出现过重启后配置丢失、插件没加载的情况,监控脚本以为"启动成功",实际系统在裸奔。
修复:重启后 grep 插件日志,确认加载成功。
八、最终状态
| 指标 | 故障峰值 | 优化后稳定值 |
|---|---|---|
| 连接数 | 45,529 | ~30,000 |
| 线程数 | 27,231 | 130-400 |
| CPU | 359% | 26-60% |
| Stealing link (ios-print) | ~7次/秒 | 0(插件拦截) |
| Timer already cancelled | 大量 | 0 |
| 进程崩溃频率 | 每53分钟一次 | 0 |
生效的配置清单
activemq.xml (mqtt+nio, URI 一行内):
maximumConnections=30000wireFormat.maxInactivityDuration=60000(60秒)transport.closeIdleConnectionTimeout=60000transport.keepAliveTime=15000transport.closeAsync=trueallowLinkStealing=falseRejectDuplicateClientIdPlugin插件(拦截 ios-print)
setenv (JVM 参数):
-Xms2048M -Xmx2048M-XX:+UseG1GC -XX:MaxGCPauseMillis=200-XX:+HeapDumpOnOutOfMemoryErrorSelectorPool.corePoolSize=8SelectorPool.maximumPoolSize=32SelectorPool.workQueueCapacity=2000
log4j2.properties:
- TransportConnection → ERROR(抑制拦截连接日志)
- TransportConnector → ERROR(抑制连接超限日志)
- MQTTProtocolConverter → ERROR(抑制 Broken pipe 堆栈)
九、经验教训
参数是否生效,看日志,不要猜。
Listening for connections是最可靠的验证方式。URI 跨行这种坑,不看日志永远发现不了。MQTT clientId 唯一是底线。共用 clientId 等于自毁,服务端再怎么优化都只是缓解。客户端修复(每台设备生成唯一 clientId)才是根本解。
NIO 线程池一定要加上限。默认无上限在高并发场景下就是自杀按钮。对于 4 核机器,32 个 NIO 线程绰绰有余。
监控脚本本身也是系统的一部分。OOM 时监控脚本不能先挂,要用最少的资源完成采集。bash 读
/proc比调用命令靠谱得多。每加一个 fix 都要验证副作用。比如加插件后要确认日志量不会爆炸,改线程池后要确认吞吐不受影响。
回滚脚本和部署脚本同样重要。每次部署脚本都配对应的回滚脚本,出问题能 30 秒回滚,比什么都强。
附:相关文件索引
| 文件 | 说明 |
|---|---|
RejectDuplicateClientIdPlugin.java | 自定义插件源码 |
deploy_reject_plugin.sh | 插件一键编译部署脚本 |
rollback_reject_plugin.sh | 插件回滚脚本 |
monitor_activemq_v4.sh | v4 监控脚本(OOM 安全轻量版) |
activemq.xml | 最终版配置(URI 单行) |
MqttClient_fixed.cs等 | 客户端修复代码(C#) |