Skip to content

快速开始:将 MQTT 数据写入 EMQX Tables

本指南通过工厂设备监控示例,介绍如何创建 EMQX Broker 和 EMQX Tables 部署、使用规则和 SQL 模板写入 MQTT 数据,以及在 Cloud Console 中查询结果。完成本指南后,您将建立一条从 MQTT 客户端到 EMQX Broker,再到 EMQX Tables 的数据链路。

本指南适用于部署在阿里云上的 EMQX Tables 基础版和 EMQX v6 的非 Serverless Broker 部署。Broker 数据集成通过 Arrow Flight SQL 写入数据,必须预先创建目标表。

准备工作

  • 准备一个 EMQX Cloud 账号和项目。
  • 确保账号可以创建 EMQX v6 的非 Serverless Broker 部署和 EMQX Tables 基础版部署。
  • 如果使用已有部署,请确保 Broker 和 Tables 位于同一项目,选择相同的云平台区域,并关联同一网络。

示例场景:工厂设备监控

工厂设备向 factory/+/metrics 主题发送 JSON 消息,包含以下字段:

字段含义示例
machine_id设备标识machine_001
production_line生产线line_a
temperature温度36.5
vibration振动值0.02
machine_status运行状态running

规则会补充 Broker 接收消息的毫秒时间戳 ts,并将上述字段写入 machine_metrics 表。

创建 Broker 和 Tables 部署

  1. 登录 Cloud Console,进入用于本示例的项目。
  2. 如果项目中没有可用的 Broker,在 EMQX Broker 区域点击新建部署,创建弹性专有版部署。在 EMQX 版本中选择 v6,在云平台中选择阿里云,并选择 Tables 支持的区域
  3. 等待 Broker 部署进入运行中状态。Tables 创建页只会显示正在运行的 Broker 所使用的可用网络。
  4. EMQX Tables 区域点击新建部署,选择基础版,并选择与 Broker 相同的云平台区域
  5. 网络关联中选择 Broker 使用的网络。Broker 和 Tables 必须关联同一网络才能使用快速配置建立内网连接。
  6. 选择 2 vCPU / 8 GB 试用规格,并按需设置部署名称。试用资格、包含额度和到期处理参见定价计费
  7. 确认试用期限和配置后,点击部署,等待 Tables 部署进入运行中状态。

创建 EMQX Tables

Tables 创建选项的详细说明参见创建 EMQX 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 语句文档

配置数据集成

在 Broker 中使用快速配置选择同一项目中的 Tables 部署。Cloud Console 会自动填充所选部署的连接信息,无需从 Tables 部署概览中复制地址或凭据。

创建连接器

  1. 进入 Broker 部署,点击数据集成
  2. 点击新建连接器,选择 EMQX Tables
  3. 保持选中快速配置,选择本项目中的目标 Tables 部署。
  4. 确认数据库名字public
  5. 点击测试连接,确认连接成功后点击新建

连接器快速配置

创建规则和写入动作

  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"

    这条 SQL 在 Broker 规则引擎中执行,用于匹配 factory/+/metrics 主题,并从 MQTT 消息中提取写入 Tables 的字段。timestamp 是 Broker 接收消息的毫秒时间戳,在规则输出中重命名为 ts

    配置规则 SQL

  3. 点击下一步,添加使用上述连接器的输出动作。

  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. 确认模板变量与规则输出字段一致。

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

SQL 模板用于参数绑定。占位符不加引号,模板末尾不加分号。ts 与目标表中的 TIMESTAMP(3) 字段对应。Arrow Flight SQL 写入方式参见添加 Arrow Flight SQL Sink,连接参数以 Cloud Console 的快速配置为准。

配置 EMQX Tables 写入动作

发布 MQTT 消息

  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. 点击发布

使用在线调试发布 MQTT 消息

在 Tables 中查询数据

  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. 点击执行查询

查询结果应包含刚发布的 machine_001 记录,温度为 36.5、振动值为 0.02、状态为 runningts 为 Broker 接收该消息的时间。

查询 machine_metrics 示例数据

如果未查到记录,请依次检查:

  • 规则是否已启用,发布主题是否与 factory/+/metrics 匹配。
  • 连接器是否连接成功。
  • 规则动作是否执行成功,以及动作错误信息。
  • machine_metrics 是否位于连接器指定的数据库中,表字段和数据类型是否与 SQL 模板一致。

至此,MQTT 消息已经由 Broker 规则提取,并通过 Arrow Flight SQL 写入 EMQX Tables。

后续步骤