MQTT Stream 设计与实现
本页说明 MQTT Stream 的内部工作机制,涵盖端到端数据管道、Stream 插件架构、自定义插件开发以及查询通道协议。
本文面向需要理解内部设计或扩展 MQTT Stream 默认行为的开发者和系统集成者。
端到端数据管道
启用 MQTT Stream 后,发布到配置主题的消息会进入独立的数据处理管道。整个过程中,正常 MQTT 转发不受影响。
启动:基于配置初始化
MQTT Stream 数据路径在启动时完全根据配置初始化。系统不会在运行时创建或修改管道。
启动期间:
- 从配置中解析 Exchange 定义、主题过滤器、Stream 插件类型和 Ringbus 参数。
- 为每条配置的数据路径创建一个 Exchange 实例和一个 Ringbus 缓冲区。
- 建立 Broker 与 Exchange 之间的本地通信通道。
- 将每条数据路径绑定到其配置的 Stream 插件。
该模型可保持运行时管道结构稳定,避免资源频繁变动,并确保高吞吐场景下的性能可预测。
消息捕获:Broker 与 Exchange 分离
Broker 和 Exchange 具有严格分离的职责。
Broker 负责 MQTT 协议处理和实时消息转发。当 PUBLISH 消息匹配配置主题时,Broker 会将其旁路复制到 MQTT Stream 管道中,不会修改原始消息或转发路径。
每条旁路复制的消息会在捕获时分配一个单调递增的毫秒级时间戳。该时间戳作为消息在整个管道中的全局 key(用于缓冲、持久化和查询),后续无需再进行排序或去重。
Exchange 接收旁路复制的消息,并仅负责缓冲、批处理、持久化和历史查询。Broker 不感知这些操作。
Ringbus 缓冲和批处理触发
逐条将高频消息写入磁盘会严重限制吞吐量。MQTT Stream 使用 Ringbus 作为内存中的第一级缓冲,用于吸收写入压力。
Ringbus 按时间戳顺序接收消息,直到达到配置容量。它有意不感知编码和存储语义:其唯一职责是快速接收数据,并按批次交给上层处理。
当缓冲区达到容量时,fullOp 配置决定后续行为。推荐模式是 FULL_RETURN(2 - Return to AIO):
- 整个批次一次性交给上层处理层。
- Ringbus 立即清空并继续接收新消息。
- 所有编码和持久化逻辑都保留在上层,使 Ringbus 保持简单,并便于支持多个后端。
Stream 插件编码和 Parquet 持久化
当 Ringbus 返回一个批次后,管道会执行以下步骤:
- 将批次规范化为
stream_data_in结构,其中每条消息包含 payload、长度和时间戳 key。 - 调用
stream_encode(streamType, ...)分发到配置的 Stream 插件。 - 插件将批次编码为可存储的表示形式,例如列式 Parquet 结构。
- 编码结果会异步写入本地 Parquet 文件。
使用 Parquet 作为存储格式时,面向批次的写入可最大化列式压缩效率,减少文件碎片,并支持快速按时间范围扫描。
查询路径
历史查询通过 Exchange 的本地查询接口处理,不会经过 MQTT 转发路径。一次查询会同时覆盖仍在内存中(Ringbus)的数据和已经持久化到磁盘的数据。Exchange 合并这些来源的数据,并将其交给 Stream 插件解码,然后再向调用方返回结果。
Stream 插件架构
Stream 插件是 MQTT Stream 中数据语义的唯一来源。核心管道不假定固定的存储格式或查询协议;这些能力都委托给插件实现。
每个插件实现三个接口:
| Interface | Invocation point | Purpose |
|---|---|---|
encode | Ringbus 满后、写入前 | 将一批消息编码为可持久化结构 |
cmd_parser | 查询开始时 | 将查询命令字符串解析为结构化参数 |
decode | 组装查询结果时 | 将持久化记录解码为最终结果格式 |
MQTT Stream 包含两个内置插件:
- RAW Stream (
streamType = 0):通用类型。以最少转换传递 MQTT Payload。适用于没有复杂语义要求的高频数据。 - SPI Stream (
streamType = 1):面向汽车和 SDV 场景。支持自定义 key 结构和查询协议。
Exchange 配置中的 streamType 字段决定整条数据路径使用哪个插件,包括编码、持久化、查询解析和结果解码。
编码流程
当 Ringbus 触发 FULL_RETURN,或显式调用 hook_last_flush 刷新剩余数据时,批次会传递给 flush_smsg_to_disk(),该函数会:
- 根据返回的消息构建
stream_data_in,并使用nni_msg_get_timestamp()作为每条记录的 key。 - 调用
stream_encode(streamType, sdata),分发到插件的encode实现。对于 RAW,raw_encode会将stream_data_in转换为stream_data_out,并调用parquet_data_alloc构建列式结构。 - 编码结果(
parquet_data*)会传递给parquet_object_alloc+parquet_write_batch_async,用于异步写入磁盘。
解码流程
当查询命令到达 Exchange 时:
stream_cmd_parser(streamType, keystr)将命令字符串解析为cmd_data结构,包含is_sync、start_key、end_key和schema。对于 RAW Stream,默认 schema 为{"ts", "data"}(参见raw_stream.c::parse_input_cmd)。- Exchange 通过
exchange_client_get_msgs_fuzz从 Ringbus 内存中检索匹配记录,并通过parquet_get_data_packets_in_range_by_column从 Parquet 文件中检索匹配记录。 - 对每个
parquet_data_ret*,调用stream_decode(streamType, ...)将列式数据转换为stream_decoded_data,即适合调用方使用的连续缓冲区。 - 结果会组装为 NNG 消息,并以增量方式(async)或聚合方式(sync)返回。
自定义 Stream 插件开发
自定义 Stream 插件无需修改核心管道即可实现。插件需要实现三个函数,并在启动时注册。
接口结构体
/* Input to encode: one batch of messages from Ringbus */
struct stream_data_in {
void **datas; /* payload pointer for each message */
uint64_t *keys; /* timestamp key for each message */
uint32_t *lens; /* payload length for each message */
uint32_t len; /* number of messages in the batch */
};
/* Output of encode: columnar structure for Parquet persistence */
struct stream_data_out {
uint32_t col_len; /* number of columns */
uint32_t row_len; /* number of rows */
uint64_t *ts; /* timestamp column */
char **schema; /* column names */
parquet_data_packet ***payload_arr; /* per-column payload */
};
/* Output of decode: continuous buffer returned to the caller */
struct stream_decoded_data {
void *data; /* decoded buffer */
uint32_t len; /* buffer length */
};
/* Output of cmd_parser: structured query parameters */
struct cmd_data {
bool is_sync; /* true = sync, false = async */
uint64_t start_key; /* query start timestamp (ms) */
uint64_t end_key; /* query end timestamp (ms) */
uint32_t schema_len; /* number of requested columns */
char **schema; /* column names to return */
};插件实现和注册
实现三个接口函数,并在启动时使用 stream_register 注册:
void *my_encode(void *data); /* data: stream_data_in* -> return: parquet_data* */
void *my_decode(void *data); /* data: parquet_data_ret* -> return: stream_decoded_data* */
void *my_cmd_parser(void *data); /* data: const char* -> return: cmd_data* */
int my_stream_init() {
int ret = 0;
char *name = malloc(strlen("my_stream") + 1);
if (name == NULL) {
return -1;
}
strcpy(name, "my_stream");
/* Use an ID that does not conflict with built-ins: 0 = RAW, 0x1 = SPI */
ret = stream_register(name, 0x2, my_decode, my_encode, my_cmd_parser);
if (ret != 0) {
free(name);
return -1;
}
return 0;
}注册后,在对应 Exchange 配置中设置 streamType = 2,即可为该数据路径启用该插件。
各接口调用位置
encode:在webhook_post.c::send_exchange_cb/hook_last_flush中调用,调用时机是RB_FULL_RETURN触发后,或显式调用hook_last_flush刷新剩余数据时。它接收一个stream_data_in批次,并必须返回可用于 Parquet 持久化的编码结构。cmd_parser:在exchange_server.c::query_cb中调用,调用时机是 NNG pair0 socket 收到查询命令时。它接收原始命令字符串,并必须返回cmd_data。decode:在query_send_sync/query_send_async中针对每个从 Parquet 检索到的parquet_data_ret*调用。它必须返回带有连续结果缓冲区的stream_decoded_data。
查询通道和协议
EMQX Edge 暴露一个本地 NNG 通道,用于查询 Exchange 数据。源码仓库中包含一个参考客户端:nng/demo/exchange_consumer/exchange_consumer.c。
该客户端会:
- 通过
nng_pair0连接到 Exchange Server,地址由exchange_url指定,默认值为tcp://127.0.0.1:10000 - 将单个命令字符串作为请求 payload 发送
- 持续读取响应帧,直到收到 2 字节 EOF 标记(
0x0B 0xAD,参见exchange_server.c::query_send_eof)
EOF 之前的每一帧都是一批解码后的查询结果。
RAW Stream 数据布局
使用 RAW Stream(streamType = 0)时,数据以 key/value 形式存储:
- key:64 位单调递增的毫秒级时间戳,在捕获时通过
nng_msg_set_timestamp设置 - value:原始 MQTT Payload,存储在 Parquet 文件的
data列中
查询命令格式
RAW Stream 使用以下命令格式:
sync-<start_key>-<end_key>
async-<start_key>-<end_key>sync:在连接关闭前返回 key 范围内的所有结果async:Exchange 拆分并调度结果增量返回,适合大时间范围查询start_key/end_key:毫秒级时间戳;end_key必须大于等于start_key
示例:查询 5 秒窗口内的记录:
$ ./demo/exchange_consumer/exchange_consumer "sync-1700000000000-1700000005000"
Received 1234 bytes
Received 5678 bytes
...raw_cmd_parser 会将其解析为 is_sync = true、start_key = 1700000000000、end_key = 1700000005000。Exchange 会从 Parquet 和 Ringbus 中读取所有匹配记录,解码后返回响应帧,直到 EOF。
要检索所有持久化数据:
$ ./demo/exchange_consumer/exchange_consumer "sync-0-9223372036854775807"自定义插件可以通过实现自己的 cmd_parser 来扩展该机制,以支持复合条件、自定义 schema 和分页。