Skip to content

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。

quick_start_add_flow

注意

保存后需要重启 EMQX Edge 使配置生效。例如,如果 EMQX Edge 运行在 Docker 中:

bash
docker restart emqx-edge

重启后,Exchange、Ring Queue 和持久化组件会自动初始化。

发布测试数据 ​

使用 MQTTX CLI 向 canudp 主题发布 20 条测试消息:

bash
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 文件。检查配置目录中是否生成文件:

bash
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 主机的系统时间可能与你的本地时间不同。

quick_start_query

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

如需通过程序查询持久化数据,请参见 MQTT Stream Query。