# 将 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 位于同一项目并关联同一网络。仅使用相同的云平台和区域不能满足集成条件。
- 了解[数据集成](./introduction.md)和[规则](./rules.md)的基本概念。

同一网络只能包含一个 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 表](./_assets/emqx_tables_create_table.png)

本示例采用时序表的时间键和分区定义。更多建表语法参见 [Datalayers CREATE 语句文档](https://docs.datalayers.cn/datalayers/latest/sql-reference/statements/create.html)。

## 创建 EMQX Tables 连接器

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

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

![创建 EMQX Tables 连接器](./_assets/emqx_tables_connector.png)

## 创建规则和写入动作

以下示例假设已经在 `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](./_assets/emqx_tables_rule_sql.png)

### 添加写入动作

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 写入动作](./_assets/emqx_tables_rule_action.png)

## 验证数据写入

### 发布测试消息

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 消息](./_assets/emqx_tables_publish_message.png)

### 查询写入结果

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 写入结果](./_assets/emqx_tables_query_result.png)

## 排查数据写入问题

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

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

## 后续步骤

- [查询 EMQX Tables 数据](../emqx_tables/emqx_tables_query_guide.md)
- [管理 EMQX Tables 部署](../emqx_tables/emqx_tables_manage_deployment.md)
