☰
ActiveMQ MQTT 连接风暴排查手记
2026/10/1 15:39:14 网站建设 项目流程

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 34819

3 万+ 连接,而当时配置的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

既然配置层面做不到,就写插件。思路很简单:

  1. 写一个 BrokerPlugin,拦截addConnection()
  2. 如果 clientId 在黑名单里(如 ios-print)且已有活跃连接,直接拒绝
  3. 拒绝时同步关闭底层 transport,减少资源泄漏
  4. 限频日志,不要每次拒绝都打日志

插件核心逻辑

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&amp;wireFormat.maxInactivityDuration=10000&amp;..."/>

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&amp;transport.useInactivityMonitor=true&amp;wireFormat.maxInactivityDuration=60000&amp;transport.closeIdleConnectionTimeout=60000&amp;transport.keepAliveTime=15000&amp;transport.soKeepAlive=true&amp;transport.tcpNoDelay=true&amp;transport.closeAsync=true&amp;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,231130-400
CPU359%26-60%
Stealing link (ios-print)~7次/秒0(插件拦截)
Timer already cancelled大量0
进程崩溃频率每53分钟一次0

生效的配置清单

activemq.xml (mqtt+nio, URI 一行内):

  • maximumConnections=30000
  • wireFormat.maxInactivityDuration=60000(60秒)
  • transport.closeIdleConnectionTimeout=60000
  • transport.keepAliveTime=15000
  • transport.closeAsync=true
  • allowLinkStealing=false
  • RejectDuplicateClientIdPlugin插件(拦截 ios-print)

setenv (JVM 参数):

  • -Xms2048M -Xmx2048M
  • -XX:+UseG1GC -XX:MaxGCPauseMillis=200
  • -XX:+HeapDumpOnOutOfMemoryError
  • SelectorPool.corePoolSize=8
  • SelectorPool.maximumPoolSize=32
  • SelectorPool.workQueueCapacity=2000

log4j2.properties:

  • TransportConnection → ERROR(抑制拦截连接日志)
  • TransportConnector → ERROR(抑制连接超限日志)
  • MQTTProtocolConverter → ERROR(抑制 Broken pipe 堆栈)

九、经验教训

  1. 参数是否生效,看日志,不要猜。Listening for connections是最可靠的验证方式。URI 跨行这种坑,不看日志永远发现不了。

  2. MQTT clientId 唯一是底线。共用 clientId 等于自毁,服务端再怎么优化都只是缓解。客户端修复(每台设备生成唯一 clientId)才是根本解。

  3. NIO 线程池一定要加上限。默认无上限在高并发场景下就是自杀按钮。对于 4 核机器,32 个 NIO 线程绰绰有余。

  4. 监控脚本本身也是系统的一部分。OOM 时监控脚本不能先挂,要用最少的资源完成采集。bash 读/proc比调用命令靠谱得多。

  5. 每加一个 fix 都要验证副作用。比如加插件后要确认日志量不会爆炸,改线程池后要确认吞吐不受影响。

  6. 回滚脚本和部署脚本同样重要。每次部署脚本都配对应的回滚脚本,出问题能 30 秒回滚,比什么都强。


附:相关文件索引

文件说明
RejectDuplicateClientIdPlugin.java自定义插件源码
deploy_reject_plugin.sh插件一键编译部署脚本
rollback_reject_plugin.sh插件回滚脚本
monitor_activemq_v4.shv4 监控脚本(OOM 安全轻量版)
activemq.xml最终版配置(URI 单行)
MqttClient_fixed.cs等客户端修复代码(C#)

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询