EMQX MatrixDB 桥接将 IoT 数据实时写入 YMatrix 超融合数据库【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx本指南介绍 EMQX 数据集成体系中的 MatrixDBYMatrix桥接能力通过emqx_bridge_matrix应用将 EMQX 与 YMatrix 超融合数据库打通借助 EMQX 规则引擎把 MQTT 消息按 SQL 模板批量写入 MatrixDB并可通过 HTTP API 对桥接进行全生命周期管理。读完本文你将掌握 MatrixDB 桥接的连接配置、动作ActionSQL 模板编写、批量参数调优以及桥接的创建、更新、查询、启停等管理操作。MatrixDBYMatrix与 EMQX 数据集成YMatrix 是 YMatrix 公司基于 PostgreSQL / Greenplum 经典开源数据库开发的一款超融合数据库产品。除了能够轻松应对时序Time Series场景之外它还支持联机事务处理OLTP与联机分析处理OLAP等经典场景适合作为 IoT 数据汇聚、分析与存储的统一底座。emqx_bridge_matrix应用正是用于连接 EMQX 与 MatrixDB 的桥接组件。用户只需在 EMQX 中创建一个规则Rule即可借助 EMQX Rules 规则引擎将 IoT 数据便捷地摄取Ingest进 MatrixDB设备通过 MQTT 上报消息到 EMQX规则引擎对消息进行匹配与处理匹配命中的消息经 MatrixDB 桥接按 SQL 模板写入数据库。该应用位于仓库的 apps/emqx_bridge_matrix 目录当前应用版本为 6.3.0见 mix.exs。应用结构与模块职责emqx_bridge_matrix是一个轻量级桥接应用源码仅包含三个 Erlang 模块职责非常清晰模块职责emqx_bridge_matrix.erlHOCON Schema 定义桥接/动作Action与连接器Connector的配置字段、命名空间与 API 示例emqx_bridge_matrix_action_info.erl向 EMQX 注册matrix动作类型action_type_name() - matrixemqx_bridge_matrix_connector_info.erl向 EMQX 注册matrix连接器类型及其配置 Schema在 mix.exs 中应用通过emqx_action_info_modules与emqx_connector_info_modules两个环境变量把上述两个*_info模块注册进 EMQX 框架从而让 EMQX 在数据集成界面与 API 中识别matrix这一桥接类型。与 PostgreSQL 桥接的深层关系从源码结构可以清楚看到MatrixDB 桥接与 PostgreSQL 桥接同出一脉。在 emqx_bridge_matrix.erl 中动作字段直接复用emqx_bridge_pgsql:fields(post, ?ACTION_TYPE, config)其中?ACTION_TYPE被定义为matrix连接器字段直接复用emqx_postgresql_connector_schema:fields({Field, ?CONNECTOR_TYPE})其中?CONNECTOR_TYPE同样被定义为matrix动作示例bridge_v2_examples/1与连接器示例connector_examples/1也分别委托给emqx_bridge_pgsql:values/1与emqx_postgresql_connector_schema:values/1生成。这一设计非常自然YMatrix 兼容 PostgreSQL 协议因此 MatrixDB 桥接直接继承了 PostgreSQL 桥接成熟的连接池、SSL 与 SQL 执行能力仅需以matrix类型注册即可。而在 emqx_bridge_matrix_connector_info.erl 中resource_callback_module() - emqx_postgresql表明MatrixDB 连接器的底层资源回调由 apps/emqx_postgresql 应用基于 epgsql 驱动提供即实际建连与 SQL 下发都由 PostgreSQL 驱动完成这进一步印证了两者协议层面的兼容性。连接器Connector配置详解连接器负责维护到 MatrixDB 的连接池。其字段定义于 emqx_postgresql_connector_schema.erl 的fields(connection_fields)并由fields(config_connector)追加公共字段与资源选项resource_opts后构成完整配置字段类型说明serverstringMatrixDB 地址格式host:port默认端口 5432YMatrix 与 PostgreSQL 一致application_namebinary应用名默认emqx用于在数据库中标识连接来源经过emqx_postgresql:validate_application_name/1校验disable_prepared_statementsboolean是否禁用预编译语句databasebinary目标数据库名必填且不能为空pool_sizepos_integer连接池大小驱动层连接池未来将让位于资源层的worker_pool_size见 emqx_connector_schema_lib.erl 中相关注释username/passwordstring数据库账号与密码密码字段经emqx_schema_secret处理auto_reconnectboolean断线自动重连开关sslobjectSSL/TLS 配置enable、verify、versions、ciphers、depth、hibernate_after、log_level、reuse_sessions、secure_renegotiate等resource_optsobject资源健康检查与生命周期选项health_check_interval、start_after_created等以上关系型数据库公共字段database、pool_size、username、password、auto_reconnect定义于 emqx_connector_schema_lib.erl 的relational_db_fields/1。一个典型的连接器配置如下来自连接器示例值见 emqx_postgresql_connector_schema.erl{ name: my_matrix_connector, type: matrix, enable: true, server: 127.0.0.1:5432, database: emqx_data, username: postgres, password: public, pool_size: 8, application_name: emqx, ssl: { enable: false, verify: verify_peer, versions: [tlsv1.3, tlsv1.2] } }动作Action与 SQL 模板动作定义了“把命中的消息如何写进 MatrixDB”。其核心是 SQL 模板字段parameters.sql该字段是一个模板字符串emqx_schema:template()格式标记为sql支持 EMQX 规则引擎的${var}占位符。默认 SQL 模板定义于 emqx_bridge_pgsql.erl 的default_sql/0insert into t_mqtt_msg(msgid, topic, qos, payload, arrived) values (${id}, ${topic}, ${qos}, ${payload}, TO_TIMESTAMP((${timestamp} :: bigint)/1000))对应的表可参照如下结构创建CREATE TABLE t_mqtt_msg ( id VARCHAR(64), topic VARCHAR(255), qos INTEGER, payload TEXT, arrived TIMESTAMP );其中${id}、${topic}、${qos}、${payload}、${timestamp}来自 EMQX 消息对象字段id为消息 IDtopic为主题qos为服务质量等级payload为消息负载timestamp为消息时间戳毫秒。TO_TIMESTAMP((${timestamp} :: bigint)/1000)将毫秒时间戳转换为数据库TIMESTAMP这是 IoT 时序数据写入数据库的常见手法。若消息负载是 JSON 且已通过规则解码为对象也可以在 SQL 模板中直接引用其内部字段例如INSERT INTO client_events(clientid, event, created_at) VALUES ( ${clientid}, ${event}, TO_TIMESTAMP((${timestamp} :: bigint)) )批量写入参数动作还带有独立的资源选项resource_opts用于控制批量写入行为定义于emqx_bridge_pgsql.erl的fields(action_resource_opts)参数默认值说明batch_size100每批写入的行数上限batch_time100ms攒批的最长等待时间时间窗内未凑满batch_size也会触发写入inflight_window—允许同时在途的最大请求数示例值为100max_buffer_bytes—缓冲数据量上限示例值为256MBrequest_ttl—请求最大存活时间示例值为45sworker_pool_size—执行写入的 worker 池大小示例值为16batch_size与batch_time的组合是吞吐与延迟的平衡点吞吐优先可调大batch_size延迟敏感可调小batch_time。一个完整的动作配置示例来自 emqx_bridge_pgsql.erl 的values/1{ name: my_action, type: matrix, enable: true, connector: my_matrix_connector, resource_opts: { batch_size: 1, batch_time: 50ms, inflight_window: 100, max_buffer_bytes: 256MB, request_ttl: 45s, worker_pool_size: 16 }, parameters: { sql: INSERT INTO client_events(clientid, event, created_at) VALUES (${clientid}, ${event}, TO_TIMESTAMP((${timestamp} :: bigint))) } }借助规则引擎将 IoT 数据摄入 MatrixDB按照 README 的说明接入流程可以概括为创建 MatrixDB 连接器配置server、database、username、password等连接参数确认连接器状态为connected创建 MatrixDB 动作为连接器配置 SQL 模板与批量参数创建规则在 EMQX 规则引擎中编写 SQL 语句匹配消息例如SELECT * FROM t/#并将动作挂接到规则的“动作”输出中验证向对应主题发布测试消息在 MatrixDB 中查询表确认数据落库。桥接的可用性还取决于 MatrixDB 侧的准备工作确保目标数据库与账号存在、目标表已按 SQL 模板创建、账号对表拥有 INSERT 权限。通过 HTTP API 管理桥接桥接的管理均通过 REST API 完成相关 API 覆盖了桥接的创建、更新、查询、启停stop/restart与列表list等操作。详细的接口说明可参考 API Docs - Bridges连接器接口见其 Connectors 一节。结合 emqx_bridge_matrix.erl 的 Schema 定义MatrixDB 相关的接口字段组织如下HTTP 方法资源说明POST/connectors创建连接器post_connector含name、typematrix与连接参数GET/PUT/connectors/{name}查询 / 更新连接器get_connector/put_connectorPOST/bridges创建桥接post_bridge_v2含name、typematrix、connector引用、parameters.sql、resource_optsGET/PUT/bridges/{name}查询 / 更新桥接get_bridge_v2/put_bridge_v2GET 额外返回status、node_status等运行状态字段POST/bridges/{name}/start与/bridges/{name}/stop启停桥接GET/bridges列出全部桥接以创建 MatrixDB 桥接为例请求体结构为{ name: my_matrix_bridge, type: matrix, connector: my_matrix_connector, enable: true, parameters: { sql: insert into t_mqtt_msg(msgid, topic, qos, payload, arrived) values (${id}, ${topic}, ${qos}, ${payload}, TO_TIMESTAMP((${timestamp} :: bigint)/1000)) }, resource_opts: { batch_size: 100, batch_time: 100ms } }查询桥接时响应会在上述字段基础上附加运行状态status如connected、node_status各节点连接状态以及actions列表便于运维观察桥接健康度。实现要点小结从源码角度回顾 MatrixDB 桥接的几个关键实现事实类型注册matrix同时作为动作类型与连接器类型注册桥接类型列表为[matrix]见 emqx_bridge_matrix_action_info.erl 与 emqx_bridge_matrix_connector_info.erlSchema 复用动作与连接器的字段定义全部委托给 PostgreSQL 桥接/连接器的 Schema 模块命名空间为bridge_matrix见 emqx_bridge_matrix.erl底层驱动连接器的资源回调模块为emqx_postgresql即基于 epgsql 的 PostgreSQL 连接与 SQL 执行链路直接为 MatrixDB 服务批量写入动作层通过batch_size/batch_time攒批降低高频 IoT 消息逐条写入的开销。依赖关系上应用仅依赖emqx_connector、emqx_resource与emqx_bridge三个伞形应用见 mix.exs充分复用 EMQX 统一的连接器与资源管理框架因此开发与维护成本极低。如需参与该模块的开发请遵循仓库的 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),仅供参考