EMQX 大规模连接集群下的会话诊断利器:emqx_session_tool 使用与原理全解
【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx
EMQX 在承载数万甚至数十万并发连接的集群中,运维人员经常需要回答一个朴素却棘手的问题:"到底哪些客户端积压了消息、哪些在丢消息?" 逐个翻页浏览客户端列表不现实,为此 EMQX 在 6.x 中新增了面向运维的诊断模块emqx_session_tool,可从远端控制台直接调用,按任意会话指标(如mqueue_len、mqueue_dropped、inflight_cnt)找出 Top-K 会话。本文基于当前仓库源码,完整讲解该模块的 API、指标清单、参数语义、集群聚合机制与底层安全设计,让你在线上集群中快速定位"问题会话"。
一、为什么需要 emqx_session_tool
EMQX 的会话信息保存在每节点的 channel registry(emqx_channel_infoETS 表,见 emqx_cm.hrl)中。当节点上有数万连接时:
- 手动分页遍历客户端列表来寻找积压会话,操作成本高、效率低;
- 常规的管理接口往往需要拉取全量数据再在客户端侧排序,内存与带宽开销大;
- 直接向连接进程发消息查询状态会给承载业务流的进程增加额外负载。
emqx_session_tool正是为解决这一痛点而生的:它在节点本地流式扫描channel registry,只保留一个有界的 Top-K 有序集合,并且读取的是 registry 中已缓存的会话指标快照,不向任何连接进程发消息,因此对线上集群足够安全(详见 emqx_session_tool.erl 的模块文档)。
二、快速上手:从远端控制台定位 Top-K 会话
emqx_session_tool的设计目标是"可从远端控制台(remote console)在运行的集群上直接调用"。最简用法:
%% 找出本地节点上消息队列最长的 20 个会话 emqx_session_tool:top_by(mqueue_len). %% 找出丢消息最多的会话 emqx_session_tool:top_by(mqueue_dropped). %% 找出飞行窗口(inflight)占用最高的会话 emqx_session_tool:top_by(inflight_cnt).top_by/1等价于top_by(Metric, #{}),即使用全部默认参数。返回结果按指标值从高到低排序,每行是一个 map,包含clientid、pid、node、metric、value等字段。
如果集群有多个节点,使用cluster_top_by/1一次性聚合所有运行节点的结果:
%% 在整个集群范围内找出 mqueue_len 最大的 Top-20 会话 emqx_session_tool:cluster_top_by(mqueue_len).每个节点的扫描只遍历本节点自身的会话集合(不存在跨节点 ETS 遍历),各节点的 Top-K 集合在发起节点合并后重新裁剪为全局 Top-K;扫描失败的节点会被跳过,其结果行不出现(详见 emqx_session_tool.erl)。
三、可排序的指标(metric)全清单
top_by/1、scan/1等函数接受的指标,可通过emqx_session_tool:available_metrics/0实时查询,它由两类指标拼接而成(见 emqx_session_tool.erl)。
第一类:会话级 gauge(?SESSION_STATS_KEYS),与 emqx_session_mem 中维护的会话统计对应:
| 指标 | 含义 |
|---|---|
subscriptions_cnt/subscriptions_max | 当前订阅数 / 订阅数上限 |
inflight_cnt/inflight_max | 当前飞行中消息数 / 飞行窗口上限 |
mqueue_len/mqueue_max | 当前消息队列长度 / 队列上限 |
mqueue_dropped | 因队列满等原因丢弃的消息数(累计) |
total_payload_bytes | 队列中累积的载荷总字节数 |
awaiting_rel_cnt/awaiting_rel_max | 等待 PUBREL 确认的消息数 / 上限 |
注意:会话统计中的durable(布尔量)与next_pkt_id(回绕的报文 ID 计数器)被有意排除——对它们排序没有实际意义(见 emqx_session_tool.erl 的注释)。
第二类:channel 报文/消息计数器(?CHANNEL_METRICS),定义于 emqx_channel.hrl:
recv_pkt、recv_msg、recv_msg.qos0、recv_msg.qos1、recv_msg.qos2、recv_msg.dropped、recv_msg.dropped.await_pubrel_timeoutsend_pkt、send_msg、send_msg.qos0、send_msg.qos1、send_msg.qos2、send_msg.dropped、send_msg.dropped.expired、send_msg.dropped.queue_full、send_msg.dropped.too_large
这些计数器可用于回答"哪些客户端发送/接收消息最多、哪些消息被丢弃(过期/队列满/过大)"等问题。
传入不在清单内的指标会直接报错:error({unsupported_metric, Metric, available_metrics()});scan/1未提供metric选项则报error({missing_required_option, metric})(见 emqx_session_tool.erl)。
四、scan_opts 选项详解
top_by/2、scan/1、cluster_top_by/2都接受同一个选项 map(类型为scan_opts()),完整定义见 emqx_session_tool.erl:
| 选项 | 默认值 | 含义 |
|---|---|---|
metric | (必填) | 用于排序的会话指标,必须是available_metrics/0之一 |
top_k | 20 | 返回的行数上限,即最终保留的 Top-K 规模 |
min_value | 1 | 排除指标值低于该阈值的会话;默认 1 意味着默认过滤掉值为 0 的会话 |
chunk | 1000 | 每次ets:select批处理的行数 |
sleep_ms | 50 | 每处理完一批后休眠的毫秒数(yield 给其他进程) |
extra_keys | [] | 附加到结果行的缓存信息字段(从缓存的 session/clientinfo/conninfo 中解析),例如created_at、username、peername、connected_at、proto_ver |
extra_stats | [] | 附加到结果行的缓存统计字段,例如mqueue_len、total_payload_bytes、inflight_cnt |
rpc_timeout | 30000 | cluster_top_by/2时每节点的 RPC 超时(毫秒);单节点scan/1忽略此项 |
一个包含丰富上下文信息的调用示例:
emqx_session_tool:top_by(mqueue_len, #{ top_k => 10, min_value => 100, %% 只关心积压超过 100 条的会话 chunk => 500, %% 每批 500 行 sleep_ms => 10, %% 批间休眠 10ms,降低对调度的影响 extra_keys => [username, peername, connected_at], extra_stats => [mqueue_len, total_payload_bytes, inflight_cnt] }).结果行结构(row(),见 emqx_session_tool.erl):
#{ clientid := binary(), %% 客户端 ID pid := pid(), %% channel 进程 PID node := node(), %% 会话所在节点 metric := atom(), %% 本次排序所用指标 value := number(), %% 该会话的指标值(已还原为正数) extras => #{atom() => term()} %% 仅当指定了 extra_keys / extra_stats 时存在 }几点细节:
- 排序键设计为"数值大者优先,数值相同时按 clientid 升序",因此并列时会得到确定性的输出(测试用例
t_top_k_tie_breaks_by_row_key对此有专门验证); extra_stats保留的是参与排序时的那个快照,即使扫描结束后会话状态发生变化,结果行的extras也维持排名时读到的一致快照(见 emqx_session_tool.erl);extra_keys在每次选中 Top-K 胜出者后解析一次(best-effort),若会话在扫描期间已断开,这些字段可能为空。
五、底层原理:流式扫描 + 有界 Top-K
scan/1的实现路径是:scan_acc_new/1构造累加器 →scan_to_end循环推进 →scan_acc_rows/1收尾(见 emqx_session_tool.erl)。其核心安全保证可以总结为四条:
- 流式扫描,绝不构建全表列表:通过
emqx_utils_stream:ets/1以chunk大小的批次对emqx_channel_infoETS 表执行ets:select,每批只投影出{ClientId, ChanPid, Stats}三元组,较大的 info map 留在表内、只在 K 个胜出者身上解析(见 emqx_session_tool.erl); - 内存有界:无论会话总数多少,中间只维护一个容量为
top_k的gb_sets有序集合(heap_offer/5只在优于最差元素时才插入),输出规模恒定; - 批次间让出调度:每处理完
chunk行休眠sleep_ms毫秒;在scan_to_end的让步点还主动erlang:garbage_collect()回收刚处理完的批次,使长扫描保持平坦的堆(见 emqx_session_tool.erl); - 不打扰连接进程:指标直接读自 registry 中缓存的 stats proplist(
proplists:get_value(Metric, Stats, 0)),全程不给 channel 进程发消息。
数据新鲜度与覆盖范围(务必知晓)
- 新鲜度:读到的是连接进程最近一次发布的 stats 快照。连接进程按所在 zone 的 stats 定时器(即 idle timeout)刷新缓存,因此指标值最多可能滞后该间隔;若某 zone 的
stats.enable为false,其连接只在注册时发布一次快照,此时 gauge 反映的是连接建立时的状态(见 emqx_session_tool.erl)。 - 范围:仅覆盖注册在本地 channel registry 中的
emqx_session_mem会话;状态存放在 DS(durable storage)中的持久化会话暂未覆盖(源码以 TODO 标注,计划在支持后为行打上engine => mem | persistent_ds标记,见 emqx_session_tool.erl)。
六、增量扫描引擎:把扫描搬进 gen_server
除了"一次跑完"的scan/1,模块还导出了一组增量扫描 API,便于让 gen_server 之类的事件驱动进程自己持有游标、逐批推进:
scan_acc_new(Opts):构造初始累加器(同scan/1的选项,sleep_ms由驱动方负责);scan_acc(Acc):推进一个chunk批次,返回{continue, Acc}或{done, Acc};scan_acc_rows(Acc):把当前 Top-K 有序集合转成结果行,可在扫描中途调用——提前中止也能拿到"迄今最好"的部分结果,这就是中止安全的体现。
对应的注释与类型定义见 emqx_session_tool.erl。测试用例t_incremental_scan_acc验证了:部分推进即可读到部分结果、推进到完成与一次性的scan/1结果完全一致、对已完成累加器再次调用scan_acc/1幂等返回{done, _}(见 emqx_session_tool_SUITE.erl)。
七、集群聚合与上层封装
cluster_top_by 的聚合实现
cluster_top_by/2的流程(见 emqx_session_tool.erl):
- 取
emqx:running_nodes()作为目标节点集; - 通过
emqx_session_tool_proto_v1:scan(Nodes, Opts, Timeout)(即erpc:multicall,见 emqx_session_tool_proto_v1.erl)并行下发扫描; - 仅收集
{ok, Rows}的成功结果,按value降序整体排序; lists:sublist(Sorted, TopK)重新裁剪为全局 Top-K。
该 bpapi 在 EMQX6.0.3引入(introduced_in() -> "6.0.3")。
scanner / collector:供 Dashboard 等上层使用的封装
仓库中还有两个配套 gen_server,构成更完整的"异步 Top 扫描"能力,由emqx_session_top_proto_v1(6.3.0引入)桥接:
- emqx_session_top_scanner.erl:节点本地的增量扫描执行器。
start_scan/1接收选项(sort、count、batch_size默认 1000、sleep_ms默认 1),用emqx_session_tool:scan_acc_new/1+scan_acc/1按定时器逐批推进,完成或取消时把结果回报给 collector;同一时刻只允许一个扫描(第二个请求返回{error, {busy, Node}})。注意此处sort选项做了映射:mqueue_length -> mqueue_len、total_payload_bytes -> total_payload_bytes。 - emqx_session_top_collector.erl:集群级收集器,
run(Opts, CompletionFun)负责在本节点启动扫描并通过erpc:multicall分发到其他节点,汇总各节点回报(固定附加extra_stats => [mqueue_len, total_payload_bytes, inflight_cnt]),按sort/count排序后调用完成回调;还提供cancel/0(尽力而为地向所有节点广播取消)与status/0(返回running/completed/failed/cancelled/idle状态,含cluster_nodes、bad_replies等诊断字段)。
也就是说,emqx_session_tool是底层"纯函数"式的扫描内核,scanner/collector 是其异步化、集群化封装——将来 Dashboard 等管理界面展示"Top 会话"视图时,走的正是这条链路。
八、测试佐证:行为即契约
模块配套的 Common Test 套件 emqx_session_tool_SUITE.erl 从多个维度固化了上述行为,可作为理解语义的活文档:
t_top_by_ranks_by_metric/t_top_k_limits_result:按指标降序、受top_k约束;t_min_value_filter:默认min_value => 1会过滤掉值为 0 的会话,显式设置阈值可收紧;t_extra_keys_populate_extras:extra_keys会按session→clientinfo→conninfo→ 顶层 的顺序在缓存 info map 中查找字段,缺失键返回undefined;t_total_payload_bytes_stats_extras:extra_stats只附加到入选 Top-K 的行,且与排序快照一致;t_unsupported_metric_errors/t_missing_metric_errors:非法指标与缺失metric的报错行为;t_real_clients_scannable:真实连接的客户端(通过emqtt建连)可经缓存 stats 被扫描到;t_cluster_top_by:两节点集群中cluster_top_by(mqueue_len, #{top_k => 3})正确合并各节点 Top-K,且每行标注会话所属节点。
九、使用注意事项小结
- 调用位置:
emqx_session_tool面向远端控制台/诊断场景,直接在运行的集群节点上调用即可; - 指标选取:先执行
emqx_session_tool:available_metrics()确认指标名(会话级 gauge 与 channel 计数器是两类不同语义); - 大集群调参:连接数很多时,可通过增大
sleep_ms、调小chunk进一步压低对线上调度的影响;RPC 超时rpc_timeout需要给足(默认 30 秒),因为扫描是"边扫描边休眠"的自限速过程; - 数据新鲜度:结果反映的是连接进程最近一次发布的 stats 快照,在
stats.enable关闭的 zone 中可能停留在连接建立时刻; - 覆盖范围:目前仅覆盖内存会话(
emqx_session_mem),DS 持久化会话尚不在扫描范围内; - 只读安全:整个扫描过程只读 ETS、不发消息、不修改任何状态,适合在生产集群反复执行。
凭借"流式扫描 + 有界 Top-K + 缓存快照"这套设计,emqx_session_tool让运维人员在十万级连接的集群上也能以毫秒级开销、恒定内存代价,快速回答"谁在积压、谁在丢消息"这一日常诊断问题。
【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考