将 MQTT 数据写入 EMQX Tables
通过 EMQX Broker 的数据集成,可以使用规则提取和转换 MQTT 消息字段,再通过 EMQX Tables 写入动作将结果保存到预先创建的时序表。本页介绍如何创建连接器、规则和写入动作,并验证数据写入结果。
本页适用于部署在阿里云上的 EMQX Tables 基础版,以及运行 EMQX v6 的非 Serverless Broker 部署。 数据通过 Arrow Flight SQL 写入目标表。
工作原理
设备向 Broker 发布 MQTT 消息后,规则引擎按照主题匹配消息并提取字段。EMQX Tables 写入动作将规则输出绑定到 SQL 模板参数,并把数据写入目标表。
数据写入流程如下:
- 设备向 Broker 发布 JSON 格式的遥测消息。
- 规则匹配 MQTT 主题并提取需要写入的字段。
- EMQX Tables 动作将规则输出绑定到 SQL 模板参数。
- EMQX Tables 将数据写入预先创建的目标表。
- 使用 Tables 的数据查询页面或数据库客户端查询数据。
前提条件
- 已创建运行中的 EMQX v6 非 Serverless Broker 部署和阿里云 EMQX Tables 基础版部署。
- Broker 和 Tables 位于同一项目并关联同一网络。仅使用相同的云平台和区域不能满足集成条件。
- 了解数据集成和规则的基本概念。
同一网络只能包含一个 Broker 和一个 Tables 部署。Broker MQTT 认证与 Tables 数据库认证使用不同的凭据。
创建目标表
以下示例使用默认数据库 public。如果部署显示其他数据库名称,请在创建目标表和查询数据时选择该数据库,并在配置连接器时填写该数据库名称。
进入 Tables 部署,在左侧导航栏中点击数据查询。
选择
public数据库,在 SQL 编辑器中输入以下建表语句:sqlCREATE TABLE IF NOT EXISTS machine_metrics ( ts TIMESTAMP(3) NOT NULL, machine_id STRING NOT NULL, production_line STRING, temperature DOUBLE, vibration DOUBLE, machine_status STRING, TIMESTAMP KEY (ts), PRIMARY KEY (machine_id, ts) ) PARTITION BY HASH (machine_id) PARTITIONS 1 ENGINE = TimeSeries;点击执行查询。执行成功后,刷新左侧数据表列表,确认
machine_metrics已创建。
ts 为毫秒精度的时间列,machine_id 用于标识设备。

本示例采用时序表的时间键和分区定义。更多建表语法参见 Datalayers CREATE 语句文档。
创建 EMQX Tables 连接器
- 进入 Broker 部署,在左侧导航栏中点击数据集成。
- 点击新建连接器,选择 EMQX Tables。
- 保持选中快速配置,然后选择目标 Tables 部署。
- 确认数据库名字与目标表所在的数据库一致。快速配置会自动填入所选 Tables 部署的数据库名称。
- 点击测试连接。
- 确认连接测试成功后,点击新建。
快速配置只列出同一项目、关联同一网络且处于运行状态的 Tables 部署,并自动填充内部连接参数。不要将这些内部连接参数用于外部数据库客户端。

创建规则和写入动作
以下示例假设已经在 public 数据库中创建 machine_metrics 表。 该表包含 ts、machine_id、production_line、temperature、vibration 和 machine_status 字段。
创建规则
连接器创建完成后,点击新建规则。
在 SQL 编辑器中输入以下规则 SQL:
sqlSELECT timestamp AS ts, payload.machine_id AS machine_id, payload.production_line AS production_line, payload.temperature AS temperature, payload.vibration AS vibration, payload.machine_status AS machine_status FROM "factory/+/metrics"该规则匹配
factory/+/metrics主题,并提取写入machine_metrics表所需的字段。timestamp是 Broker 接收消息的毫秒时间戳,在规则输出中重命名为ts。

添加写入动作
点击下一步,进入新建输出动作。
在使用连接器中选择已创建的 EMQX Tables 连接器。
确认动作类型为 EMQX Tables。
在 SQL 模板中输入:
sqlINSERT INTO machine_metrics ( ts, machine_id, production_line, temperature, vibration, machine_status ) VALUES ( ${ts}, ${machine_id}, ${production_line}, ${temperature}, ${vibration}, ${machine_status} )确认 SQL 模板中的列名与目标表一致,占位符与规则输出字段一致。
点击确认完成动作配置,再完成规则创建。
SQL 模板使用参数绑定。占位符不加引号,模板末尾不加分号。ts 应与目标表中时间字段的精度一致。

验证数据写入
发布测试消息
进入 Broker 部署,点击诊断工具 -> 在线调试。
按页面提示使用自动生成的认证信息,或使用已在访问控制 -> 认证中添加的用户名和密码连接部署。
在消息区域中,将主题设置为
factory/line_a/metrics,将 QoS 设置为 QoS 0,并选择 JSON 格式。输入以下 Payload:
json{ "machine_id": "machine_001", "production_line": "line_a", "temperature": 36.5, "vibration": 0.02, "machine_status": "running" }点击发布。
在 Broker 的规则列表中打开目标规则,确认规则命中次数和动作成功次数增加。

查询写入结果
进入 Tables 部署,点击数据查询。
选择
public数据库,在 SQL 编辑器中输入:sqlSELECT ts, machine_id, production_line, temperature, vibration, machine_status FROM machine_metrics ORDER BY ts DESC LIMIT 10;点击执行查询。
确认查询结果中包含
machine_001记录,且字段值与发布的消息一致。

排查数据写入问题
| 现象 | 检查项 |
|---|---|
| 无法选择 Tables 部署 | Tables 是否位于同一项目、关联同一网络并处于运行状态 |
| 连接测试失败 | Broker 和 Tables 是否关联同一网络,目标数据库是否存在 |
| 规则未命中 | 规则是否启用,消息主题是否匹配 |
| 动作执行失败 | SQL 模板、占位符、目标表字段和字段类型是否一致 |
| 查询不到数据 | 查询的数据库和表是否与连接器及 SQL 模板一致 |
修正配置后重新发布消息,再次检查动作统计和 Tables 查询结果。