Skip to content

将 MQTT 数据写入 EMQX Tables ​

通过 EMQX Broker 的数据集成,可以使用规则提取和转换 MQTT 消息字段,再通过 EMQX Tables 写入动作将结果保存到预先创建的时序表。本页介绍如何创建连接器、规则和写入动作,并验证数据写入结果。

本页适用于部署在阿里云上的 EMQX Tables 基础版,以及运行 EMQX v6 的非 Serverless Broker 部署。 数据通过 Arrow Flight SQL 写入目标表。

工作原理 ​

设备向 Broker 发布 MQTT 消息后,规则引擎按照主题匹配消息并提取字段。EMQX Tables 写入动作将规则输出绑定到 SQL 模板参数,并把数据写入目标表。

数据写入流程如下:

  1. 设备向 Broker 发布 JSON 格式的遥测消息。
  2. 规则匹配 MQTT 主题并提取需要写入的字段。
  3. EMQX Tables 动作将规则输出绑定到 SQL 模板参数。
  4. EMQX Tables 将数据写入预先创建的目标表。
  5. 使用 Tables 的数据查询页面或数据库客户端查询数据。

前提条件 ​

  • 已创建运行中的 EMQX v6 非 Serverless Broker 部署和阿里云 EMQX Tables 基础版部署。
  • Broker 和 Tables 位于同一项目并关联同一网络。仅使用相同的云平台和区域不能满足集成条件。
  • 了解数据集成和规则的基本概念。

同一网络只能包含一个 Broker 和一个 Tables 部署。Broker MQTT 认证与 Tables 数据库认证使用不同的凭据。

创建目标表 ​

以下示例使用默认数据库 public。如果部署显示其他数据库名称,请在创建目标表和查询数据时选择该数据库,并在配置连接器时填写该数据库名称。

  1. 进入 Tables 部署,在左侧导航栏中点击数据查询。

  2. 选择 public 数据库,在 SQL 编辑器中输入以下建表语句:

    sql
    CREATE 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;
  3. 点击执行查询。执行成功后,刷新左侧数据表列表,确认 machine_metrics 已创建。

ts 为毫秒精度的时间列,machine_id 用于标识设备。

创建 machine_metrics 表

本示例采用时序表的时间键和分区定义。更多建表语法参见 Datalayers CREATE 语句文档。

创建 EMQX Tables 连接器 ​

  1. 进入 Broker 部署,在左侧导航栏中点击数据集成。
  2. 点击新建连接器,选择 EMQX Tables。
  3. 保持选中快速配置,然后选择目标 Tables 部署。
  4. 确认数据库名字与目标表所在的数据库一致。快速配置会自动填入所选 Tables 部署的数据库名称。
  5. 点击测试连接。
  6. 确认连接测试成功后,点击新建。

快速配置只列出同一项目、关联同一网络且处于运行状态的 Tables 部署,并自动填充内部连接参数。不要将这些内部连接参数用于外部数据库客户端。

创建 EMQX Tables 连接器

创建规则和写入动作 ​

以下示例假设已经在 public 数据库中创建 machine_metrics 表。 该表包含 ts、machine_id、production_line、temperature、vibration 和 machine_status 字段。

创建规则 ​

  1. 连接器创建完成后,点击新建规则。

  2. 在 SQL 编辑器中输入以下规则 SQL:

    sql
    SELECT
        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。

配置规则 SQL

添加写入动作 ​

  1. 点击下一步,进入新建输出动作。

  2. 在使用连接器中选择已创建的 EMQX Tables 连接器。

  3. 确认动作类型为 EMQX Tables。

  4. 在 SQL 模板中输入:

    sql
    INSERT INTO machine_metrics (
        ts,
        machine_id,
        production_line,
        temperature,
        vibration,
        machine_status
    )
    VALUES (
        ${ts},
        ${machine_id},
        ${production_line},
        ${temperature},
        ${vibration},
        ${machine_status}
    )
  5. 确认 SQL 模板中的列名与目标表一致,占位符与规则输出字段一致。

  6. 点击确认完成动作配置,再完成规则创建。

SQL 模板使用参数绑定。占位符不加引号,模板末尾不加分号。ts 应与目标表中时间字段的精度一致。

配置 EMQX Tables 写入动作

验证数据写入 ​

发布测试消息 ​

  1. 进入 Broker 部署,点击诊断工具 -> 在线调试。

  2. 按页面提示使用自动生成的认证信息,或使用已在访问控制 -> 认证中添加的用户名和密码连接部署。

  3. 在消息区域中,将主题设置为 factory/line_a/metrics,将 QoS 设置为 QoS 0,并选择 JSON 格式。

  4. 输入以下 Payload:

    json
    {
      "machine_id": "machine_001",
      "production_line": "line_a",
      "temperature": 36.5,
      "vibration": 0.02,
      "machine_status": "running"
    }
  5. 点击发布。

  6. 在 Broker 的规则列表中打开目标规则,确认规则命中次数和动作成功次数增加。

使用在线调试发布 MQTT 消息

查询写入结果 ​

  1. 进入 Tables 部署,点击数据查询。

  2. 选择 public 数据库,在 SQL 编辑器中输入:

    sql
    SELECT ts, machine_id, production_line, temperature, vibration, machine_status
    FROM machine_metrics
    ORDER BY ts DESC
    LIMIT 10;
  3. 点击执行查询。

  4. 确认查询结果中包含 machine_001 记录,且字段值与发布的消息一致。

查询 EMQX Tables 写入结果

排查数据写入问题 ​

现象检查项
无法选择 Tables 部署Tables 是否位于同一项目、关联同一网络并处于运行状态
连接测试失败Broker 和 Tables 是否关联同一网络,目标数据库是否存在
规则未命中规则是否启用,消息主题是否匹配
动作执行失败SQL 模板、占位符、目标表字段和字段类型是否一致
查询不到数据查询的数据库和表是否与连接器及 SQL 模板一致

修正配置后重新发布消息,再次检查动作统计和 Tables 查询结果。

后续步骤 ​