EMQX 连接速率限制热更新原理:监听器更新后限流立即生效的机制解析
2026/9/23 15:52:10 网站建设 项目流程

EMQX 连接速率限制热更新原理:监听器更新后限流立即生效的机制解析

【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx

导读

本文围绕 EMQX 的一条限流修复记录展开:监听器(Listener)配置更新完成后,连接速率限制(connection rate limits)必须立即生效。在旧实现中,内部限流器的部分状态不会随配置变更同步更新,典型表现是调大max_conn_burst后,实际生效的限流仍然比预期严格。文章将结合当前仓库源码,从限流子系统架构、配置更新触发链、令牌桶实现细节与测试验证四个层面,讲清楚"为什么旧行为会滞后"以及"现在如何做到即时生效",帮助读者正确配置与排查 EMQX 监听器级连接限流。

关联变更说明

本仓库changes/ee/fix-15794.en.md记录了该修复的原始描述:

Ensure that any changes to connection rate limits take effect immediately after the listener update has completed. Previously, parts of internal limiter state were not directly affected by configuration changes. For example, after increasing the burst rate, the effective rate limit could appear stricter than expected.

翻译过来即:确保连接速率限制的任何变更,在监听器更新完成之后立即生效。此前,内部限流器状态的一部分不会直接受到配置变更的影响。例如在提高 burst 速率之后,实际生效的限流可能表现得比预期更严格。

这条修复涉及三个核心问题域:

  1. 监听器更新完成(listener update)后,限流配置何时、由谁重新加载;
  2. 已建立的限流客户端(client)如何感知新配置,而不继续沿用旧状态;
  3. burst(突发)令牌为何在旧行为下"看起来更严格",新实现如何消除这一偏差。

EMQX 限流子系统总体架构

限流功能在仓库中的入口模块是 emqx_limiter.erl,其配套文档 apps/emqx/src/emqx_limiter/README.md 对整体设计做了精炼说明:

  • Limiter(限流器):一个令牌桶模型(token bucket)实体,通过全局唯一 ID{Group, Name}标识,例如{{zone, default}, messages}
  • Client(限流客户端):连接到某个 limiter 的消费方,通过emqx_limiter_client:try_consume/2消费令牌;
  • 两种限流语义
    • Shared limiter(共享限流器):连接到同一 limiter 的所有客户端共享同一个桶,协作消费令牌;
    • Exclusive limiter(独占限流器):每个连接到 limiter 的客户端各自持有独立的桶,只被该客户端独占消费。

一个 limiter 由两个参数刻画:

  • rate:令牌生成速率,如1000/10s10MB/h
  • burst:突发令牌量,如10000/h,用于在更长时间窗口内授予额外令牌,以应对突发流量。

监听器限流的命名空间

在 emqx_limiter.erl 中定义了各类限流名称宏:

-define(CHANNEL_LIMITS, [messages, bytes, subscribes]). -define(SESSION_LIMITS, [delivery_bytes, delivery_messages]). -define(CLIENT_LIMITS, ?CHANNEL_LIMITS ++ ?SESSION_LIMITS). -define(LISTENER_LIMITS, [max_conn]). -define(ZONE_LIMITS, [max_conn, messages, bytes]).

其中与"连接速率限制"直接相关的是LISTENER_LIMITS中的max_conn,对应监听器配置项max_conn_rate/max_conn_burst。监听器创建时,emqx_limiter:create_listener_limiters/2(L117-L122)会同时建立两个 group:

create_listener_limiters(ListenerId, ListenerConfig) -> ListenerLimiters = listener_limiter_options(ListenerConfig), ClientLimiters = client_limiter_options(ListenerConfig), ok = create_group(shared, listener_group(ListenerId), ListenerLimiters), ok = create_group(exclusive, channel_group(ListenerId), ClientLimiters).
  • listener_group(ListenerId){listener, ListenerId},承载shared类型的max_conn连接数限流;
  • channel_group(ListenerId){channel, ListenerId},承载exclusive类型的消息/字节/订阅限流。

也就是说,连接速率限制走的是共享限流器,而单连接内的消息速率等走的是独占限流器,二者语义不同,更新路径也有差异,这正是理解该修复的关键背景。

配置更新触发链:监听器更新后发生了什么

配置回调入口

监听器配置更新由 emqx_listeners.erl 的pre_config_update/3post_config_update/5处理。当通过 API 或集群配置下发更新监听器时,最终进入update_listener/4(L405-L416):

update_listener(Type, Name, OldConf, NewConf) -> ListenerId = listener_id(Type, Name), case is_running(Type, ListenerId, NewConf) of true -> ok = emqx_limiter:update_listener_limiters(ListenerId, NewConf), do_update_running_listener(Type, Name, OldConf, NewConf); false -> ok = maybe_unregister_ocsp_stapling_refresh(Type, Name, NewConf), restart_listener(Type, Name, OldConf, NewConf) end.

关键点在is_running/3分支:当监听器处于运行状态时,先调用emqx_limiter:update_listener_limiters(ListenerId, NewConf)更新限流器,再执行do_update_running_listener/4完成监听器本身的运行时更新。这样限流配置的刷新与监听器更新处于同一事务流程中,且限流更新排在前面,从而保证"更新完成即生效"。

注释中还特别指出一个边界情况:更新目标监听器可能并未运行(例如此前一次失败的更新导致其限流器被删除),此时走restart_listener/4重新启动,限流器由启动流程重建。这说明更新逻辑对"运行中"与"未运行"两类场景都做了兜底。

限流组更新:仅在配置变化时重注册

emqx_limiter:update_listener_limiters/2(L124-L129)会把监听器配置解析为限流选项,并分别更新 listener 与 channel 两个 group:

update_listener_limiters(ListenerId, ListenerConfig) -> ListenerLimiters = listener_limiter_options(ListenerConfig), ClientLimiters = client_limiter_options(ListenerConfig), ok = update_group(listener_group(ListenerId), ListenerLimiters), ok = update_group(channel_group(ListenerId), ClientLimiters).

update_group/2(L212-L225)的实现体现了"按需更新"的优化:

update_group(Group, Options) -> case emqx_limiter_registry:find_group(Group) of undefined -> error({limiter_group_not_found, Group}); {Module, OldOptions} -> Diff = lists:foldl(fun lists:delete/2, OldOptions, Options), Diff =/= [] andalso begin ok = emqx_limiter_registry:register_group(Group, Module, Options), ok = Module:update_group(Group, Options) end, ok end.

它先对比新旧选项,只有存在差异时才重新注册 group 并通知对应限流模块执行update_group/2回调;配置无变化时则跳过,避免不必要的开销。

注册表:新配置如何"实时"到达已有客户端

persistent_term 缓存

限流组配置缓存于 emqx_limiter_registry.erl,本质是一个gen_server+persistent_term的组合:

  • register_group/3(L68-L84)将{module, limiter_options}写入persistent_term,key 为{?MODULE, Group}
  • find_group/1(L94-L101)从persistent_term读取;
  • get_limiter_options/1(L103-L110)按{Group, Name}精确取出某个限流器的完整选项。

persistent_term是 OTP 中"读操作极快"的进程内共享存储,非常适合限流这种每次消费令牌都要读取配置的高频场景。README 中也明确说明connect/1是"backed by a persistent term lookup"的轻量操作。

独占限流器:每次消费都读取最新配置

对连接内的消息/字节等独占限流器,emqx_limiter_exclusive.erl 的try_consume/2(L83-L98)每次消费前都实时调用emqx_limiter_registry:get_limiter_options/1取回当前配置

try_consume(#{limiter_id := LimiterId} = State0, Amount) -> LimiterOptions = emqx_limiter_registry:get_limiter_options(LimiterId), Result = case try_consume(State0, Amount, LimiterOptions) of ...

由于注册表数据来自persistent_term,一旦update_group/2重注册完成,后续所有try_consume调用立即读到新参数——这就是"配置变更即时生效"在独占限流器侧的核心机制。模块注释(L49-L51)也说明了为何create_group/update_group/delete_group对独占限流器是 no-op:桶的状态在客户端进程侧,限流器自身只是配置的持有者

共享限流器:注册表驱动的桶

连接速率(max_conn)对应的共享限流器在 emqx_limiter_shared.erl 中实现。与独占限流器不同,共享桶的状态以原子值(atomic_value)形式存在于进程内共享区域,消费通过try_consume_accumulated_burst/2等原子操作完成,保证多客户端并发消费的一致性。无论哪种实现,配置读取都统一走emqx_limiter_registry,因此配置更新路径一致

burst 令牌为何曾"看起来更严格"

回到修复描述中的现象:"after increasing the burst rate, the effective rate limit could appear stricter than expected"。

结合emqx_limiter_exclusive.erl的状态结构(L36-L41)可以推断旧行为的问题所在:客户端状态中包含tokensburst_tokenslast_timelast_burst_time运行时状态,而rate/burst属于配置参数。若配置变更只更新了注册表中的参数,却没有正确处理桶状态与旧参数之间的衔接,就会出现以下偏差:

  • 客户端此前基于旧 burst 窗口累积了较少的burst_tokens
  • 配置调大 burst 后,若客户端仍按旧窗口时间点(last_burst_time)判断是否可补充突发令牌,就会继续表现为"突发额度不足";
  • 结果即"实际限流比预期严格"。

新机制通过每次消费都从注册表读取最新 options,并基于最新参数计算补充量(见try_consume_regular/4try_consume_burst/5中对capacity/interval/burst_capacity/burst_interval的使用),配合update_group/2在监听器更新流程中的前置调用,保证新配置在更新完成时点即对后续所有消费生效。

配置格式与参数说明

监听器级连接限流的配置项由 emqx_limiter_schema.erl 定义:make_mqtt_limiters_schema/2(L61-L77)按Name ++ "_rate"Name ++ "_burst"的命名规则生成字段。parse_rate/1(L175-L201)支持的正则语法为:

^(\d+)(kb|mb|gb|b|)(/(\d*)([mshd]{1,2}))?$

即速率/突发字符串支持infinity、纯数字(10,按 1 秒窗口解释)、10/2s10/500ms以及带容量单位的形式(如5kb/1m,用于字节类限流)。时间单位支持d/h/m/s/ms

一个典型的监听器连接限流配置如下:

listeners.tcp.default { bind = "0.0.0.0:1883" max_connections = 1024000 max_conn_rate = "1000/s" # 常规连接速率:每秒 1000 个 max_conn_burst = "10000/m" # 突发连接额度:每分钟最多额外 10000 个 }

参数含义:

配置项类型含义
max_conn_rateinfinity{N, T}字符串常规窗口内的连接建立速率,格式如1000/s5/500ms
max_conn_burst{N, T}字符串突发窗口内的额外连接额度,如10000/m;不可为infinity
max_connections整数监听器最大并发连接数(对应max_conn的总量上限,非速率)

在 emqx_limiter.erl 的config/2(L294-L322)中,配置被解析为三种限流选项之一:

  • #{capacity => infinity}:未配置限流(unlimited);
  • #{capacity, interval, burst_capacity => 0}:仅常规速率(limited);
  • #{capacity, interval, burst_capacity, burst_interval}:常规速率 + 突发(limited_with_burst)。

测试验证:更新后立即生效

面向修复的专项测试

emqx_listeners_limits_SUITE.erl 中的t_max_conn_rate_update/1(L117-L162)直接对应本修复,其验证路径为:

  1. max_conn_rate = "1/1s"max_conn_burst = "1/1m"启动监听器;
  2. 建立 2 个连接耗尽常规 + 突发额度,第 3 个连接被拒绝;
  3. 等待 1 秒冷却,允许 1 个新连接,再拒绝 1 个(突发未恢复);
  4. 通过emqx:update_config([listeners, Type, Name], {update, #{<<"max_conn_burst">> => <<"20/1m">>}})在线调大突发额度;
  5. 更新完成后立即并发建立 10 个连接,全部成功。
%% Update the limit, allowing for much higher bursts: ?assertMatch( {ok, _}, emqx:update_config( [listeners, Type, Name], {update, #{<<"max_conn_burst">> => <<"20/1m">>}} ) ), %% Connection burst should be allowed right after: Clients3 = emqx_utils:pmap( fun(_) -> emqtt_connect("127.0.0.1", Port, Config) end, lists:seq(1, 10) ), ?assertEqual([pong], lists:usort([emqtt:ping(C) || C <- Clients3])),

这个用例在三个协议族(tcpwswss,见 L15-L20)上分别运行,覆盖了共享限流器侧"更新后突发连接立即放行"的行为。

限流子系统集成测试

emqx_limiter_SUITE.erl 则覆盖更广的限流场景,例如:

  • t_max_conn_listener/1(L52-L61):对监听器设置max_conn_rate = "2/500ms"后,断言触发esockd_limiter_consume_pause事件(连接接受被节流);
  • t_max_conn_zone/1(L63-L72):zone 级max_conn_rate行为验证;
  • t_max_message_rate_listener/1等:监听器/zone 的消息速率与字节速率限流,对应独占限流器路径(事件limiter_exclusive_try_consume)。

这些测试通过 snabbkaffe 事件断言验证限流触发与恢复,构成对"配置生效"行为的自动化保障。

与 esockd 的衔接:接受连接时的限流回调

连接速率限流的实际执行点在 emqx_esockd_limiter.erl,它实现esockd_generic_limiter行为:

  • create/1(L47-L53)把限流客户端直接作为状态保存;
  • consume/2(L59-L77)调用emqx_limiter_client:try_consume/2:成功则放行新连接;失败则返回{pause, 100, ...}让 esockd 暂停接受连接 100ms,并记录listener_accept_throttled_due_to_quota_exceeded告警日志。

由于限流客户端每次消费都从注册表读取最新配置,监听器配置更新的那一刻起,新到达的连接立即按新限流参数接受或节流,无需重启监听器。esockd 层的暂停间隔为 100ms(?PAUSE_INTERVAL),配合共享桶的原子消费,实现了对瞬时连接洪峰的平滑抑制。

实践要点与排查建议

  1. 在线调整 burst 可立即生效:通过 Dashboard 或emqx:update_config([listeners, <type>, <name>], {update, #{<<"max_conn_burst">> => <<"20/1m">>}})调整监听器限流,无需重启监听器,更新完成后新连接立即按新参数执行。
  2. 限流配置仅在变化时重注册update_group/2会对新旧配置做差集比较,无变化则不触发重注册,因此频繁下发相同配置不会产生额外开销。
  3. zone 级与监听器级限流叠加生效create_listener_limiter/3(emqx_limiter.erl)会将 zone 级与 listener 级的max_conn限流组合成emqx_limiter_composite,两者都满足才放行。若调整后行为异常,应同时检查zones.<zone>.mqtt.limiter.max_conn_rate/burst与监听器自身的限流配置。
  4. 观察节流日志:当连接被限流时,日志会出现listener_accept_throttled_due_to_quota_exceeded告警,可用于定位限流是否触发、以及是否需要上调 burst 参数。
  5. 升级后行为差异:本修复属于行为修正。升级前若观察到"调大 burst 后限流仍偏严格",升级后应消失;若仍存在,可借助t_max_conn_rate_update同类用例(在apps/emqx/test/emqx_listeners_limits_SUITE.erl)复现验证。

小结

本修复的实质是打通"监听器配置更新 → 限流注册表重载 → 客户端实时读取新配置"这条链路:emqx_listeners:update_listener/4在运行中监听器上先调用emqx_limiter:update_listener_limiters/2update_group/2检测到差异后重注册persistent_term中的限流选项,随后所有限流客户端在每次令牌消费时都能立即读到最新参数,从而消除了"调大 burst 后实际限流仍偏严格"的滞后现象。无论共享限流器(连接速率)还是独占限流器(消息/字节速率),配置生效路径均统一由 emqx_limiter_registry.erl 承载,保证一致性。

如需深入阅读,推荐按以下路径继续探索:

  • 限流子系统设计文档:apps/emqx/src/emqx_limiter/README.md
  • 限流入口模块:apps/emqx/src/emqx_limiter/src/emqx_limiter.erl
  • 监听器更新触发链:apps/emqx/src/emqx_listeners.erl
  • 连接速率限流回调:apps/emqx/src/emqx_limiter/src/emqx_esockd_limiter.erl
  • 专项测试:apps/emqx/test/emqx_listeners_limits_SUITE.erl

【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询