后端物联网消息队列通信【免费下载链接】emqxThe 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 池池中的每个 workeremqx_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 probes) clear or overwrite anothers flag. -define(OPTVAR_SUB_OK(SOURCE_RES_ID, WORKER_ID), {?MODULE, subscription_ok, SOURCE_RES_ID, WORKER_ID} ).健康检查路径on_get_channel_status/3→check_workers/3→emqx_resource_pool:common_health_check_workers/2会读取每个 worker 的 optvar 状态参见 emqx_bridge_gcp_pubsub_impl_consumer.erl 与 emqx_bridge_gcp_pubsub_consumer_worker.erlhealth_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 的注释探针结束随即销毁临时池。问题由此产生探针临时池的 worker 索引同样从 1 开始编号与运行中数据源池的 worker 索引重复由于旧键只含索引、不含池标识探针池 worker 与运行池 worker共享同一批 optvar 键探针池销毁时其清理逻辑clear_optvar/2执行optvar:unset把运行中数据源 worker 的健康状态键一并清除此后运行中数据源的健康检查读不到状态超时后报告disconnectedreasontimeout且因 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/2与terminate/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.erlpr_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.钩子执行两步操作重启受影响连接器若系统里仍存在旧格式的索引键说明还有旧版 beam 启动的 worker则找出所有该类型的、状态非stopped的 connector 实例并逐个restart让它们在新代码下重新启动 worker、按新格式发布健康状态清扫遗留索引键把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 的验证步骤订阅gcp_pubsub_consumer_worker_subscription_ready事件等待 worker 就绪创建 connector 与 source确认其状态为connected调用probe_source_api即 Test Connection 干跑探针期望返回 204再次执行健康检查health_check_channel断言状态仍为connected通过get_source_api再次确认 REST 接口返回{status: connected}。该用例直接复现了 issue #18190 的故障场景并验证修复后探针池的销毁不再影响运行池。2. 热升级钩子t_post_upgrade_pr_18193测试用例 t_post_upgrade_pr_18193/1 模拟了旧版 worker 遗留索引键的升级场景创建两个 connector/source 对其中一对会被禁用手工向 optvar 写入旧格式索引键{emqx_bridge_gcp_pubsub_consumer_worker, subscription_ok, 1}模拟旧版 worker 发布的状态禁用另一对连接器调用emqx_post_upgrade:pr_18193_gcp_pubsub_consumer_worker_optvars(vsn)断言受影响的 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 连接器或数据源点击测试连接后某运行中的数据源变为disconnectedreasontimeout且无法自愈即为该缺陷的典型症状旧版本中临时恢复手段是手动禁用再启用该数据源与修复说明描述一致。升级后的自动修复热升级到修复版本后无需人工逐个重启数据源——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.mdConsumer 桥接实现emqx_bridge_gcp_pubsub_impl_consumer.erlConsumer workeroptvar 键与健康检查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赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX GCP Pub/Sub 生产者连接器健康检查失败诊断unhealthy_target 与状态码解析EMQX GCP Pub/Sub 生产者连接器健康检查失败诊断unhealthy_target 与状态码解析 导读本文围绕 EMQX 开源仓库中 apps/后端物联网消息队列通信LunaTranslator免费的游戏翻译工具五分钟跑通视觉小说实时翻译LunaTranslator免费的游戏翻译工具五分钟跑通视觉小说实时翻译 LunaTranslator 是一款完全开源免费的视觉小说翻译工具它的职责是把游后端物联网消息队列通信EMQX MongoDB Connector 健康检查修复解析find 权限不足不再误判连接断开EMQX MongoDB Connector 健康检查修复解析 find 权限不足不再误判连接断开 导读 本文围绕 EMQX 仓库变更记录 changes/e后端物联网消息队列通信上一篇Houston零配置Meteor管理工具让Django Admin体验在Meteor中重生下一篇floccus缓存机制CacheTree如何提升同步性能创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
