ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

EMQX RabbitMQ 连接器多节点 servers 配置:连接级故障转移与连接池旋转详解

EMQX RabbitMQ 连接器多节点 servers 配置:连接级故障转移与连接池旋转详解 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),仅供参考
返回列表