ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

EMQX 连接速率限制热更新原理:监听器更新后限流立即生效的机制解析 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 速率之后实际生效的限流可能表现得比预期更严格。这条修复涉及三个核心问题域监听器更新完成listener update后限流配置何时、由谁重新加载已建立的限流客户端client如何感知新配置而不继续沿用旧状态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/10s、10MB/hburst突发令牌量如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/2L117-L122会同时建立两个 groupcreate_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/3与post_config_update/5处理。当通过 API 或集群配置下发更新监听器时最终进入update_listener/4L405-L416update_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/2L124-L129会把监听器配置解析为限流选项并分别更新 listener 与 channel 两个 groupupdate_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/2L212-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_serverpersistent_term的组合register_group/3L68-L84将{module, limiter_options}写入persistent_termkey 为{?MODULE, Group}find_group/1L94-L101从persistent_term读取get_limiter_options/1L103-L110按{Group, Name}精确取出某个限流器的完整选项。persistent_term是 OTP 中读操作极快的进程内共享存储非常适合限流这种每次消费令牌都要读取配置的高频场景。README 中也明确说明connect/1是backed by a persistent term lookup的轻量操作。独占限流器每次消费都读取最新配置对连接内的消息/字节等独占限流器emqx_limiter_exclusive.erl 的try_consume/2L83-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可以推断旧行为的问题所在客户端状态中包含tokens、burst_tokens、last_time、last_burst_time等运行时状态而rate/burst属于配置参数。若配置变更只更新了注册表中的参数却没有正确处理桶状态与旧参数之间的衔接就会出现以下偏差客户端此前基于旧 burst 窗口累积了较少的burst_tokens配置调大 burst 后若客户端仍按旧窗口时间点last_burst_time判断是否可补充突发令牌就会继续表现为突发额度不足结果即实际限流比预期严格。新机制通过每次消费都从注册表读取最新 options并基于最新参数计算补充量见try_consume_regular/4与try_consume_burst/5中对capacity/interval/burst_capacity/burst_interval的使用配合update_group/2在监听器更新流程中的前置调用保证新配置在更新完成时点即对后续所有消费生效。配置格式与参数说明监听器级连接限流的配置项由 emqx_limiter_schema.erl 定义make_mqtt_limiters_schema/2L61-L77按Name _rate与Name _burst的命名规则生成字段。parse_rate/1L175-L201支持的正则语法为^(\d)(kb|mb|gb|b|)(/(\d*)([mshd]{1,2}))?$即速率/突发字符串支持infinity、纯数字10按 1 秒窗口解释、10/2s、10/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/s、5/500msmax_conn_burst{N, T}字符串突发窗口内的额外连接额度如10000/m不可为infinitymax_connections整数监听器最大并发连接数对应max_conn的总量上限非速率在 emqx_limiter.erl 的config/2L294-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/1L117-L162直接对应本修复其验证路径为以max_conn_rate 1/1s、max_conn_burst 1/1m启动监听器建立 2 个连接耗尽常规 突发额度第 3 个连接被拒绝等待 1 秒冷却允许 1 个新连接再拒绝 1 个突发未恢复通过emqx:update_config([listeners, Type, Name], {update, #{max_conn_burst 20/1m}})在线调大突发额度更新完成后立即并发建立 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])),这个用例在三个协议族tcp、ws、wss见 L15-L20上分别运行覆盖了共享限流器侧更新后突发连接立即放行的行为。限流子系统集成测试emqx_limiter_SUITE.erl 则覆盖更广的限流场景例如t_max_conn_listener/1L52-L61对监听器设置max_conn_rate 2/500ms后断言触发esockd_limiter_consume_pause事件连接接受被节流t_max_conn_zone/1L63-L72zone 级max_conn_rate行为验证t_max_message_rate_listener/1等监听器/zone 的消息速率与字节速率限流对应独占限流器路径事件limiter_exclusive_try_consume。这些测试通过 snabbkaffe 事件断言验证限流触发与恢复构成对配置生效行为的自动化保障。与 esockd 的衔接接受连接时的限流回调连接速率限流的实际执行点在 emqx_esockd_limiter.erl它实现esockd_generic_limiter行为create/1L47-L53把限流客户端直接作为状态保存consume/2L59-L77调用emqx_limiter_client:try_consume/2成功则放行新连接失败则返回{pause, 100, ...}让 esockd 暂停接受连接 100ms并记录listener_accept_throttled_due_to_quota_exceeded告警日志。由于限流客户端每次消费都从注册表读取最新配置监听器配置更新的那一刻起新到达的连接立即按新限流参数接受或节流无需重启监听器。esockd 层的暂停间隔为 100ms?PAUSE_INTERVAL配合共享桶的原子消费实现了对瞬时连接洪峰的平滑抑制。实践要点与排查建议在线调整 burst 可立即生效通过 Dashboard 或emqx:update_config([listeners, type, name], {update, #{max_conn_burst 20/1m}})调整监听器限流无需重启监听器更新完成后新连接立即按新参数执行。限流配置仅在变化时重注册update_group/2会对新旧配置做差集比较无变化则不触发重注册因此频繁下发相同配置不会产生额外开销。zone 级与监听器级限流叠加生效create_listener_limiter/3emqx_limiter.erl会将 zone 级与 listener 级的max_conn限流组合成emqx_limiter_composite两者都满足才放行。若调整后行为异常应同时检查zones.zone.mqtt.limiter.max_conn_rate/burst与监听器自身的限流配置。观察节流日志当连接被限流时日志会出现listener_accept_throttled_due_to_quota_exceeded告警可用于定位限流是否触发、以及是否需要上调 burst 参数。升级后行为差异本修复属于行为修正。升级前若观察到调大 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),仅供参考
返回列表