# MQTT Stream 用户指南

本页介绍如何配置和操作 MQTT Stream，包括 Flow 设置、数据持久化和历史查询。

有关内部架构和 Stream 插件开发，请参见 [MQTT Stream 设计与实现](mqtt-stream-design.md)。

## MQTT Stream 的工作原理

配置 MQTT Stream 后，符合条件的 MQTT 消息会通过旁路通道采集，不影响正常消息转发，并被路由到本地数据管道。

从用户视角看，工作流如下：

1. EMQX Edge 正常接入并转发 MQTT 消息。
2. 匹配配置主题的消息会被复制到内部 Exchange。
3. 消息首先使用 Ring Queue 缓存在内存中。
4. 当缓冲区达到配置阈值时，会触发批处理。批次数据由 Stream 插件编码，并写入本地存储。
5. 后续可按时间范围查询和回放历史数据。

该工作流的设计目标是：

- **高性能**：避免逐条消息写磁盘，降低 I/O 开销
- **高可靠性**：即使网络故障也能保留数据
- **可查询性**：支持按时间范围检索历史数据

## 配置 MQTT Stream

每个配置的 Stream 称为一个 Flow，用于定义完整的数据管道：捕获哪个 MQTT 主题、如何缓冲数据以及将数据持久化到哪里。你可以通过 Dashboard 配置 Flow，也可以直接编辑配置文件。

### 通过 Dashboard 配置

在 Dashboard 中进入 **MQTT Stream**，点击 **Flow** 标签页。点击 **Add** 打开 Add Flow 表单。

::: tip 注意
配置变更仅在重启 EMQX Edge 后生效。
:::

按以下说明填写表单字段。

#### Basic

| Field | Description |
|---|---|
| **Instance Name** | 该 Flow 实例的可选显示名称。 |
| **Flow URL** | Flow 的内部通信地址。它不是 MQTT 客户端使用的 Broker 地址。除非另有说明，请保留默认值 `tcp://127.0.0.1:10000`。 |
| **Query Frequency** | 限制该 Flow 的查询频率。值 N 表示每 5 秒最多允许 N 次查询。例如，`5` 表示每 5 秒最多 5 次查询。留空则使用系统默认值。 |

#### Flow

| Field | Description |
|---|---|
| **Topic** | 要捕获的 MQTT 主题。只有发布到该主题的消息会进入 MQTT Stream 管道。不支持 `/` 字符。 |
| **Stream Type** | 用于数据编码和解码的 Stream 插件。`0` = RAW Stream；`1` = SPI Stream；`2` = CANP Stream。如果后端注册了自定义插件，请输入其插件 ID。 |

#### Ring Queue

| Field | Description |
|---|---|
| **Capacity** | Ring Queue 可容纳的最大消息数。达到该阈值时，会触发 Full Operation 行为。 |
| **Full Operation** | 定义 Ring Queue 达到容量时的处理方式。当前仅支持 `2 - Return to AIO and Write to File`，且不可修改。 |

#### Parquet

| Field | Description |
|---|---|
| **Compression** | Parquet 文件压缩算法。`zstd` 是默认值，推荐用于在速度和压缩率之间取得较好平衡。其他选项包括：`uncompressed`、`snappy`、`gzip`、`brotli`、`lz4`。 |
| **Directory** | Parquet 文件存储目录。必填。默认值为 `/tmp/edge-parquet`。如果目录不存在，系统会自动创建。 |
| **File Name Prefix** | 生成的 Parquet 文件名前缀，可选。请使用普通文件名字符，避免路径分隔符。 |
| **File Count** | 保留的 Parquet 文件最大数量。达到限制后会执行文件轮转。 |
| **Total Written File Size** | 所有已写入 Parquet 文件的总大小限制。可从下拉框中选择单位（MB、GB 等）。 |

#### Parquet 加密

启用 **Enable Encryption** 后，会为 Parquet 配置生成加密配置块，用于控制文件加密。以下字段将可用：

| Field | Description |
|---|---|
| **Key ID** | 密钥检索元数据，用于标识正在使用的加密密钥。 |
| **Encryption Key** | Parquet 文件的加密密钥。支持原始 AES 密钥，以及后端支持的 Base64/wrapped key。 |
| **Encryption Type** | 加密算法。支持值：`AES_GCM_CTR_V1` 和 `AES_GCM_V1`。默认使用 `AES_GCM_V1` 语义。 |

![Add Flow form](./assets/mqtt-stream-add-flow.png)

### 保存并重启

点击 **Save** 保存 Flow 配置，然后重启 EMQX Edge 使变更生效。重启后：

- MQTT Broker 开始监听客户端连接。
- Exchange、Ring Queue 和持久化组件会自动初始化。
- MQTT Stream 数据管道开始生效。

### 通过配置文件配置

除了 Dashboard，也可以直接在 NanoMQ 配置文件中配置 MQTT Stream。这是生产部署或容器化环境中常见的做法，适用于将配置作为代码管理的场景。

以下示例为 `canudp` 主题配置一个带 Ring Queue 和 Parquet 持久化的 Exchange：

```hocon
exchange_client.mq1 {
    exchange_url = "tcp://127.0.0.1:10000"
    limit_frequency = 5

    exchange {
        topic = "canudp"
        name = "canudp"
        streamType = 0

        ringbus {
            name = "ringbus"
            cap = 1000
            fullOp = 2  # Only supported value: 2 = RB_FULL_RETURN (return all messages to upper layer)
        }

        parquet {
            compress = zstd
            dir = "/tmp/edge-parquet"
            file_name_prefix = "edge"
            file_count = 1500
            file_size = 100MB
        }
    }
}
```

关键配置参数如下：

| Parameter | Description |
|---|---|
| `exchange_url` | Exchange server 的内部地址。必须与 Dashboard 中配置的 Flow URL 一致。 |
| `exchange.topic` | 要捕获的 MQTT 主题。 |
| `exchange.streamType` | Stream 编码类型。`0` = RAW，`1` = SPI，`2` = CANP。 |
| `ringbus.cap` | Ring Queue 容量，单位为消息数。 |
| `ringbus.fullOp` | Ring Queue 满时的行为。必须设置为 `2`（Return to AIO and Write to File）。 |
| `parquet.dir` | Parquet 文件写入目录。如果目录不存在，EMQX Edge 会自动创建。 |
| `parquet.file_name_prefix` | 生成的 Parquet 文件名前缀，可选。请使用普通文件名字符，避免路径分隔符。 |

编辑配置文件后，重启 EMQX Edge 使变更生效。

## Stream Types

`streamType` 参数决定 MQTT 消息 Payload 写入 Parquet 时如何编码，并定义生成的 Parquet Schema。内置三种 Stream 类型：

| streamType | Name | Parquet Schema | Use Case |
|:---:|---|---|---|
| `0` | RAW | `ts`, `data` | 通用 MQTT 消息捕获；Payload 原样存储 |
| `1` | SPI | `ts` + 每种 packet type ID 对应的动态列（4 位十六进制） | 具有多个 ECU packet type 的车辆 SPI 总线数据 |
| `2` | CANP | `ts` + 每个 CAN ID 对应的动态列（4 位十六进制） | 具有多个 CAN ID 的 CAN/CAN-FD 总线数据 |

### RAW Stream（streamType = 0）

RAW 是默认 Stream 类型。每条 MQTT 消息会写入为一行，包含两个固定列：

- `ts`：消息时间戳（uint64，毫秒）
- `data`：原始 Payload 字节（二进制）

当你不需要解析消息结构，并希望原样存储 Payload 时，可使用 RAW。

### SPI Stream（streamType = 1）

SPI Stream 解析车辆 SPI 总线帧，并将每种 packet type 存储为单独的 Parquet 列。预期帧格式如下：

```text
| 0x55 (1B) | type+id (2B) | update+len (1B) | payload (len B) |
              type (4 bit) | id (12 bit) | update (1 bit) | len (7 bit)
```

Parquet Schema 由 `ts` 加上数据中观察到的每个 packet type ID 对应的一列组成。列名为 4 位十六进制字符串，例如 `0055`、`0057`。SPI Stream 适用于车载网关通过共享 SPI 总线从多个 ECU 采集数据的场景。

### CANP Stream（streamType = 2）

CANP Stream 解析 CAN 协议帧，并将每个 CAN ID 存储为单独的 Parquet 列。预期帧格式如下：

```text
| ts (8B) | len (4B) | [ tsdiff (1B) | busid (1B) | canid (2B) | len (2B) | payload (len B) ] ... |
```

Parquet Schema 由 `ts` 加上数据中观察到的每个 CAN ID 对应的一列组成。列名为 4 位十六进制字符串，例如 `0123`、`0456`。CANP Stream 还支持 v2.1 紧凑格式，可将具有相同 CAN ID 的多帧打包到同一个单元格中。CANP Stream 适用于车端边缘节点上的 CAN/CAN-FD 总线数据采集。

## 验证 MQTT Stream

本节通过发送测试数据并观察持久化输出来验证 MQTT Stream。

### 发送测试数据

使用 MQTTX CLI 向配置主题发布消息。以下示例使用单个持久连接向 `canudp` 主题发布消息：

```bash
mqttx bench pub -t "canudp" -h 127.0.0.1 -p 1883 -m "message" -L 1500 -c 1
```

消息数量必须超过 Ring Queue 的 **Capacity** 值，才能触发批量刷新到磁盘。使用单个持久连接（`-c 1`）可避免连接频繁创建和销毁影响 Ring Queue 管道。请将 `canudp` 替换为你的 Flow 中配置的主题。

### Parquet 文件命名

生成的 Parquet 文件遵循以下命名约定：

```text
{prefix}_{topic}-{start_ts}~{end_ts}_{seq}_{hash}.parquet
```

| Field | Description |
|---|---|
| `prefix` | 配置中的 `file_name_prefix` 值。 |
| `topic` | Exchange 主题名称。 |
| `start_ts` | 文件中最早消息的时间戳（毫秒）。 |
| `end_ts` | 文件中最新消息的时间戳（毫秒）。 |
| `seq` | 单调递增的文件序号。 |
| `hash` | 用于去重和完整性校验的内容哈希。 |

示例：

```text
edge_canudp-1781145590395~1781145948876_0_3ff7819a161625c40d495038fd85c248.parquet
```

### 查看日志和持久化数据

当 Ring Queue 达到容量时：

- 触发批处理。
- 缓冲批次会被编码并写入本地 Parquet 文件。

相关日志条目可确认 MQTT Stream 管道已激活。

在配置目录中会生成 Parquet 文件。每个文件通常对应 Ring Queue 返回的一个批次。

### 验证检查清单

满足以下条件时，说明 MQTT Stream 工作正常：

- MQTT 消息发布成功。
- Ring Queue 达到阈值并触发批处理。
- 配置目录中出现 Parquet 文件。

可使用以下工具检查 Parquet 文件：

- `parquet-tools`
- Pandas / PyArrow (Python)
- Apache Spark 或 Flink

## 查询历史数据

数据持久化后，MQTT Stream 支持基于时间范围的历史查询，可用于故障排查、回放和分析。查询通过 Dashboard 执行，不影响实时 MQTT 消息转发。

### 在 Dashboard 中查询数据

在 Dashboard 中进入 **MQTT Stream**，点击 **Query** 标签页。

填写查询字段：

| Field | Description |
|---|---|
| **Topic** | 要查询的主题。必填。不支持 `/` 字符。 |
| **Time Range** | 必填。使用日期时间选择器选择开始和结束时间，或选择预设范围：Last 5 minutes、Last 15 minutes、Last 1 hour、Last 24 hours 或 Last 7 days。 |
| **Schema** | 可选。指定要返回的列，以逗号分隔。留空则返回所有可用列。可用列取决于该主题配置的 Stream 类型。 |

点击 **Exact Search** 执行查询。匹配指定主题和时间范围的结果会显示在结果区域。

![Query results](./assets/mqtt-stream-query-results.png)

Payload 会以 Base64 编码字符串显示。点击任意行的 **View**，并切换到 **Decoded Text**，即可查看原始消息内容。

![Payload Viewer](./assets/mqtt-stream-payload-viewer.png)

#### 不同 Stream 类型的 Schema 值

Schema 是列过滤器：设置后，仅返回指定列；留空时，返回所有可用列。

无论 Stream 类型如何，`ts` 列（毫秒级时间戳）始终可用。其他列取决于 Flow 的 `streamType`：

| Stream Type | Available Columns | Example Schema Value |
|---|---|---|
| RAW (`0`) | `ts`, `data` | `data` 或 `ts,data` |
| SPI (`1`) | `ts` + 数据中存在的 packet type ID（4 位十六进制） | `ts,0055,0057` |
| CANP (`2`) | `ts` + 数据中存在的 CAN ID（4 位十六进制） | `ts,0123,0456` |

对于 SPI 和 CANP 类型，可用列名取决于写入数据时观察到的 packet type ID 或 CAN ID。要发现可用列，请先在不指定 schema 的情况下运行查询，并检查返回字段名。

### 通过 REST API 查询

历史数据也可以通过 REST API 进行程序化查询。

完整 API 参考请参见 [MQTT Stream Query](../api/v4.md#mqtt-stream-query)。

**Endpoint**

```bash
GET /api/v4/mqtt_stream
```

**必填参数**

| Parameter | Type | Description |
|---|---|---|
| `topic` | String | 要查询的 MQTT 主题。必须与 Exchange 主题名称一致。 |
| `start_ts` | uint64 | 时间范围开始，包含该时间点（毫秒）。 |
| `end_ts` | uint64 | 时间范围结束，包含该时间点（毫秒）。 |

**可选参数**

| Parameter | Type | Default | Description |
|---|---|---|---|
| `schema` | String | 所有列 | 要返回的列列表，以逗号分隔，例如 `data` 或 `ts,0055`。 |
| `limit` | uint64 | `10000` | 返回记录的最大数量。 |
| `offset` | uint64 | `0` | 分页偏移量。 |

**响应**

```json
{
    "count": 20,
    "dataArray": [
        {"ts": 1781145590395, "data": "bWVzc2FnZSAx"},
        {"ts": 1781145596370, "data": "bWVzc2FnZSAy"}
    ]
}
```

- `count`：应用 `offset` 和 `limit` 后返回的记录总数。
- `dataArray`：结果行数组。列值为 Base64 编码。如果指定了 `schema`，仅返回请求的列。

**请求示例**

查询 RAW Stream 的所有列：

```bash
curl -i --basic -u admin:public -X GET \
  "http://localhost:8081/api/v4/mqtt_stream?topic=canudp&start_ts=1781145590395&end_ts=1781145948876"
```

查询指定 SPI packet type 列：

```bash
curl -i --basic -u admin:public -X GET \
  "http://localhost:8081/api/v4/mqtt_stream?topic=canudp&start_ts=1781145590395&end_ts=1781145948876&schema=ts,0055,0057"
```

分页查询结果（每页 100 条的第二页）：

```bash
curl -i --basic -u admin:public -X GET \
  "http://localhost:8081/api/v4/mqtt_stream?topic=canudp&start_ts=0&end_ts=9999999999999&limit=100&offset=100"
```

**错误码**

| HTTP Status | Meaning |
|---|---|
| `200 OK` | 查询成功。 |
| `204 No Content` | 指定时间范围内没有匹配数据。 |
| `400 Bad Request` | 缺少必填参数、时间戳格式无效或 schema 值无效。 |

### 典型使用场景

- 在设备或系统故障后回放历史数据。
- 补偿网络中断造成的数据缺口。
- 导出本地持久化数据用于离线分析。

## 故障排查

**EMQX Edge 启动失败并报错 `NULL s->ex_node!`**

配置中缺少 `ringbus.name` 或 `ringbus.cap` 字段，或字段为空。这两个字段都是必填项。缺少 `name` 会导致 Ring Queue 初始化失败，进而阻止 Exchange 正确绑定。请同时设置这两个字段并重启。

**未生成 Parquet 文件**

检查以下事项：

1. EMQX Edge 编译时已启用 `-DENABLE_PARQUET=ON`。
2. `parquet.dir` 路径可写。
3. 如果设置了 `parquet.file_name_prefix`，请确认它只包含普通文件名字符，且不包含路径分隔符。
4. 发布的消息数超过了 `ringbus.cap`。Ring Queue 未满时不会刷新。

**REST API 返回 `400 Bad Request`**

确认已提供 `topic`、`start_ts` 和 `end_ts`，且 `start_ts` 小于或等于 `end_ts`。`topic` 值必须与 Flow 中配置的 Exchange 主题名称完全一致。
