MQTT Stream 快速上手
本快速上手以典型车端边缘部署为例,演示如何在 EMQX Edge 中启用 MQTT Stream、将 MQTT 数据本地持久化,并验证历史查询。
完成本指南后,你将能够确认:
- MQTT Stream 已正确启用。
- 指定主题的 MQTT 消息已成功持久化到磁盘。
- 可以按时间范围查询历史数据。
场景:车端边缘数据持久化和回放
车端边缘部署是 MQTT Stream 的典型使用场景。
在典型车载环境中,车载网关连接多个 CAN / CAN-FD 总线,并持续从电池、动力系统、整车控制单元等子系统采集高频传感器和诊断数据。这些数据会作为 MQTT 消息发布到本地 EMQX Edge 实例。
在该场景中,MQTT Stream 用于在车辆端构建可靠的本地数据管道:
- 旁路捕获指定主题的 MQTT 消息,例如
canudp - 将高频数据按批次缓存在本地
- 将数据按批次持久化到本地磁盘
- 按时间窗口回放历史数据,用于故障排查和分析
本快速上手使用最小配置端到端验证完整工作流。
前置条件
开始前,请确保:
- 已安装并运行 EMQX Edge 1.5.0 或更高版本。本指南示例使用 Docker,但也适用于任何受支持的安装方式。
- 已安装 MQTTX CLI。安装说明请参见 MQTTX CLI。
在 Dashboard 中添加 Flow
在 Dashboard 中进入 MQTT Stream,点击 Flow 标签页。点击 Add 打开 Add Flow 表单。
填写以下字段:
Basic
- Instance Name:输入该 Flow 实例的名称,例如
canudp-flow。 - Flow URL:
tcp://127.0.0.1:10000(保持默认值)
Flow
- Topic:
canudp - Stream Type:
0
Ring Queue
- Capacity:
10 - Full Operation:
2 - Return to AIO and Write to File
Parquet
- Directory:
/tmp/edge-parquet(或主机上的任意绝对路径)
其他字段保持默认值。点击 Save。

注意
保存后需要重启 EMQX Edge 使配置生效。例如,如果 EMQX Edge 运行在 Docker 中:
docker restart emqx-edge重启后,Exchange、Ring Queue 和持久化组件会自动初始化。
发布测试数据
使用 MQTTX CLI 向 canudp 主题发布 20 条测试消息:
mqttx bench pub -t "canudp" -h 127.0.0.1 -p 1883 -m "message" -L 20 -c 1消息数量必须超过 Ring Queue 的 Capacity 值,才能触发批量刷新到磁盘。使用单个持久连接(-c 1)可避免连接频繁创建和销毁影响 Ring Queue 管道。
验证持久化数据
当 Ring Queue 阈值达到后,缓冲消息会按批次写入本地 Parquet 文件。检查配置目录中是否生成文件:
docker exec emqx-edge ls /tmp/edge-parquet如果存在 Parquet 文件,说明 MQTT Stream 正常工作。
查询历史数据
数据持久化后,你可以在 Dashboard 中查询。
进入 MQTT Stream,点击 Query 标签页。在 Topic 字段中输入 canudp,设置覆盖测试数据发布时间段的 Time Range,并将 Schema 留空。点击 Exact Search。
TIP
如果没有结果,请扩大时间范围。EMQX Edge 主机的系统时间可能与你的本地时间不同。

如果返回结果,说明完整 MQTT Stream 管道已正常工作:从消息捕获、Ring Queue 批处理,到 Parquet 持久化和 Dashboard 查询均已打通。Payload 会以 Base64 编码字符串显示。点击任意行的 View,并切换到 Decoded Text,即可查看原始消息内容。
如需通过程序查询持久化数据,请参见 MQTT Stream Query。