EMQX RabbitMQ 连接器多节点 servers 配置连接级故障转移与连接池旋转详解【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx导读本文围绕 EMQX 开源仓库中 RabbitMQ 桥接连接器的一项增强特性展开连接器配置支持多节点servers列表如rmq1:5672,rmq2:5672并在建立连接时按序尝试各节点实现故障转移同时通过连接池 worker 起始节点轮转避免多连接同时打向单一节点。读完本文你将掌握该特性的配置方式、与旧版server/port配置的兼容规则、底层源码实现原理以及对应的单元测试与集成测试验证方法。该功能对应的变更记录见 changes/ee/feat-17933.en.md并已进入 6.1.4 与 6.2.3 的发布记录changes/6.1.4.en.md、changes/6.2.3.en.md。一、功能概述在引入该特性之前EMQX 的 RabbitMQ 连接器只支持单一 RabbitMQ 节点地址serverport。当 RabbitMQ 以集群方式部署时单点配置意味着所有连接都依赖同一个节点且无法在连接建立时对集群内其他节点进行兜底。本次变更带来的核心能力有三点多节点servers列表连接器配置中可通过逗号分隔的servers字段一次性声明多个 RabbitMQ 节点例如rmq1:5672,rmq2:5672。连接级故障转移connect-time failover建立 AMQP 连接时按列表顺序逐个尝试只要某个节点可用即成功建连全部失败才报错。连接池启动偏移旋转rotated pool start offsets当pool_size 1时连接池中的每个 worker 从不同的列表位置开始尝试连接从而把初始连接请求分散到不同节点上。同时为保持向后兼容当servers未配置时旧版server/port配置依然生效。二、配置方式多节点 servers 与旧配置兼容1.servers字段RabbitMQ 连接器的 HOCON 配置 Schema 定义在 apps/emqx_bridge_rabbitmq/src/emqx_bridge_rabbitmq_connector_schema.erl 的fields(connector)中{servers, emqx_schema:servers_sc( #{ aliases [server], default localhost, desc ?DESC(servers) }, emqx_bridge_rabbitmq_client:host_options() )}, {port, ?HOCON( emqx_schema:port_number(), #{default 5672, desc ?DESC(port)} )},关键点servers通过emqx_schema:servers_sc/2生成 Schema声明了aliases [server]即旧的server字段是servers的别名Schema 层会把server归一化为规范的servers字段这一点在客户端模块注释%%serveris normalized to the canonicalserversfield by the schema中有明确说明见 emqx_bridge_rabbitmq_client.erl。默认值为localhostport默认5672。解析选项host_options()返回#{default_port DefaultPort, ssrf_check true}即列表中未显式书写端口的节点会使用默认端口并且对主机名解析启用 SSRF 检查。2. 旧版server/port兼容规则当servers未设置时用户依然可以按旧方式配置server rmq-legacy port 5671此时server作为别名会被归一化到servers配合port作为默认端口参与解析。测试用例 emqx_bridge_rabbitmq_client_tests.erl 明确验证了schema_test_中的两类输入都会被接受servers rmq1:5672,rmq2:5672多节点列表正常通过server rmq-legacy, port 5671旧格式正常通过并归一化为servers。3. 完整连接器配置参数结合 emqx_bridge_rabbitmq_connector_schema.erl连接器支持以下配置项字段类型默认值说明serversstring逗号分隔 host[:port] 列表localhost多节点列表可带别名serverport端口号5672未显式书写端口时的默认端口usernamebinary必填RabbitMQ 用户名passwordbinary必填RabbitMQ 密码使用密钥混淆保护pool_size正整数8连接池大小timeout时长5s连接超时virtual_hostbinary/RabbitMQ vhostheartbeat时长30sAMQP 心跳间隔sslobject#{enable false}TLS 配置复用emqx_connector_schema_lib:ssl_fields()Schema 中给出的完整示例值connector_example_values/0为#{ name rabbitmq_connector, type rabbitmq, enable true, servers 127.0.0.1:5672, username guest, password ******, pool_size 8, timeout 5s, virtual_host /, heartbeat 30s, ssl #{enable false} }4. 多节点配置示例HOCONservers rmq1:5672,rmq2:5672,rmq3:5673 username emqx password secret pool_size 8 timeout 5s virtual_host / heartbeat 30s ssl { enable false }通过 HTTP API 创建连接器时POST/api/v5/connectors可写成{ name: rabbitmq_connector, type: rabbitmq, enable: true, servers: rmq1:5672,rmq2:5672, username: guest, password: public, pool_size: 8, timeout: 5s, virtual_host: /, heartbeat: 30s, ssl: {enable: false} }提示连接器本身不直接收发数据需配合规则引擎 Action生产者或 Source消费者使用。Action/Source 的参数 Schema 见 emqx_bridge_rabbitmq_pubsub_schema.erl例如exchange、routing_key、delivery_mode、wait_for_publish_confirmations等。三、源码级实现原理1. servers 解析emqx_schema:servers_sc与parse_serversservers_sc/2定义在 apps/emqx/src/emqx_schema.erl由三阶段组成converterconvert_servers/1把 HOCON 值归一化为逗号分隔字符串。这一步处理了一个典型的 HOCON 陷阱——host.domain.name:80这类字符串在未加引号时可能被 HOCON 解析成嵌套 mapconvert_servers会将其还原为host:port对同时会去除逗号两侧的空格并把字符串数组旧格式[s1:80,s2:80]转换为逗号分隔形式。validator在配置加载时调用parse_servers/2验证每个host[:port]是否可解析保证非法地址在启动阶段即被拒绝。runtime parsing由各使用模块在运行时再次解析。RabbitMQ 客户端在运行时通过 emqx_bridge_rabbitmq_client.erl 的parse_servers/2完成解析parse_servers(BinServers, DefaultPort) - [ {emqx_utils_conv:str(Host), Port} || #{hostname : Host, port : Port} - emqx_schema:parse_servers(BinServers, host_options(DefaultPort)) ].emqx_schema:parse_servers/2emqx_schema.erl支持两种输入形态逗号分隔字符串或字符串数组兼容旧 Schema最终输出[{Host, Port}]元组列表。端口优先级规则可以从单元测试 emqx_bridge_rabbitmq_client_tests.erl 归纳为输入servers配置port解析结果rmq1:5672,rmq2:56731111[{rmq1,5672},{rmq2,5673}]内联端口优先rmq-legacy5671[{rmq-legacy,5671}]无端口时用默认端口rmq1,rmq2:56735671[{rmq1,5671},{rmq2,5673}]混合场景2. 连接级故障转移do_start_connection顺序尝试故障转移的核心逻辑在 emqx_bridge_rabbitmq_client.erldo_start_connection([], _AmqpParamsBase, Tried) - {error, #{reason all_nodes_failed, tried lists:reverse(Tried)}}; do_start_connection([{Host, Port} | Rest], AmqpParamsBase, Tried) - Params AmqpParamsBase#amqp_params_network{host Host, port Port}, case amqp_connection:start(Params) of {ok, Conn} - {ok, Conn}; {error, Reason} - ?SLOG(warning, #{ msg rabbitmq_connection_node_failed, host Host, port Port, reason Reason }), do_start_connection(Rest, AmqpParamsBase, [{Host, Port, Reason} | Tried]) end.工作机制可以概括为从列表头部开始逐个用amqp_connection:start/1尝试建连当前节点失败时记录rabbitmq_connection_node_failed级别的warning 日志包含 host、port 和原因然后继续尝试下一个节点只要有一个节点成功立即返回{ok, Conn}不再继续全部节点失败时返回{error, #{reason all_nodes_failed, tried [...]}}tried中记录每个被尝试节点及其失败原因便于排障。需要说明的是这里的故障转移发生在连接建立阶段connect-time。连接建立后的运行期断线由 AMQP 客户端自身的重连机制以及连接器的auto_reconnect2 秒间隔见 emqx_bridge_rabbitmq_connector.erl 的?AUTO_RECONNECT_INTERVAL_S共同保障。3. 连接池启动偏移旋转rotate_servers连接池场景下如果所有 worker 都从列表第一个节点开始尝试会导致初始连接请求全部集中到同一节点。rotate_servers/2emqx_bridge_rabbitmq_client.erl通过按 worker id 轮转列表起点解决这一问题rotate_servers(Servers, WorkerId) when is_integer(WorkerId), WorkerId 0 - Offset (WorkerId - 1) rem length(Servers), {Left, Right} lists:split(Offset, Servers), Right Left; rotate_servers(Servers, _WorkerId) - Servers.旋转算法为偏移量Offset (WorkerId - 1) rem length(Servers)将列表按偏移量拆成Left、Right两段后拼接为Right Left。单元测试 emqx_bridge_rabbitmq_client_tests.erl 给出了 3 节点列表[a, b, c]的旋转结果WorkerId旋转后起始顺序1[a, b, c]偏移 0不旋转2[b, c, a]3[c, a, b]5[b, c, a](5-1) rem 3 1与 WorkerId 2 同偏移这样在pool_size 8、servers含 3 个节点时各 worker 的首选节点均匀分布在这 3 个节点上当首选节点不可用时各 worker 的失败转移顺序也各不相同进一步分散了重试压力。4. 与 ecpool 连接池的整合连接器实现 emqx_bridge_rabbitmq_connector.erl 同时实现了emqx_resource与ecpool_worker两个 behaviouron_start/2调用emqx_resource_pool:start(InstanceId, ?MODULE, Options)创建连接池Options 中包含pool_size来自配置与auto_reconnect2 秒connect/1是 ecpool worker 的回调负责为每个 worker 建立 AMQP 连接connect(Options) - Config proplists:get_value(config, Options), WorkerId proplists:get_value(ecpool_worker_id, Options, 1), ... Servers0 emqx_bridge_rabbitmq_client:servers_from_config(Config), Servers emqx_bridge_rabbitmq_client:rotate_servers(Servers0, WorkerId), AmqpParamsBase #amqp_params_network{ ssl_options to_ssl_options(Config), username Username, password Password, connection_timeout Timeout, virtual_host VirtualHost, heartbeat Heartbeat }, case emqx_bridge_rabbitmq_client:start_connection(Servers, AmqpParamsBase) of {ok, RabbitMQConn} - {ok, RabbitMQConn}; {error, Reason} - ... % 记录 rabbitmq_connector_connection_failed 错误日志 end.从源码结构可以推断出完整的连接建立调用链emqx_resource_pool:start └─ ecpool 创建 pool_size 个 worker └─ connect/1每个 worker 调用一次 ├─ servers_from_config/1 → 解析 servers 列表 ├─ rotate_servers/2 → 按 WorkerId 旋转起始节点 └─ start_connection/2 → 顺序尝试节点实现 failover此外连接器通过init_secret/0初始化凭据混淆密钥确保密码不会明文出现在日志中emqx_bridge_rabbitmq_connector.erl。四、测试验证1. 单元测试eunit 测试文件 emqx_bridge_rabbitmq_client_tests.erl 覆盖了本特性的全部核心路径servers_from_config_test_验证内联端口优先、默认端口回退、混合写法三种解析规则rotate_servers_test验证 3 节点列表在 WorkerId 1/2/3/5 下的旋转结果start_connection_failover_test用meck模拟amqp_connection:start/1令{bad,1}返回{error, econnrefused}、{good,5672}返回成功断言最终返回{ok, rabbitmq_conn_stub}并通过meck:history验证尝试顺序为[{bad,1},{good,5672}]即严格按列表顺序故障转移start_connection_all_failed_test所有节点均失败时断言返回{error, #{reason : all_nodes_failed, tried : [{a,1,econnrefused},{b,2,econnrefused}]}}schema_test_验证新servers格式与旧server/port格式都能通过 Schema 校验。2. 集成测试emqx_bridge_rabbitmq_action_SUITE.erl 中的t_multi_node_connect_failover提供了端到端验证t_multi_node_connect_failover(TCConfig) - {201, _} create_connector_api(TCConfig, #{ servers rabbitmq:1,rabbitmq:5672 }), ...该用例故意把第一个节点配置为不可达端口rabbitmq:1第二个节点使用真实可用的rabbitmq:5672然后通过规则引擎发布一条multi-node-failover消息最终断言消息成功投递到 RabbitMQ。这直接证明了当列表头部节点不可达时连接器会自动尝试后续节点并成功完成数据流转。五、适用场景与注意事项推荐场景RabbitMQ 以集群多节点方式部署希望连接器建连时不依赖单一节点连接池较大默认pool_size 8希望各 worker 的初始连接与失败重试分散到不同节点避免惊群效应需要在 RabbitMQ 节点滚动升级、单节点短暂不可用期间保证 EMQX 数据集成链路仍能初始化。使用注意故障转移发生在连接建立时建连成功后如果该连接后续断开AMQP 客户端会尝试重连原节点如需在运行期切换节点依赖的是客户端重连与连接器的 2 秒auto_reconnect机制而非本特性的顺序尝试逻辑。列表顺序即优先级servers中的节点顺序决定了连接尝试的先后顺序建议把最稳定/最优先的节点放在前面。端口省略规则列表中省略端口的节点使用连接器的port配置默认 5672作为默认端口显式书写的端口优先。旧配置兼容server作为servers的别名依然可用升级现有配置无需强制迁移两者同时存在时以servers的归一化结果为准。SSRF 防护节点解析启用了ssrf_check主机名解析会经过 SSRF 安全检查见 emqx_bridge_rabbitmq_client.erl。六、版本记录本特性由 PR #17933 引入已收录于以下版本变更记录changes/6.1.4.en.mdchanges/6.2.3.en.md原始变更描述与本文核心一致RabbitMQ connector supports a multi-nodeserverslist (e.g.rmq1:5672,rmq2:5672) with connect-time failover and rotated pool start offsets. Legacyserver/portremain whenserversis unset.RabbitMQ 连接器支持多节点servers列表带连接时故障转移与旋转的连接池起始偏移当servers未设置时保留旧版server/port配置。如需深入源码可重点阅读以下文件配置 Schemaapps/emqx_bridge_rabbitmq/src/emqx_bridge_rabbitmq_connector_schema.erl客户端解析与建连apps/emqx_bridge_rabbitmq/src/emqx_bridge_rabbitmq_client.erl连接器资源实现apps/emqx_bridge_rabbitmq/src/emqx_bridge_rabbitmq_connector.erl单元测试apps/emqx_bridge_rabbitmq/test/emqx_bridge_rabbitmq_client_tests.erl集成测试apps/emqx_bridge_rabbitmq/test/emqx_bridge_rabbitmq_action_SUITE.erl通用 servers 解析基础设施apps/emqx/src/emqx_schema.erl【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
AI网关实战:从部署到核心功能,小团队接入大模型的最佳实践 1. 这个 37K Star 的项目到底解决了什么问题先说个我自己踩过的坑。去年我给团队搭内部 AI 服务,接了大模型 API,一开始觉得挺简单:不就是 HTTP 请求嘛,拿着 Key 调一下,返回结果就完事了。结果真正上线两周࿰… · 2026/9/23 3:33:27
OpenCode与Claude Code本质区别:API网关vs本地调度器 1. 这不是“选哪个更好”的测评,而是两个工具在真实开发流中的角色错位OpenCode 和 Claude Code 这两个名字最近频繁出现在开发者群、技术论坛和 VS Code 插件市场评论区里,但很多人点开安装、配置、跑起来之后才发现——它们根本不是同一类东西。我过去… · 2026/9/23 3:33:21
LLM工具调用速记:生产级Function Calling与MCP工程实践 1. 什么是“LLM工具调用速记”:不是语法口诀,而是工程现场的肌肉记忆你打开一个Agent项目,刚写完一段prompt,准备让模型调用数据库查询接口——结果模型返回了一段看似合理但根本无法执行的JSON,字段名拼错、required参… · 2026/9/23 3:33:21
3分钟搞懂系统截图快捷键原理,面试不再挂科 3分钟搞懂系统截图快捷键原理,面试不再挂科 面试时被问“系统截图快捷键底层是怎么实现的”,你脑子里是不是只剩“Ctrl+Shift+S”?别慌,这题卡住很多人。今天这篇文章带你一文搞懂,从用户按下按键到图片存盘,全链路拆解,让你下次回答能直… · 2026/9/23 4:15:40
科技内容为何在短视频平台爆发?从1.4万亿次观看看全民科技热潮 你有没有在深夜刷抖音时,点进一个标题叫“为什么AI画手多了一根手指”的视频,结果一路刷完了评论区几百条吵架式讨论?这不是你的错觉,而是科技内容正式从小众爱好走向大众茶余饭桌的标志。2025年,抖音上科技类内容的观… · 2026/9/23 4:15:40
Biome Markdown 格式化器有序列表编号重排机制深度解析 开发工具Lint格式化静态分析代码质量前端 【免费下载链接】biome A toolchain for web projects, aimed to provide functionalities to maintain them. Biome offers formatter and linter, usable via CLI and LSP. 项目地址: https://gitcode.com/gh_mirrors/bi/… · 2026/9/23 4:15:40
存在主义视角下的焦虑本质与转化方法 1. 焦虑的本质与哲学解读焦虑(Angst)这个词在德语中有着特殊的哲学含义,它不同于普通的恐惧或担忧。我第一次深入理解这个概念是在研读存在主义哲学著作时,那种醍醐灌顶的感觉至今难忘。焦虑不是简单的负面情绪,而是人… · 2026/9/23 4:15:40
n76备考保姆级教程:告别配置地狱,5天搞定证书 n76备考保姆级教程:告别配置地狱,5天搞定证书 配置环境就卡半天,代码跑不通,报错日志看得人眼瞎。这种痛苦每个想考n76的朋友都经历过。 别慌,这篇保姆级教程带你避开90%的坑。 坑的现象:为什么你总是卡在环境配置上 现象描述:… · 2026/9/23 4:15:40
纯前端3D时空渲染引擎:浏览器内实现60帧高性能可视化 1. 从标题拆解这个项目的真实面貌1.1 这个标题到底在说什么“在浏览器里开间谍卫星”这个说法听起来很唬人,但拆开来看,它描述的其实是一类非常具体的技术产品形态:一个完全跑在浏览器里的三维时空数据可视化引擎。所谓“间谍卫星”是一种比喻… · 2026/9/23 4:15:34
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29