MQTT Stream 用户指南
本页介绍如何配置和操作 MQTT Stream,包括 Flow 设置、数据持久化和历史查询。
有关内部架构和 Stream 插件开发,请参见 MQTT Stream 设计与实现。
MQTT Stream 的工作原理
配置 MQTT Stream 后,符合条件的 MQTT 消息会通过旁路通道采集,不影响正常消息转发,并被路由到本地数据管道。
从用户视角看,工作流如下:
- EMQX Edge 正常接入并转发 MQTT 消息。
- 匹配配置主题的消息会被复制到内部 Exchange。
- 消息首先使用 Ring Queue 缓存在内存中。
- 当缓冲区达到配置阈值时,会触发批处理。批次数据由 Stream 插件编码,并写入本地存储。
- 后续可按时间范围查询和回放历史数据。
该工作流的设计目标是:
- 高性能:避免逐条消息写磁盘,降低 I/O 开销
- 高可靠性:即使网络故障也能保留数据
- 可查询性:支持按时间范围检索历史数据
配置 MQTT Stream
每个配置的 Stream 称为一个 Flow,用于定义完整的数据管道:捕获哪个 MQTT 主题、如何缓冲数据以及将数据持久化到哪里。你可以通过 Dashboard 配置 Flow,也可以直接编辑配置文件。
通过 Dashboard 配置
在 Dashboard 中进入 MQTT Stream,点击 Flow 标签页。点击 Add 打开 Add Flow 表单。
注意
配置变更仅在重启 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 语义。 |

保存并重启
点击 Save 保存 Flow 配置,然后重启 EMQX Edge 使变更生效。重启后:
- MQTT Broker 开始监听客户端连接。
- Exchange、Ring Queue 和持久化组件会自动初始化。
- MQTT Stream 数据管道开始生效。
通过配置文件配置
除了 Dashboard,也可以直接在 NanoMQ 配置文件中配置 MQTT Stream。这是生产部署或容器化环境中常见的做法,适用于将配置作为代码管理的场景。
以下示例为 canudp 主题配置一个带 Ring Queue 和 Parquet 持久化的 Exchange:
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 列。预期帧格式如下:
| 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 列。预期帧格式如下:
| 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 主题发布消息:
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 文件遵循以下命名约定:
{prefix}_{topic}-{start_ts}~{end_ts}_{seq}_{hash}.parquet| Field | Description |
|---|---|
prefix | 配置中的 file_name_prefix 值。 |
topic | Exchange 主题名称。 |
start_ts | 文件中最早消息的时间戳(毫秒)。 |
end_ts | 文件中最新消息的时间戳(毫秒)。 |
seq | 单调递增的文件序号。 |
hash | 用于去重和完整性校验的内容哈希。 |
示例:
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 执行查询。匹配指定主题和时间范围的结果会显示在结果区域。

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

不同 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。
Endpoint
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 | 分页偏移量。 |
响应
{
"count": 20,
"dataArray": [
{"ts": 1781145590395, "data": "bWVzc2FnZSAx"},
{"ts": 1781145596370, "data": "bWVzc2FnZSAy"}
]
}count:应用offset和limit后返回的记录总数。dataArray:结果行数组。列值为 Base64 编码。如果指定了schema,仅返回请求的列。
请求示例
查询 RAW Stream 的所有列:
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 列:
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 条的第二页):
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 文件
检查以下事项:
- EMQX Edge 编译时已启用
-DENABLE_PARQUET=ON。 parquet.dir路径可写。- 如果设置了
parquet.file_name_prefix,请确认它只包含普通文件名字符,且不包含路径分隔符。 - 发布的消息数超过了
ringbus.cap。Ring Queue 未满时不会刷新。
REST API 返回 400 Bad Request
确认已提供 topic、start_ts 和 end_ts,且 start_ts 小于或等于 end_ts。topic 值必须与 Flow 中配置的 Exchange 主题名称完全一致。