EMQX GCP Pub/Sub 消费端连接测试引发的健康状态误判修复(18193)深度解析
2026/9/24 19:57:23 网站建设 项目流程
  • 后端
  • 物联网
  • 消息队列
  • 通信

【免费下载链接】emqx

The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles

项目地址:https://gitcode.com/gh_mirrors/em/emqx
点击查看免费下载

导读

本文围绕 EMQX 仓库中 changes/ee/fix-18193.en.md 所记录的缺陷修复展开,剖析一个影响 GCP Pub/Sub Consumer(消费端)数据源的隐蔽问题:当用户在连接器或数据源上执行"Test Connection"(连接测试/探针)之后,正在运行的消费端数据源可能被错误地标记为disconnected(原因timeout),且一直保持该状态,直到手动禁用再重新启用。读完本文,你将掌握该问题的根因(探针临时 worker 池与运行中数据源共享健康状态簿记导致的状态串扰)、修复思路(按 worker 池隔离 optvar 键空间),以及配套热升级钩子的工作原理与对应的回归测试验证方式。

问题描述与影响范围

修复说明(见 changes/ee/fix-18193.en.md)原文指出:

在 GCP Pub/Sub Consumer 连接器或数据源上使用"Test Connection"后,一个正在运行的 GCP Pub/Sub Consumer 数据源可能显示为disconnected(原因timeout),并一直保持该状态,直到手动禁用并重新启用。

该缺陷影响6.1.3 和 6.2.2两个版本,同时也在 6.2.3 的变更日志中登记(见 changes/6.2.3.en.md)。从测试代码来看,该问题的追踪编号为 issue #18190,修复 PR 为 #18193。

关键语义澄清:这里的"健康状态误判"并不会真的中断消息拉取,而是让dashboard/API 层面的健康检查(health check)报告错误状态,进而误导运维人员认为数据源不可用。

根因分析:连接测试临时 worker 池污染运行中数据源的健康状态

背景:GCP Pub/Sub Consumer 的 worker 池与健康检查机制

EMQX 的 GCP Pub/Sub Consumer 数据源以连接器(connector)+ 数据源(source)的两层模型实现,核心实现位于 emqx_bridge_gcp_pubsub_impl_consumer.erl。每个数据源会通过emqx_resource_pool:start/3启动一个 worker 池,池中的每个 worker(emqx_bridge_gcp_pubsub_consumer_worker)负责:

  • 创建/校验 Pub/Sub 订阅(ensure_subscription_exists/1);
  • 长轮询拉取消息(do_pull_async/1);
  • 上报订阅状态。

worker 的健康状态通过optvar(进程间可变状态存储,见 emqx_utils 中 optvar 模块)发布。核心键定义在 emqx_bridge_gcp_pubsub_consumer_worker.erl:

%% Must be scoped by the source resource id: worker indices repeat across pools, and a %% bare-index key would let one pool (e.g. a probe's) clear or overwrite another's flag. -define(OPTVAR_SUB_OK(SOURCE_RES_ID, WORKER_ID), {?MODULE, subscription_ok, SOURCE_RES_ID, WORKER_ID} ).

健康检查路径(on_get_channel_status/3check_workers/3emqx_resource_pool:common_health_check_workers/2)会读取每个 worker 的 optvar 状态,参见 emqx_bridge_gcp_pubsub_impl_consumer.erl 与 emqx_bridge_gcp_pubsub_consumer_worker.erl:

health_check(SourceResId, WorkerId, HCTimeout) -> case optvar:read(?OPTVAR_SUB_OK(SourceResId, WorkerId), HCTimeout) of {ok, Status} -> Status; timeout -> timeout end.

optvar:read在超时时间内读不到状态时,健康检查即返回timeout,最终表现为数据源disconnected

缺陷链条:索引键碰撞导致探针池清理误删运行池状态

问题出在修复前的旧实现上。修复前的 worker 健康状态键只以 ecpool worker 索引(从 1 开始的整数)为键,即形如{emqx_bridge_gcp_pubsub_consumer_worker, subscription_ok, WorkerId},没有把 source 资源 ID 纳入键空间。

"Test Connection"(探针,probe)的执行流程是:EMQX 为探针创建一个临时的 worker 池,用与真实数据源相同的配置做一次干跑(dry-run)验证(见测试 t_probe_does_not_disturb_running_source/1 中 "Test Connection": a dry-run probe of the same source config 的注释),探针结束随即销毁临时池。

问题由此产生:

  1. 探针临时池的 worker 索引同样从 1 开始编号,与运行中数据源池的 worker 索引重复
  2. 由于旧键只含索引、不含池标识,探针池 worker 与运行池 worker共享同一批 optvar 键
  3. 探针池销毁时,其清理逻辑clear_optvar/2执行optvar:unset把运行中数据源 worker 的健康状态键一并清除
  4. 此后运行中数据源的健康检查读不到状态,超时后报告disconnected(reasontimeout),且因 worker 不会自行重新发布该状态,故障一直持续到手动禁用/重启用。

该分析在测试模块的-doc注释中有明确印证(见 emqx_bridge_gcp_pubsub_consumer_SUITE.erl):

每个池的 worker 索引都从 1 重新开始,因此仅含索引的键会与探针临时池共享;探针池销毁后清除了运行池的标志,使运行中的数据源卡在disconnected(健康检查超时)直到手动重启(issue #18190)。

修复方案:按池隔离 optvar 键空间

修复的核心思路非常清晰——让每个 worker 池各自维护独立的健康状态簿记,探针池的创建与清理不再影响运行中的数据源

具体改动体现在 emqx_bridge_gcp_pubsub_consumer_worker.erl:健康状态键从"仅 worker 索引"改为以 source 资源 ID 限定作用域的复合键:

-define(OPTVAR_SUB_OK(SOURCE_RES_ID, WORKER_ID), {?MODULE, subscription_ok, SOURCE_RES_ID, WORKER_ID} ).

由于探针池与运行池的source 资源 ID 必然不同(探针使用独立的临时资源 ID),两者写出的 optvar 键不再碰撞;探针池销毁时只会清理自己键空间内的状态,运行中数据源的健康状态完好无损。

配套地,worker 的清理与写入路径都改用复合键:

  • 写入:connect/1启动后optvar:set(?OPTVAR_SUB_OK(SourceResId, WorkerId), subscription_ok),见 emqx_bridge_gcp_pubsub_consumer_worker.erl;
  • 读取:health_check/3用同样的复合键读取,见 emqx_bridge_gcp_pubsub_consumer_worker.erl;
  • 清理:clear_optvar/2terminate/2中的clear_optvar调用,见 emqx_bridge_gcp_pubsub_consumer_worker.erl 与 emqx_bridge_gcp_pubsub_consumer_worker.erl。

另外值得一提的是,worker 池的停止路径stop_consumers1/2会按lists:seq(1, PoolSize)逐 worker 执行clear_optvar(见 emqx_bridge_gcp_pubsub_impl_consumer.erl)。修复后由于键已按 source 资源 ID 隔离,这一步清理只会作用于本池,天然安全。

热升级钩子:让旧版本启动的 worker 也获得新簿记

对于已经运行在旧版本代码(修复前 beam)上的 worker,它们写入的仍是旧格式的"索引键"。修复后的健康检查读取新格式的"资源 ID 复合键",读不到旧 worker 发布的状态,健康检查将一直失败——除非重启这些 worker。

为此修复附带了一个热升级(hot-upgrade)后置钩子,实现在 emqx_post_upgrade.erl:

pr_18193_gcp_pubsub_consumer_worker_optvars(_FromVsn) -> ... HasIndexKeyedOptvars = lists:any(IsIndexKeyed, optvar:list_all()), ConnResIds = [ ConnResId || HasIndexKeyedOptvars, ConnResId <- EMQXResource:list_instances_by_type(ConnImpl), case EMQXResource:get_instance(ConnResId) of {ok, _, #{status := stopped}} -> false; {ok, _, _} -> true; _ -> false end ], lists:foreach(fun EMQXResource:restart/1, ConnResIds), lists:foreach( fun(K) -> case IsIndexKeyed(K) of true -> optvar:unset(K); false -> ok end end, optvar:list_all() ), ok.

钩子执行两步操作:

  1. 重启受影响连接器:若系统里仍存在旧格式的索引键(说明还有旧版 beam 启动的 worker),则找出所有该类型的、状态非stopped的 connector 实例并逐个restart,让它们在新代码下重新启动 worker、按新格式发布健康状态;
  2. 清扫遗留索引键:把optvar:list_all()中所有旧格式的索引键unset,避免残留脏数据。

钩子同样被导出在 emqx_post_upgrade.erl 的导出列表,供升级流程统一调度。修复说明中的 "A hot-upgrade hook is included so consumers started by older versions are restarted to pick up the new bookkeeping" 正是对这一机制的概括。

回归测试验证

仓库为该修复提供了两个层面的回归测试,均在 emqx_bridge_gcp_pubsub_consumer_SUITE.erl 中:

1. 探针不干扰运行中数据源(t_probe_does_not_disturb_running_source)

测试用例 t_probe_does_not_disturb_running_source/1 的验证步骤:

  1. 订阅gcp_pubsub_consumer_worker_subscription_ready事件,等待 worker 就绪;
  2. 创建 connector 与 source,确认其状态为connected
  3. 调用probe_source_api(即 "Test Connection" 干跑探针),期望返回 204;
  4. 再次执行健康检查health_check_channel,断言状态仍为connected
  5. 通过get_source_api再次确认 REST 接口返回{"status": "connected"}

该用例直接复现了 issue #18190 的故障场景,并验证修复后探针池的销毁不再影响运行池。

2. 热升级钩子(t_post_upgrade_pr_18193)

测试用例 t_post_upgrade_pr_18193/1 模拟了"旧版 worker 遗留索引键"的升级场景:

  1. 创建两个 connector/source 对,其中一对会被禁用;
  2. 手工向 optvar 写入旧格式索引键{emqx_bridge_gcp_pubsub_consumer_worker, subscription_ok, 1},模拟旧版 worker 发布的状态;
  3. 禁用另一对连接器;
  4. 调用emqx_post_upgrade:pr_18193_gcp_pubsub_consumer_worker_optvars("vsn")
  5. 断言:受影响的 connector 经历stop_enter → start重启流程;系统中不再残留任何旧格式索引键。

(测试模块开头还定义了旧格式键的宏?OPTVAR_SUB_OK(X) = {emqx_bridge_gcp_pubsub_consumer_worker, subscription_ok, X},用于模拟旧版本键格式,见 emqx_bridge_gcp_pubsub_consumer_SUITE.erl。)

此外,基于 gRPC 的 consumer 变体同样有对应的探针回归用例t_probe_does_not_disturb_running_source(见 emqx_bridge_gcp_pubsub_consumer_grpc_SUITE.erl),说明修复覆盖了 REST/gRPC 两条消费通道。

运维视角:如何判断与规避

  • 版本核查:如果你正运行 EMQX 6.1.3 或 6.2.2,且使用 GCP Pub/Sub Consumer 数据源,应尽快升级到包含 #18193 修复的版本(如 6.2.3 及以后)。修复详情可参见 changes/6.2.3.en.md。
  • 现象识别:若在 dashboard 对 GCP Pub/Sub 连接器或数据源点击"测试连接"后,某运行中的数据源变为disconnected(reasontimeout)且无法自愈,即为该缺陷的典型症状;旧版本中临时恢复手段是手动禁用再启用该数据源(与修复说明描述一致)。
  • 升级后的自动修复:热升级到修复版本后无需人工逐个重启数据源——pr_18193_gcp_pubsub_consumer_worker_optvars钩子会自动重启受影响连接器并清扫遗留状态键。

小结

fix-18193是一个典型的"状态簿记键空间设计缺陷"案例:连接测试(探针)临时 worker 池与运行中数据源池共享了未加作用域的健康状态键,探针池销毁时的清理误删了运行池的状态标志。修复通过将 optvar 键从"仅 worker 索引"改为"source 资源 ID + worker 索引"的复合键实现池级隔离,并辅以热升级钩子重启旧版 worker、清扫遗留脏键,最终由两个回归测试用例锁定行为。该修复同时也提醒数据集成插件开发者:任何跨进程共享的可变状态键,都必须显式携带所属实例的作用域标识

参考资料

  • 修复记录:changes/ee/fix-18193.en.md
  • 版本变更日志:changes/6.2.3.en.md
  • Consumer 桥接实现:emqx_bridge_gcp_pubsub_impl_consumer.erl
  • Consumer worker(optvar 键与健康检查):emqx_bridge_gcp_pubsub_consumer_worker.erl
  • 热升级钩子:emqx_post_upgrade.erl
  • 回归测试:emqx_bridge_gcp_pubsub_consumer_SUITE.erl、emqx_bridge_gcp_pubsub_consumer_grpc_SUITE.erl
  • 资源管理基础设施:emqx_resource.erl、emqx_resource_pool.erl
  • 后端
  • 物联网
  • 消息队列
  • 通信

【免费下载链接】emqx

The most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles

项目地址:https://gitcode.com/gh_mirrors/em/emqx
点击查看免费下载

相关推荐

上一篇:Houston:零配置Meteor管理工具,让Django Admin体验在Meteor中重生
下一篇:floccus缓存机制:CacheTree如何提升同步性能

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

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

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

立即咨询