EMQX Redis Bridge 实战指南用规则引擎把 IoT 数据写入 Redis【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqxEMQX 的emqx_bridge_redis应用负责把 EMQX 与 Redis 连接起来用户只需在规则引擎中编写一条规则即可将 IoT 消息实时写入 Redis 的字符串、List、Hash 等数据结构用于缓存、去重、设备状态存储等场景。本文围绕该应用系统讲解其工作原理、三种 Redis 部署形态单机 / Sentinel / Cluster的连接器配置、动作命令模板语法、批量写入与错误处理机制并结合仓库源码给出可复现的配置与验证方法。一、应用定位EMQX 与 Redis 之间的数据通道Redis 是一个内存型数据结构存储常被用作分布式键值数据库、缓存与消息中间件并支持可选的数据持久化。emqx_bridge_redis 就是专门用于打通 EMQX 与 Redis 的桥接应用用户创建一条规则通过 EMQX 规则引擎 把消息转发到 Redis 动作Action动作再通过 Redis 连接器Connector把命令下发到 Redis 实例。从仓库结构看这条链路分三层规则层规则引擎的 SQL 输出消息触发 Redis 动作。动作层emqx_bridge_redis负责把消息按command_template模板渲染成具体的 Redis 命令如LPUSH、HMSET。连接器层emqx_redis应用负责建立并维护到 Redis 的连接池、执行命令、健康检查同时被 Redis 桥接和emqx_auth用于权限校验两个场景复用。连接器的真正实现在独立的 emqx_redis 应用中其模块注释明确写道该应用“为 Redis 桥接插入消息、为emqx_auth检查用户权限提供 API”可见 Redis 连接器是桥接与认证功能共享的底层基础组件。二、支持的 Redis 部署形态与连接器配置连接器配置由 emqx_bridge_redis_schema.erl 定义其parameters字段是一个联合类型支持三种形态对应emqx_redis.erl中redis_single_connector、redis_sentinel_connector、redis_cluster_connector三套字段定义2.1 单机模式singleconnectors.redis.my_connector { enable true description My redis single connector parameters { redis_type single server 127.0.0.1:6379 pool_size 8 database 1 username test password ****** } ssl { enable false } }单机模式使用server字段指定单个 Redis 地址emqx_schema:servers_sc会自动补齐默认端口见emqx_redis.erl中的?REDIS_HOST_OPTIONS默认端口为 Redis 标准端口 6379。2.2 哨兵模式sentinelconnectors.redis.my_sentinel_connector { enable true description My redis sentinel connector parameters { redis_type sentinel servers 127.0.0.1:6379,127.0.0.2:6379 sentinel myredismaster pool_size 8 database 1 sentinel_username sentinel_user sentinel_password ****** username test password ****** } ssl { enable false } }哨兵模式的关键点servers是逗号分隔的哨兵节点列表EMQX 会依次连接并询问主节点地址。sentinel是哨兵监控的 master 名称如myredismaster必须填写且required true。哨兵自身的认证与 Redis 主从的认证是两套独立凭证sentinel_username/sentinel_password用于连接哨兵username/password用于连接实际的 master。从源码看连接启动时会先做哨兵连通性校验emqx_redis.erl中的validate_sentinel_connection/4逐个连接哨兵节点执行SENTINEL get-master-addr-by-name master名只有拿到 master 地址才算校验通过全部失败则返回{error, {sentinel_error, ...}}其中IDONTKNOW被映射为sentinel_master_unreachable。2.3 集群模式clusterconnectors.redis.my_cluster_connector { enable true description My redis cluster connector parameters { redis_type cluster servers 127.0.0.1:6379,127.0.0.2:6379 pool_size 8 username test password ****** } ssl { enable false } }集群模式同样使用逗号分隔的servers列表至少给出一个可路由的种子节点即可由eredis_cluster库负责槽位slot路由。需要注意集群模式不支持database字段——emqx_redis.erl中redis_cluster_connector的字段定义显式执行了lists:keydelete(database, 1, redis_fields())且database/1函数对cluster类型返回空列表。集群模式与批量动作不兼容emqx_redis.erl的do_cmd/3中对集群批量命令留有% TODO Cluster mode is currently incompatible with batching注释且 emqx_bridge_redis_action_info.erl 在从集群连接器转换动作配置时会强制把batch_size改为 1、batch_time改为0ms。2.4 连接器通用参数无论哪种形态连接器都支持以下公共参数定义在emqx_redis.erl的redis_fields/0参数类型默认值说明redis_typeenumsinglesingle/sentinel/cluster三选一server/serversstring无单机用server哨兵/集群用逗号分隔的serverspool_sizeinteger无测试中常见 1~8连接池大小来自emqx_connector_schema_lib:pool_size/1usernamestring可选Redis 6 ACL 用户名passwordstring可选认证密码databaseinteger0逻辑数据库编号仅 single/sentinel 可用auto_reconnectboolean见emqx_connector_schema_lib是否自动重连sslobject{enable false}TLS 配置enabletrue时支持verify、cacertfile、certfile、keyfile等SSL 开启后底层会通过emqx_tls_lib:to_client_opts/1生成客户端 TLS 选项测试套件中常见配置为verify verify_none搭配客户端证书。三、动作配置command_template 命令模板动作Action是规则引擎与 Redis 之间的最终执行单元其核心配置是command_template——一条由字符串数组组成的 Redis 命令模板数组第一项是命令名后续项是参数。以下示例来自 emqx_bridge_redis_schema.erl 的官方示例bridges.redis.my_action { enable true connector my_connector_name description My action parameters { command_template [LPUSH, MSGS, ${payload}] } resource_opts { batch_size 1 } }3.1 模板语法与占位符command_template的每个元素都会被 emqx_bridge_redis_connector.erl 用emqx_placeholder进行预处理和运行时渲染支持${payload}、${topic}、${clientid}、${username}等规则引擎内置变量以及 SQLSELECT出来的任意字段例如${payload.temperature}。支持字符串与变量的拼接例如MSGS/${topic}会被渲染为带前缀的 Redis key。变量占位符也可用于命令名本身即模板第一项实现动态命令。校验规则见emqx_bridge_redis.erl的is_command_template_valid/1必须是非空字符串列表否则报错 “the value of the field command_template should be a nonempty list of strings (templates for Redis command and arguments)”。以测试套件 emqx_bridge_redis_SUITE.erl 中的典型模板为例parameters.command_template [RPUSH, MSGS/${topic}, ${payload}]配合规则SELECT * FROM t/#设备向t/device1发布消息后Redis 中会生成 keyMSGS/t/device1并用RPUSH把消息 payload 追加到 List 尾部——这正是把 Redis 当作 MQTT 消息队列/日志缓冲的常见用法。3.2 Hash 写入的高级用法map_to_redis_hset_args对于需要把消息转成 Redis Hash 的场景规则 SQL 中可以使用map_to_redis_hset_args(payload)函数。emqx_bridge_redis_connector.erl的proc_tmpl/2对该函数结果做了特殊处理当占位符渲染结果以map_to_redis_hset_args打头时会将其余部分展开为多个命令参数而不是一个整体字符串从而把形如{a: 1, b: 2}的 JSON payload 直接展开成HMSET key a 1 b 2的参数序列。测试用例t_map_to_redis_hset_args验证了该链路SELECT map_to_redis_hset_args(payload) as payload FROM t/#parameters.command_template [HMSET, t_map_to_redis_hset_args, ${payload}]设备发布{a: 1, b: 2}后Redis 上会得到包含a、b两个字段的 Hash。3.3 动作资源选项resource_opts动作层支持独立的资源选项默认值定义在 schema 中参数默认值说明batch_size100批量写入时一批最多包含的消息条数batch_time100ms批量写入的最大攒批等待时间query_modesyncsync同步或async异步来自测试矩阵max_retries、retry_interval继承通用动作选项失败重试策略health_check_interval、health_check_timeout继承通用选项健康检查周期与超时在同步模式下动作会调用连接器的on_query/3单条或on_batch_query/3批量在query_mode async下EMQX 会先把消息放入队列再按batch_size/batch_time攒批执行。测试套件对batch_size 5, batch_time 100ms与batch_size 1, batch_time 0ms即不攒批两种配置都做了矩阵覆盖。四、与规则引擎配合的完整数据链路把连接器、动作与规则串起来端到端流程如下创建 Redis 连接器上文三种形态任选其一EMQX 建立连接池并通过PING健康检查。创建 Redis 动作绑定连接器配置command_template。创建规则SELECT payload, topic, clientid FROM t/#在规则的“动作”里挂载刚创建的 Redis 动作。设备发布消息到t/#规则引擎触发动作连接器把渲染后的 Redis 命令发送到 Redis。该链路在t_rule_action测试中完整走通设备用emqtt客户端发布消息测试先执行LRANGE清空目标 key再在发布后轮询断言LRANGE能读回完整 payload同时验证了 TCP/TLS × 单机/哨兵/集群 × 同步/异步 × 是否攒批的 24 种组合emqx_bridge_redis_SUITE.erl。五、命令执行、批处理与错误分类5.1 底层命令执行连接器层把渲染后的命令封装为{cmd, Command}单条或{cmds, Commands}批量交给emqx_redis:on_query/3单机/哨兵通过ecpool:pick_and_do从连接池取连接调用eredis:q单条或eredis:qp批量 pipeline。集群直接调用eredis_cluster:q按槽位路由批量命令暂为逐个执行后聚合结果。批量结果wrap_qp_result只有在所有子命令都成功时才返回{ok, Results}只要有一条失败就返回{error, Results}。5.2 可恢复错误与不可恢复错误emqx_redis.erl的do_query会对错误做分类不可恢复错误unrecoverableERR unknown command ...命令不存在或invalid_cluster_command集群下不支持的命令。这类错误重试没有意义动作会直接失败。可恢复错误recoverableno_connection、ecpool_empty等会触发重试/缓存机制。测试用例t_permanent_error专门构造了[BAD, COMMAND, ${payload}]模板断言发送结果恒为{error, _}验证了不可恢复错误的处理路径。5.3 断连与重放保障t_check_replay测试演示了断连期间消息不丢失的行为通过 toxiproxy 人为切断 Redis 连接连续发布 5 条消息期间批量发送返回{error, _}且触发重试恢复连接后消息成功写入日志事件redis_bridge_connector_send_done从result : {error, _}翻转为result : {ok, _}。这说明在同步 攒批模式下EMQX 会对发送失败的消息执行重试直到 Redis 恢复可用。5.4 健康检查与状态单机/哨兵对连接池每个 worker 执行PINGdo_get_status全部{ok, _}才判定connected。集群检查eredis_cluster_monitor是否还有 worker没有则直接判定disconnected否则对全部 worker 执行PING。特殊处理若PING返回NOAUTH Authentication required.说明凭证缺失或错误会被标记为{disconnected, {unhealthy_target, ...}}并产生emqx_redis_auth_required_error事件。测试t_auth_error_username_password、t_auth_error_password_only验证了错误密码会导致创建连接器后状态为disconnected且原因以{unhealthy_target,开头。六、通过 HTTP API 管理 Redis 桥接emqx_bridge_redis遵循 EMQX 统一的桥接 v2 管理 API实现位于emqx_bridge_v2_api本应用的 schema 在 emqx_bridge_redis_schema.erl 中定义了对应的请求/响应字段支持创建连接器POST /api/v5/connectors类型redis查询/更新/删除连接器GET/PUT/DELETEemqx_bridge_redis_connector_info.erl通过api_schema(Method)把 HTTP 方法映射到post_connector、put_connector、get_connector等字段集合创建动作POST /api/v5/bridges类型redis携带connector引用启停与状态查询响应中带status如connected/disconnected与node_status各节点状态例如emqxlocalhost - connectedGET 接口的返回示例来自 schema 中的action_example(get){ type: redis, name: my_action, status: connected, node_status: [ { node: emqxlocalhost, status: connected } ] }创建连接器时若 Redis 不可达接口仍会返回201但status为disconnected见t_create_disconnected测试随后 EMQX 会按auto_reconnect策略持续尝试重连。七、源码结构与进一步探索如果你想深入源码可以按以下路径阅读关注点文件连接器/动作的 HOCON schema 与示例apps/emqx_bridge_redis/src/emqx_bridge_redis_schema.erl动作参数与模板校验apps/emqx_bridge_redis/src/emqx_bridge_redis.erl模板渲染、批量发送、通道管理apps/emqx_bridge_redis/src/emqx_bridge_redis_connector.erl连接池、命令执行、健康检查apps/emqx_redis/src/emqx_redis.erlredis-cli 风格命令行解析split/1复刻 Redis 的sdssplitargsapps/emqx_redis/src/emqx_redis_command.erl全矩阵集成测试TCP/TLS × 三种形态 × 同步/异步 × 攒批apps/emqx_bridge_redis/test/emqx_bridge_redis_SUITE.erl其中emqx_redis_command:split/1值得一提它完整复刻了 Redis 自身redis-cli的sdssplitargs解析逻辑支持双引号/单引号包裹、\xHH十六进制转义和\n \r \t \b \a特殊字符为命令模板的解析提供了与 Redis 客户端一致的语义。八、小结emqx_bridge_redis让“规则引擎 → Redis”的接入变得高度模板化三种部署形态共用一套连接器框架command_template支持任意 Redis 命令与占位符组合map_to_redis_hset_args让 JSON 消息一键写入 Hash批量/重试机制保证了 Redis 短暂不可用时的消息不丢失。结合 emqx_bridge_redis_SUITE.erl 的矩阵测试与 emqx_redis.erl 的实现细节你可以在此基础上按需扩展出缓存刷新、设备状态快照、消息归档等各类 Redis 集成场景。如需参与该应用的开发与贡献请参考仓库根目录的 CONTRIBUTING.md。【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考