Skip to content

MQTT Stream ​

MQTT Stream 是 EMQX Edge 内置的数据持久化和查询功能。配置主题上的消息会被旁路捕获、缓存在内存中,并按批次写入磁盘,不影响实时 MQTT 转发。持久化后的数据可直接在 Dashboard 中按时间范围查询。

为什么需要 MQTT Stream ​

在边缘计算场景中,通常会产生大量高频 MQTT 消息,例如工业遥测数据、车辆 CAN 总线数据和设备状态指标。

如果只依赖实时上行到云端,通常会面临以下挑战:

  • 网络不稳定或不可用时,关键数据可能丢失。
  • 事件发生后,难以完整还原事件前后的数据。
  • 边缘端缺少统一的本地缓冲和持久化机制。

MQTT Stream 旨在解决这些问题。借助 MQTT Stream,EMQX Edge 可以在边缘端建立完整的本地数据管道:

MQTT 接入 -> 本地缓冲 -> 批量持久化 -> 历史查询

该管道独立于外部系统运行,并且在弱网或离线条件下仍可持续工作,为故障排查、分析和系统优化提供可靠的数据基础。

核心功能 ​

  • 非侵入式捕获:匹配的消息会被旁路复制到数据管道中,不影响正常 MQTT 转发。
  • 内存 Ring Buffer:消息在写入前先按批次缓存在内存中,避免逐条消息写磁盘。
  • 批量 Parquet 持久化:批次数据会被异步编码并写入本地 Parquet 文件,适合按时间范围扫描和下游分析。
  • 压缩和加密:Parquet 文件支持多种压缩算法,并可选择启用静态加密。
  • Dashboard 查询:可直接在 Dashboard 中按主题和时间窗口查询历史数据。
  • 离线运行:无论网络是否可用,数据都可被缓冲和持久化。
  • 可插拔 Stream 插件:编码、存储 Schema 和查询行为由插件定义。内置 RAW、SPI 和 CANP 插件,并支持自定义插件。

典型使用场景 ​

  • 弱网或离线网络中的可靠数据保留:在连接不稳定的部署中,MQTT Stream 会在本地缓冲并持久化消息,防止数据丢失。

  • 事件后的数据分析和故障排查:通过按时间顺序存储高频 MQTT 数据,用户可以回放并分析历史数据,用于诊断设备或系统故障。

  • 高频边缘数据的本地归档:对于不适合实时上传到云端的数据,MQTT Stream 可在保证数据完整性的同时实现批量缓冲和结构化本地存储。

  • 端到端边云数据管道:本地持久化的数据后续可用于分析、导出或转发到云端系统,支持从边缘采集到反馈优化的闭环数据流。

了解更多 ​

MQTT Stream 从 EMQX Edge 1.5.0 开始作为商业增值功能提供。如果你需要在边缘端实现轻量级本地持久化和消息流管理,请联系我们了解更多信息。