后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载本篇技术指南基于 EMQX 开源仓库中的变更记录 fix-16936.en.md深入剖析 Azure Blob Storage Action 在聚合aggregated模式下健康检查超时问题的成因、修复方案与底层实现。读者将了解到 EMQX 桥接组件连接器健康检查 通道Channel健康检查两级机制的真实调用链掌握聚合模式健康检查触发的完整子流程以及erlazure:list_blobs列表请求限制max_results在其中发挥的关键作用。问题背景聚合模式下健康检查为什么会超时变更记录 fix-16936.en.md 描述的问题非常聚焦Fixed an issue where the health check of an Azure Blob Storage Action in aggregate mode could timeout if the container contained too many blobs. 修复了一个问题当容器中包含过多 blob 时聚合模式下的 Azure Blob Storage Action 健康检查可能超时。要理解这个问题必须先理解 EMQX 中 Azure Blob Storage 动作的通道健康检查是如何实现的。在 emqx_bridge_azure_blob_storage_connector.erl 中聚合模式通道的状态检查入口是channel_status/2channel_status(#{mode : aggregated} ActionState, ConnState) - #{driver_state : DriverState} ConnState, #{container : Container, aggreg_id : AggregId} ActionState, %% NOTE: This will effectively trigger uploads of buffers yet to be uploaded. Timestamp erlang:system_time(second), ok emqx_connector_aggregator:tick(AggregId, Timestamp), ok check_schema_reference_valid(ActionState), ok check_container_accessible(DriverState, Container), ok check_aggreg_upload_errors(AggregId), ?status_connected.从源码看一次聚合模式通道的健康检查会依次执行四个步骤tick/2向聚合器发送一次时钟信号实际上会触发尚未上传的缓冲数据执行上传源码注释明确指出这一点check_schema_reference_valid/1校验聚合容器选项如 Parquet 的 Schema Registry 引用是否仍然有效check_container_accessible/2验证目标容器是否可访问——这正是本次修复的核心位置check_aggreg_upload_errors/1检查聚合上传是否累积了错误若有则抛出{unhealthy_target, ErrorMessage}使通道被标记为不健康。其中第三步check_container_accessible/2的实现为check_container_accessible(DriverState, Container) - do_list_blobs(DriverState, Container).它直接委托给do_list_blobs/2。在修复之前这一步需要对容器内的 blob 进行列表枚举当容器内 blob 数量非常多时一次完整列举需要与 Azure 服务端进行多轮分页交互健康检查耗时随之线性增长最终超出健康检查的请求时限health_check_interval/request_ttl表现为健康检查超时。修复方案将 blob 列表请求限制为单条结果修复后的do_list_blobs/2位于 emqx_bridge_azure_blob_storage_connector.erldo_list_blobs(DriverState, Container) - try erlazure:list_blobs(DriverState, Container, [{max_results, 1}]) of {L, _} when is_list(L) - ok catch _:_ - error end.这次修复包含两个关键点{max_results, 1}限制列表规模向erlazure:list_blobs/3传入max_results 1健康检查只需确认容器存在且可被列举而不再关心容器内有多少 blob。无论容器中包含 1 个还是数百万个 blob健康检查的耗时都保持恒定单次请求从根本上消除了blob 数量越多、健康检查越慢的放大效应。这与连接器级健康检查health_check/1的做法保持一致——后者同样使用erlazure:list_containers(DriverState, [{max_results, 1}])只列举 1 个容器即可判定连通性health_check(DriverState) - case erlazure:list_containers(DriverState, [{max_results, 1}]) of {error, Reason} - {?status_disconnected, Reason}; {L, _} when is_list(L) - ?status_connected end.异常兜底将erlazure:list_blobs/3的调用包裹在try ... catch中。一旦列举请求抛异常例如凭据失效、容器被删除、网络中断直接返回error原子由上层channel_status/2链式求值使健康检查失败通道被标记为?status_disconnected避免异常向上抛出污染资源管理器的状态机。值得说明的是聚合模式上传本身并不依赖本次列举结果——上传路径是通过process_append/2、process_write/1与process_complete/1逐块block写入并最终put_block_list提交的。因此健康检查中列举 blob的唯一目的就是以最小代价验证容器可达性限制为单条结果不会影响功能正确性这正是该修复能够成立的前提。深入EMQX 中 Azure Blob Storage Action 的两级健康检查机制本次修复所涉及的健康检查实际上是两层结构理解它们的区别有助于排查同类问题。连接器级健康检查Connector Health Check连接器Connector负责管理与 Azure 存储账号的底层会话。其健康检查回调on_get_status/2直接调用health_check/1见上节代码通过erlazure:list_containers/2列举存储账号下的容器来验证连接器配置中resource_opts.health_check_interval决定检查频率示例默认值为45s见 emqx_bridge_azure_blob_storage_connector_schema.erl返回?status_connected或{?status_disconnected, Reason}。通道级健康检查Channel Health Check通道Channel即具体 Action。聚合模式下执行上文channel_status/2的四个步骤而直接direct模式则简单得多——它不检查容器因为连接器健康检查已经验证了客户端可用于列举容器源码注释原话channel_status(#{mode : direct}, _ConnState) - %% Theres nothing in particular to check for in this mode; the connector health check %% already verifies that were able to use the client to list containers. ?status_connected;因此本次超时问题仅存在于聚合模式direct 模式通道健康检查不触碰容器列举聚合模式则因check_container_accessible/2全量列举 blob 而可能超时。健康检查的副作用设计聚合模式的通道健康检查并非只读探活——tick/2会触发待上传缓冲数据的实际上传check_aggreg_upload_errors/1会把最近一次上传失败反映为通道不健康。源码注释也提示多次上传失败会导致通道在连续的多次健康检查中被标记为不健康3 upload failures will cause the channel to be marked as unhealthy for 3 consecutive health checks。这种设计让健康检查同时充当了故障传播的载体配置较长的health_check_interval时需要考虑故障反馈的延迟。聚合模式 Action 的配置参考要复现或验证上述健康检查行为需要配置一个聚合模式的 Azure Blob Storage Action。以下字段来自 emqx_bridge_azure_blob_storage_action_schema.erl 的aggreg_parameters与aggregation定义配置项类型默认值说明parameters.modeaggregated必填聚合上传模式parameters.aggregation.container.typecsv/json_lines/parquet必填聚合容器格式决定上传文件的 content-typetext/csv/application/jsonl/ 其它parameters.aggregation.time_intervalduration秒1h聚合时间窗口parameters.aggregation.max_records正整数1_000_000单个聚合文件最大记录数parameters.containerstring必填Azure Blob Storage 目标容器名注意与聚合container是两个不同概念parameters.blob模板字符串必填blob 名称模板如${action}/${node}/${datetime.rfc3339}/${sequence}parameters.max_block_sizebytesize4000mb单块最大字节数必须 ≤ 4000 MiB由block_size_validator/1强制校验parameters.min_block_sizebytesize10mb触发块写入的最小缓冲字节数resource_opts.batch_size整数100批量投递条数注释建议聚合动作使用宽松批处理默认值resource_opts.batch_timeduration10ms批量投递时间窗口resource_opts.health_check_intervalduration30s健康检查周期Schema 中给出的聚合模式完整配置示例action_example(put, aggregated)见 emqx_bridge_azure_blob_storage_action_schema.erl如下{ enable true connector my_connector parameters { mode aggregated aggregation { container { type csv column_order [a, b] } time_interval 4s max_records 10000 } container mycontainer blob ${action}/${node}/${datetime.rfc3339}/${sequence} } resource_opts { batch_time 10ms batch_size 100 health_check_interval 30s inflight_window 100 query_mode sync request_ttl 45s worker_pool_size 16 } }连接器侧的配置emqx_bridge_azure_blob_storage_connector_schema.erl包含account_name必填Azure 存储账号名account_key必填secret 类型存储账号密钥必须是合法的 Base64 编码值validate_account_key/1会在启动时校验错误提示为bad account key; must be a valid base64 encoded valueendpoint可选自定义终结点默认由erlazure:new/1依据账号名推导resource_opts通用资源选项其中health_check_interval示例值为45s。在聚合模式健康检查的四个步骤中check_schema_reference_valid/1会再次校验聚合容器选项调用emqx_connector_aggreg_delivery:validate_container_opts/1因此当container.type parquet且引用了 Schema Registry 时健康检查还会连带验证 schema 引用的有效性。如何验证该修复仓库中的测试套件 emqx_bridge_azure_blob_storage_SUITE.erl 覆盖了健康检查相关的行为测试配置中大量使用health_check_interval 1s来缩短健康检查周期以便在用例执行窗口内触发多次检查list_blobs/2辅助函数通过erlazure:list_blobs(Client, str(ContainerName))直接读取容器内容用于断言上传结果t_parquet_bad_reference_health_check/0见 emqx_bridge_azure_blob_storage_SUITE.erl专门验证通道健康检查期间校验 Schema Registry 引用这一行为与check_schema_reference_valid/1一一对应。运行该套件需要本地可用的 Azure Blob Storage 兼容服务仓库 CI 使用 Azurite 模拟器见.ci/docker-compose-file/docker-compose-azurite.yaml并配置account_name/account_key环境。对于希望自行验证大容器健康检查不超时的读者可以在 Azurite 或真实存储账号中预先写入大量 blob再以1s的health_check_interval观察通道状态是否持续保持connected——修复后由于列表请求固定为单条结果健康检查耗时与 blob 数量无关。小结本次修复fix-16936针对聚合模式 Azure Blob Storage Action 的健康检查超时问题核心改动体现在 emqx_bridge_azure_blob_storage_connector.erl 的do_list_blobs/2将erlazure:list_blobs/3的max_results固定为1使健康检查从全量列举 blob降级为单条可达性探测并用try/catch兜底异常。这一改动的设计原则同样体现在连接器级健康检查list_containers限量为 1上健康检查只应回答目标是否可达而不是目标里有什么将探测开销与数据规模解耦是桥接组件健康检查的通用最佳实践。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX Oracle Action 健康检查机制详解复杂 SQL 模板下健康检查失败的修复与源码剖析EMQX Oracle Action 健康检查机制详解复杂 SQL 模板下健康检查失败的修复与源码剖析 本文聚焦 EMQX 企业版EE变更记录 fix 1后端物联网消息队列通信EMQX 新增 resource_opts.health_check_timeoutConnector/Action/Source 健康检查超时机制详解与升级影响EMQX 新增 resource_opts.health_check_timeout Connector/Action/Source 健康检查超时机制详解与升后端物联网消息队列通信EMQX Kafka / Pulsar Connector 健康检查超时状态机修复从 disconnected 到 connecting 的演进EMQX Kafka / Pulsar Connector 健康检查超时状态机修复从 disconnected 到 connecting 的演进 导读 本文以后端物联网消息队列通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考