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 用户指南:配置和操作 MQTT Stream,包括持久化和查询。
- MQTT Stream 快速上手:通过最小配置快速验证数据持久化和时间范围查询。
- MQTT Stream 设计与实现:了解内部架构、数据管道和 Stream 插件模型。
MQTT Stream 从 EMQX Edge 1.5.0 开始作为商业增值功能提供。如果你需要在边缘端实现轻量级本地持久化和消息流管理,请联系我们了解更多信息。