Skip to content

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 返回一个批次后,管道会执行以下步骤:

  1. 将批次规范化为 stream_data_in 结构,其中每条消息包含 payload、长度和时间戳 key。
  2. 调用 stream_encode(streamType, ...) 分发到配置的 Stream 插件。
  3. 插件将批次编码为可存储的表示形式,例如列式 Parquet 结构。
  4. 编码结果会异步写入本地 Parquet 文件。

使用 Parquet 作为存储格式时,面向批次的写入可最大化列式压缩效率,减少文件碎片,并支持快速按时间范围扫描。

查询路径 ​

历史查询通过 Exchange 的本地查询接口处理,不会经过 MQTT 转发路径。一次查询会同时覆盖仍在内存中(Ringbus)的数据和已经持久化到磁盘的数据。Exchange 合并这些来源的数据,并将其交给 Stream 插件解码,然后再向调用方返回结果。

Stream 插件架构 ​

Stream 插件是 MQTT Stream 中数据语义的唯一来源。核心管道不假定固定的存储格式或查询协议;这些能力都委托给插件实现。

每个插件实现三个接口:

InterfaceInvocation pointPurpose
encodeRingbus 满后、写入前将一批消息编码为可持久化结构
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(),该函数会:

  1. 根据返回的消息构建 stream_data_in,并使用 nni_msg_get_timestamp() 作为每条记录的 key。
  2. 调用 stream_encode(streamType, sdata),分发到插件的 encode 实现。对于 RAW,raw_encode 会将 stream_data_in 转换为 stream_data_out,并调用 parquet_data_alloc 构建列式结构。
  3. 编码结果(parquet_data*)会传递给 parquet_object_alloc + parquet_write_batch_async,用于异步写入磁盘。

解码流程 ​

当查询命令到达 Exchange 时:

  1. 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)。
  2. Exchange 通过 exchange_client_get_msgs_fuzz 从 Ringbus 内存中检索匹配记录,并通过 parquet_get_data_packets_in_range_by_column 从 Parquet 文件中检索匹配记录。
  3. 对每个 parquet_data_ret*,调用 stream_decode(streamType, ...) 将列式数据转换为 stream_decoded_data,即适合调用方使用的连续缓冲区。
  4. 结果会组装为 NNG 消息,并以增量方式(async)或聚合方式(sync)返回。

自定义 Stream 插件开发 ​

自定义 Stream 插件无需修改核心管道即可实现。插件需要实现三个函数,并在启动时注册。

接口结构体 ​

c
/* 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 注册:

c
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 使用以下命令格式:

text
sync-<start_key>-<end_key>
async-<start_key>-<end_key>
  • sync:在连接关闭前返回 key 范围内的所有结果
  • async:Exchange 拆分并调度结果增量返回,适合大时间范围查询
  • start_key / end_key:毫秒级时间戳;end_key 必须大于等于 start_key

示例:查询 5 秒窗口内的记录:

bash
$ ./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。

要检索所有持久化数据:

bash
$ ./demo/exchange_consumer/exchange_consumer "sync-0-9223372036854775807"

自定义插件可以通过实现自己的 cmd_parser 来扩展该机制,以支持复合条件、自定义 schema 和分页。