Skip to content

MQTT Stream 用户指南 ​

本页介绍如何配置和操作 MQTT Stream,包括 Flow 设置、数据持久化和历史查询。

有关内部架构和 Stream 插件开发,请参见 MQTT Stream 设计与实现。

MQTT Stream 的工作原理 ​

配置 MQTT Stream 后,符合条件的 MQTT 消息会通过旁路通道采集,不影响正常消息转发,并被路由到本地数据管道。

从用户视角看,工作流如下:

  1. EMQX Edge 正常接入并转发 MQTT 消息。
  2. 匹配配置主题的消息会被复制到内部 Exchange。
  3. 消息首先使用 Ring Queue 缓存在内存中。
  4. 当缓冲区达到配置阈值时,会触发批处理。批次数据由 Stream 插件编码,并写入本地存储。
  5. 后续可按时间范围查询和回放历史数据。

该工作流的设计目标是:

  • 高性能:避免逐条消息写磁盘,降低 I/O 开销
  • 高可靠性:即使网络故障也能保留数据
  • 可查询性:支持按时间范围检索历史数据

配置 MQTT Stream ​

每个配置的 Stream 称为一个 Flow,用于定义完整的数据管道:捕获哪个 MQTT 主题、如何缓冲数据以及将数据持久化到哪里。你可以通过 Dashboard 配置 Flow,也可以直接编辑配置文件。

通过 Dashboard 配置 ​

在 Dashboard 中进入 MQTT Stream,点击 Flow 标签页。点击 Add 打开 Add Flow 表单。

注意

配置变更仅在重启 EMQX Edge 后生效。

按以下说明填写表单字段。

Basic ​

FieldDescription
Instance Name该 Flow 实例的可选显示名称。
Flow URLFlow 的内部通信地址。它不是 MQTT 客户端使用的 Broker 地址。除非另有说明,请保留默认值 tcp://127.0.0.1:10000。
Query Frequency限制该 Flow 的查询频率。值 N 表示每 5 秒最多允许 N 次查询。例如,5 表示每 5 秒最多 5 次查询。留空则使用系统默认值。

Flow ​

FieldDescription
Topic要捕获的 MQTT 主题。只有发布到该主题的消息会进入 MQTT Stream 管道。不支持 / 字符。
Stream Type用于数据编码和解码的 Stream 插件。0 = RAW Stream;1 = SPI Stream;2 = CANP Stream。如果后端注册了自定义插件,请输入其插件 ID。

Ring Queue ​

FieldDescription
CapacityRing Queue 可容纳的最大消息数。达到该阈值时,会触发 Full Operation 行为。
Full Operation定义 Ring Queue 达到容量时的处理方式。当前仅支持 2 - Return to AIO and Write to File,且不可修改。

Parquet ​

FieldDescription
CompressionParquet 文件压缩算法。zstd 是默认值,推荐用于在速度和压缩率之间取得较好平衡。其他选项包括:uncompressed、snappy、gzip、brotli、lz4。
DirectoryParquet 文件存储目录。必填。默认值为 /tmp/edge-parquet。如果目录不存在,系统会自动创建。
File Name Prefix生成的 Parquet 文件名前缀,可选。请使用普通文件名字符,避免路径分隔符。
File Count保留的 Parquet 文件最大数量。达到限制后会执行文件轮转。
Total Written File Size所有已写入 Parquet 文件的总大小限制。可从下拉框中选择单位(MB、GB 等)。

Parquet 加密 ​

启用 Enable Encryption 后,会为 Parquet 配置生成加密配置块,用于控制文件加密。以下字段将可用:

FieldDescription
Key ID密钥检索元数据,用于标识正在使用的加密密钥。
Encryption KeyParquet 文件的加密密钥。支持原始 AES 密钥,以及后端支持的 Base64/wrapped key。
Encryption Type加密算法。支持值:AES_GCM_CTR_V1 和 AES_GCM_V1。默认使用 AES_GCM_V1 语义。

Add Flow form

保存并重启 ​

点击 Save 保存 Flow 配置,然后重启 EMQX Edge 使变更生效。重启后:

  • MQTT Broker 开始监听客户端连接。
  • Exchange、Ring Queue 和持久化组件会自动初始化。
  • MQTT Stream 数据管道开始生效。

通过配置文件配置 ​

除了 Dashboard,也可以直接在 NanoMQ 配置文件中配置 MQTT Stream。这是生产部署或容器化环境中常见的做法,适用于将配置作为代码管理的场景。

以下示例为 canudp 主题配置一个带 Ring Queue 和 Parquet 持久化的 Exchange:

hocon
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
        }
    }
}

关键配置参数如下:

ParameterDescription
exchange_urlExchange server 的内部地址。必须与 Dashboard 中配置的 Flow URL 一致。
exchange.topic要捕获的 MQTT 主题。
exchange.streamTypeStream 编码类型。0 = RAW,1 = SPI,2 = CANP。
ringbus.capRing Queue 容量,单位为消息数。
ringbus.fullOpRing Queue 满时的行为。必须设置为 2(Return to AIO and Write to File)。
parquet.dirParquet 文件写入目录。如果目录不存在,EMQX Edge 会自动创建。
parquet.file_name_prefix生成的 Parquet 文件名前缀,可选。请使用普通文件名字符,避免路径分隔符。

编辑配置文件后,重启 EMQX Edge 使变更生效。

Stream Types ​

streamType 参数决定 MQTT 消息 Payload 写入 Parquet 时如何编码,并定义生成的 Parquet Schema。内置三种 Stream 类型:

streamTypeNameParquet SchemaUse Case
0RAWts, data通用 MQTT 消息捕获;Payload 原样存储
1SPIts + 每种 packet type ID 对应的动态列(4 位十六进制)具有多个 ECU packet type 的车辆 SPI 总线数据
2CANPts + 每个 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 列。预期帧格式如下:

text
| 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 列。预期帧格式如下:

text
| 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 主题发布消息:

bash
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 文件遵循以下命名约定:

text
{prefix}_{topic}-{start_ts}~{end_ts}_{seq}_{hash}.parquet
FieldDescription
prefix配置中的 file_name_prefix 值。
topicExchange 主题名称。
start_ts文件中最早消息的时间戳(毫秒)。
end_ts文件中最新消息的时间戳(毫秒)。
seq单调递增的文件序号。
hash用于去重和完整性校验的内容哈希。

示例:

text
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 标签页。

填写查询字段:

FieldDescription
Topic要查询的主题。必填。不支持 / 字符。
Time Range必填。使用日期时间选择器选择开始和结束时间,或选择预设范围:Last 5 minutes、Last 15 minutes、Last 1 hour、Last 24 hours 或 Last 7 days。
Schema可选。指定要返回的列,以逗号分隔。留空则返回所有可用列。可用列取决于该主题配置的 Stream 类型。

点击 Exact Search 执行查询。匹配指定主题和时间范围的结果会显示在结果区域。

Query results

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

Payload Viewer

不同 Stream 类型的 Schema 值 ​

Schema 是列过滤器:设置后,仅返回指定列;留空时,返回所有可用列。

无论 Stream 类型如何,ts 列(毫秒级时间戳)始终可用。其他列取决于 Flow 的 streamType:

Stream TypeAvailable ColumnsExample Schema Value
RAW (0)ts, datadata 或 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

bash
GET /api/v4/mqtt_stream

必填参数

ParameterTypeDescription
topicString要查询的 MQTT 主题。必须与 Exchange 主题名称一致。
start_tsuint64时间范围开始,包含该时间点(毫秒)。
end_tsuint64时间范围结束,包含该时间点(毫秒)。

可选参数

ParameterTypeDefaultDescription
schemaString所有列要返回的列列表,以逗号分隔,例如 data 或 ts,0055。
limituint6410000返回记录的最大数量。
offsetuint640分页偏移量。

响应

json
{
    "count": 20,
    "dataArray": [
        {"ts": 1781145590395, "data": "bWVzc2FnZSAx"},
        {"ts": 1781145596370, "data": "bWVzc2FnZSAy"}
    ]
}
  • count:应用 offset 和 limit 后返回的记录总数。
  • dataArray:结果行数组。列值为 Base64 编码。如果指定了 schema,仅返回请求的列。

请求示例

查询 RAW Stream 的所有列:

bash
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 列:

bash
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 条的第二页):

bash
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 StatusMeaning
200 OK查询成功。
204 No Content指定时间范围内没有匹配数据。
400 Bad Request缺少必填参数、时间戳格式无效或 schema 值无效。

典型使用场景 ​

  • 在设备或系统故障后回放历史数据。
  • 补偿网络中断造成的数据缺口。
  • 导出本地持久化数据用于离线分析。

故障排查 ​

EMQX Edge 启动失败并报错 NULL s->ex_node!

配置中缺少 ringbus.name 或 ringbus.cap 字段,或字段为空。这两个字段都是必填项。缺少 name 会导致 Ring Queue 初始化失败,进而阻止 Exchange 正确绑定。请同时设置这两个字段并重启。

未生成 Parquet 文件

检查以下事项:

  1. EMQX Edge 编译时已启用 -DENABLE_PARQUET=ON。
  2. parquet.dir 路径可写。
  3. 如果设置了 parquet.file_name_prefix,请确认它只包含普通文件名字符,且不包含路径分隔符。
  4. 发布的消息数超过了 ringbus.cap。Ring Queue 未满时不会刷新。

REST API 返回 400 Bad Request

确认已提供 topic、start_ts 和 end_ts,且 start_ts 小于或等于 end_ts。topic 值必须与 Flow 中配置的 Exchange 主题名称完全一致。