EMQX 会话上限超限后的重连恢复机制解析——基于 v5.8.5 行为修复 #14654【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx导读本文围绕 EMQX 仓库变更记录 fix-14654.en.md 所描述的行为修复展开当 Broker 达到最大会话数max session limit之后只要客户端携带的旧会话尚未过期、也尚未被清理该客户端依然可以成功重连。文章将结合emqx_cm、emqx_channel、emqx_stats等模块的源码讲清会话上限从何而来重连时会话如何被恢复陈旧会话记录如何被清理三条主线帮助你理解 MQTT 会话生命周期与 EMQX 连接准入逻辑的底层原理并给出可落地的配置与观测方法。一、变更条目速览一次关于会话上限与重连的行为修复关联文档 changes/ee/fix-14654.en.md 全文仅一句话Clients can now reconnect even if the max session limit is exceeded, as long as their old sessions are not yet expired or cleaned up.同一修复也记录在 changes/v5.8.5.en.md 的发布说明中[#14654] Clients can now reconnect successfully even if the maximum session limit has been reached, as long as their previous sessions remain active (i.e., not expired or cleaned up).这条变更表达了三层含义场景Broker 已经达到会话/连接上限正常情况下新的连接请求会被拒绝例外如果发起连接的客户端此前已有会话且该会话既未到期未超过session_expiry_interval也未被清理则这次重连应当被放行本质重连恢复一个已存在会话本质上不应被视为新建会话而占用配额因此在会话上限判定与连接准入逻辑上需要与之区分。下文将逐层拆解其背后的机制并给出源码级证据。二、先厘清概念EMQX 中的会话Session与连接Connection在展开修复逻辑之前需要区分两个容易混淆的概念连接Connection客户端与 Broker 之间的 TCP/TLS 网络通道对应一个 channel 进程连接断开即消亡。会话SessionMQTT 协议层面的状态集合包括订阅关系、未确认的 QoS 消息inflight、消息队列mqueue等由clientid唯一标识。会话的生命周期由session_expiry_interval控制客户端在 CONNECT 报文中携带该值若大于 0连接断开后会话仍保留持久会话在到期前客户端可用相同clientid重连并恢复状态若为 0连接断开即会话销毁。服务端侧的默认与上限配置位于 apps/emqx/src/emqx_schema.erl配置项默认值说明mqtt.session_expiry_interval2h会话过期时间类型为 durationmqtt.max_session_expiry_intervalinfinity客户端可请求的会话过期时间上限客户端重连时EMQX 需要判断该clientid是否已有未过期的会话如果有则执行会话恢复session present true而非新建会话。三、会话上限从何而来统计指标与限流基础设施3.1 会话数量的统计指标EMQX 通过emqx_stats统一维护计数器与峰值会话相关的指标定义在 apps/emqx/src/emqx_stats.erlsessions.count/sessions.max本节点会话数及其峰值cluster_sessions.count/cluster_sessions.max集群范围内会话数及其峰值disconnected_sessions.count已断开但会话尚未过期的会话数。这些指标由 apps/emqx/src/emqx_cm.erl 中的?CHAN_STATS声明映射到具体的 ETS 表{?CHAN_TAB, channels.count, channels.max}, {?CHAN_TAB, sessions.count, sessions.max}, {?CHAN_CONN_TAB, connections.count, connections.max}, {?CHAN_LIVE_TAB, live_connections.count,live_connections.max}, {?CHAN_REG_TAB, cluster_sessions.count,cluster_sessions.max}周期性统计函数stats_fun/0apps/emqx/src/emqx_cm.erl读取各 ETS 表大小并通过emqx_stats:setstat/3更新计数与峰值——只有当前值超过历史峰值时才更新*_max见 apps/emqx/src/emqx_stats.erl。这意味着sessions.max记录的是历史峰值而非硬性配额。3.2 真正的上限连接准入限流Broker 层面真正会拒绝新连接的硬性上限来自监听器限流limiter体系监听器级max_conn通过 emqx_ranch_limiter 实现达到上限时allow_capacity/1判定失败allow_limiter/1对超过连接速率的情况返回{close_connections_for, PAUSE_INTERVAL}即暂停接收新连接区域/客户端级限流max_conn、messages、bytes、delivery_messages、delivery_bytes、subscribes等限流器定义在 emqx_limiter_schema.erl会话相关限流名称见 emqx_limiter.erl-define(CHANNEL_LIMITS, [messages, bytes, subscribes]). -define(SESSION_LIMITS, [delivery_bytes, delivery_messages]). -define(LISTENER_LIMITS, [max_conn]). -define(ZONE_LIMITS, [max_conn, messages, bytes]).需要说明的是变更条目中max session limit更具体的实现例如企业版的会话配额能力不在当前开源仓库可见范围内从开源代码可以确认的是会话/连接达到上限时准入逻辑会拒绝新连接而本修复正是为已有未过期会话的重连打开放行通道。四、重连时发生了什么open_session的完整路径4.1 从 CONNECT 到会话打开客户端重连时channel 进程在 apps/emqx/src/emqx_channel.erl 中调用emqx_cm:open_session/4并根据返回结果决定 CONNACKcase emqx_cm:open_session(CleanStart, ClientInfo, ConnInfo, MaybeWillMsg) of {ok, #{session : Session, present : false}} - ok emqx_cm:register_channel(ClientId, self(), ConnInfo), ... handle_out(connack, {?RC_SUCCESS, sp(false), AckProps}, ensure_connected(NChannel)); {ok, #{session : Session, present : true, replay : ReplayContext}} - ok emqx_cm:register_channel(ClientId, self(), ConnInfo), ... handle_out(connack, {?RC_SUCCESS, sp(true), AckProps}, ensure_connected(NChannel)); {error, client_id_unavailable} - ReasonString THROTTLED, handle_out(connack, {?RC_SERVER_BUSY, ReasonString}, Channel); {error, Reason} - ... handle_out(connack, ?RC_UNSPECIFIED_ERROR, Channel) end.关键点present true表示该clientid存在可恢复的旧会话CONNACK 中spsession present置为 1并携带ReplayContext用于消息重放{error, client_id_unavailable}时客户端会收到RC_SERVER_BUSY0x83服务器繁忙附带的 reason string 为THROTTLED——这正是连接被节流拒绝的出口。4.2open_session内部注册表查询与节流判定核心实现在 apps/emqx/src/emqx_cm.erlopen_session(CleanStart, ClientInfo #{clientid : ClientId}, ConnInfo, MaybeWillMsg) - Pids emqx_cm_registry:lookup_all_channels(ClientId), {Local, Remote} lists:partition(fun(Pid) - node(Pid) : node() end, Pids), {LocalAlive, LocalDown} lists:partition(fun erlang:is_process_alive/1, Local), LocalStillDown drop_stale_local_rows(ClientId, LocalDown), case LocalStillDown of [_ | _] - %% At least one old session is in the middle of getting cleaned up. %% i.e. emqx_cm_pool is busy handling the async clean_down messages. %% Do not accept this client ID in this node. ?SLOG(warning, #{msg clientid_registration_throttled, ...}), {error, client_id_unavailable}; [] - do_open_session(CleanStart, ClientInfo, ConnInfo, MaybeWillMsg) end.处理流程分三步从全局 channel 注册表查出该clientid关联的所有 channel 进程按本节点/远端节点和存活/已死分组对本地已死进程调用drop_stale_local_rows/2做陈旧记录甄别与清理若清理后仍有正在清理中的旧会话则返回{error, client_id_unavailable}触发节流否则进入do_open_session执行真正的会话打开。4.3 会话打开新建还是恢复do_open_session/4apps/emqx/src/emqx_cm.erl依据CleanStart分支CleanStart true先discard_session丢弃旧会话并emqx_session:destroy然后emqx_session:create新建——这是全新开始语义CleanStart false调用emqx_session:open若存在可恢复会话则返回present true否则返回present false也即允许会话不存在的普通新连接。emqx_session:openapps/emqx/src/emqx_session.erl按 zone 配置构造会话参数max_subscriptions、max_awaiting_rel等并通过emqx_session_mem:open/4尝试接管既有会话apps/emqx/src/emqx_session_mem.erl 中先调用emqx_cm:takeover_session_begin(ClientId)找到旧 channel 则导出其会话状态、恢复resume、按新连接调整 inflight 窗口resize_inflight并应用新配置apply_conf——这就是旧会话未过期即可无缝恢复的底层实现。五、修复核心陈旧记录的甄别与清理drop_stale_local_rows5.1 两种已死进程行两种处置重连的成败取决于drop_stale_local_rows/2如何处置本地已死进程行apps/emqx/src/emqx_cm.erldrop_stale_local_rows(ClientId, LocalDownPids) - lists:filter( fun(DeadPid) - case do_get_chann_conn_mod(ClientId, DeadPid) of undefined - _ emqx_cm_registry:unregister_channel({ClientId, DeadPid}), false; _ConnMod - true end end, LocalDownPids ).情形 Atombstone墓碑记录——注册表中残留的进程行指向一个本节点从未注册过、或早已清理完毕的 pid此时chan-conn表中查不到对应的conn_mod条目do_get_chann_conn_mod返回undefined见 apps/emqx/src/emqx_cm.erl。这类记录是纯垃圾数据drop_stale_local_rows直接调用emqx_cm_registry:unregister_channel将其移除并返回 false不阻塞让重连继续。情形 B正在清理中的 channel——进程刚退出DOWN消息还排队在emqx_cm_pool中等待异步执行clean_down此时chan-conn表仍有对应条目返回 true 保留在LocalStillDown中触发client_id_unavailable节流。5.2 修复前后行为对比修复前本地已死进程行一律视为清理中重连被节流拒绝RC_SERVER_BUSY/THROTTLED即使旧会话本身完全有效——于是出现明明会话还没过期客户端却连不上的异常修复后tombstone 被主动清理只有真正的清理中场景才短暂节流旧会话未过期、未被清理的客户端可以正常重连并恢复会话。这一甄别逻辑与 apps/emqx/src/emqx_cm_registry.erl 的注册/注销流程相衔接register_channel通过mria写入集群共享的?CHAN_REG_TABunregister_channel2在删除记录的同时写入历史标记insert_hist_d供注册时的历史比对使用。5.3 集群视角的辅助约束在集群环境下重连还会受会话锁session locker约束。apps/emqx/src/emqx_cm_locker.erl 基于ekka_locker对clientid加锁锁获取失败同样返回{error, client_id_unavailable}锁策略由broker.session_locking_strategy配置决定local | leader | quorum | all。也就是说client_id_unavailable是准入被拒的统一出口而 #14654 的修复消除了其中陈旧记录误伤这一类误拒。六、会话未过期、未被清理的成立条件要让重连放行真正生效旧会话必须同时满足两个条件未过期CONNECT 携带的session_expiry_interval 0且尚未到期由 apps/emqx/src/emqx_schema.erl 中的mqtt.session_expiry_interval默认2h与mqtt.max_session_expiry_interval默认infinity共同约束。过期后会话被销毁重连只能新建会话。未被清理clean_down流程尚未执行完毕。进程退出后EMQX 通过 apps/emqx/src/emqx_cm.erl 中的emqx_cm_pool异步批量执行clean_down/1逐 pid 完成emqx_broker_helper:clean_down/2与do_unregister_channel。这一异步窗口正是清理中节流存在的技术原因。观测上disconnected_sessions.countapps/emqx/src/emqx_cm.erl统计连接已断开但会话未过期的数量——它由?CHAN_CONN_TAB与?CHAN_LIVE_TAB的表大小之差计算且被钳制在非负值会话接管瞬间连接行先删、存活行后删避免瞬时负值。该指标可用于判断当前集群中待恢复会话的规模。七、测试与验证会话接管语义的单元证据仓库测试用例印证了会话未过期即可被接管恢复的核心语义apps/emqx/test/emqx_channel_SUITE.erl 的t_handle_call_takeover_begin当expiry_interval 0时接管begin直接导致 channel 关闭会话随连接结束当把expiry_interval设为60000后接管begin返回会话数据会话得以交接同文件t_handle_call_takeover_endapps/emqx/test/emqx_channel_SUITE.erl验证会话在接管结束后存活并交由新连接继续使用apps/emqx/test/emqx_broker_SUITE.erl 的t_connected_client_count_transient_takeover覆盖接管风暴takeover storm场景下live_connections.count的瞬时波动与最终收敛。这些用例说明只要会话未过期重连走的始终是恢复 接管路径而不是新建路径——这正是 #14654 允许超限后重连的语义基础。八、运维视角配置与观测清单8.1 关键配置项配置路径默认值作用mqtt.session_expiry_interval2h会话保留时长决定未过期窗口mqtt.max_session_expiry_intervalinfinity服务端允许的最大会话过期时间broker.session_locking_strategy集群相关会话锁策略影响多节点重连的准入listenermax_conn/ limitermax_conn视配置监听器连接/速率上限触发RC_SERVER_BUSY的根源之一相关源码位置emqx_schema.erl、emqx_cm_locker.erl、emqx_ranch_limiter.erl。8.2 观测指标sessions.count/sessions.max本节点会话数与峰值cluster_sessions.count/cluster_sessions.max集群会话数与峰值disconnected_sessions.count等待恢复的会话规模日志关键词clientid_registration_throttled节流触发、more_than_one_channel_found注册表异常、listener_accept_refused_reached_max_connections监听器达上限。8.3 实践建议为关键业务客户端设置合理的session_expiry_interval避免因会话过早过期导致重连退化为新建会话从而在超限场景下被拒绝关注disconnected_sessions.count与clientid_registration_throttled日志的规模若节流频繁说明clean_down队列积压或会话数量逼近上限应结合业务评估扩容或调低会话保留时长升级到包含 #14654 修复的版本v5.8.5 及之后后tombstone 型陈旧记录会被自动清理无需人工介入若仍观察到会话未过期却连不上可从监听器max_conn、会话锁策略与集群注册表一致性三个方向排查。结语#14654 是一次小而关键的行为修复它在会话/连接已达上限的严格准入与旧会话未过期应当恢复的 MQTT 语义之间找到了正确的平衡点。通过drop_stale_local_rows对 tombstone 与清理中两种陈旧记录加以甄别EMQX 在 v5.8.5 之后既避免了残留注册数据对重连的误伤也保留了RC_SERVER_BUSY节流对瞬态竞态的防护。理解这条链路有助于在实际部署中精准定位重连失败类问题并合理设计会话保留策略。【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考