# 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 中数据语义的唯一来源。核心管道不假定固定的存储格式或查询协议；这些能力都委托给插件实现。

每个插件实现三个接口：

| 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()`，该函数会：

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 和分页。
