diff --git a/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md new file mode 100644 index 00000000..24598520 --- /dev/null +++ b/docs/superpowers/specs/2026-06-29-vehicle-ingest-redesign.zh.md @@ -0,0 +1,616 @@ +# 车辆接入重设计方案 + +## 背景 + +当前 GB32960 和 JT808 链路已经能够完成接收、解析、发布和查询,但职责边界被拆散在多条路径里: + +- Netty Handler 生成 `RawFrame`,Dispatcher 同时发出原始归档事件和标准化车辆事件。 +- Kafka Envelope 只携带 raw archive 引用,原始 bytes 由 archive sink 另行存储。 +- GB32960 历史查询依赖 raw archive 记录,再回读文件重新解析。 +- JT808 位置历史又单独物化到 TDengine。 +- DuckDB/Parquet、本地 raw archive、Kafka、TDengine 分别承担一部分职责,导致查询语义和持久化边界不好推理。 + +新设计保留已有协议解析器和部署经验,但重新定义接入、持久化、存储和查询边界。 + +## 目标 + +1. 保存每一条入站协议原始帧,包括可拿到 bytes 的异常帧。 +2. 支持按车辆、协议、消息类型、时间范围查询 GB32960 和 JT808 历史上报数据。 +3. 新增字段、协议消息、厂商扩展时,不需要重写核心管线。 +4. 支撑 10,000 辆车生产写入,以及 1,000 QPS 历史查询。 +5. 明确 ACK 和持久化规则,避免“已经 ACK 但数据悄悄丢失”。 + +## 非目标 + +- 不把 Kafka 作为长期历史查询数据库。 +- 不把大块原始 bytes 直接塞进 Kafka topic。 +- 不继续把 DuckDB/Parquet 作为 GB32960/JT808 的生产热查询库。 +- 不做一个掩盖协议差异的万能接口。内部可以统一,外部 API 保持协议语义清晰。 + +## 推荐架构 + +采用 `Raw Archive + Kafka + TDengine` 作为核心架构。 + +```mermaid +flowchart LR + GB["GB32960 TCP"] --> Edge["协议接入边界"] + JT["JT808 TCP"] --> Edge + Edge --> RawStore["Raw Archive"] + RawStore --> RawTopic["Kafka vehicle.raw-frame.v1"] + Edge --> Decode["协议解析器"] + Decode --> FactTopic["Kafka vehicle.decoded-fact.v1"] + RawTopic --> Writer["History Writer"] + FactTopic --> Writer + Writer --> TD["TDengine"] + Writer --> Redis["实时状态"] + API["历史 API"] --> TD + API --> RawStore + RealtimeAPI["实时 API"] --> Redis +``` + +架构里只有三类事实: + +- `RawFrameFact`:每一条入站帧一条记录,不管解析成功还是失败。 +- `DecodedFact`:从 raw frame 派生出来的零条或多条标准事实。 +- `RawArchiveObject`:不可变原始 bytes,通过 `raw_uri` 寻址。 + +raw archive 是排障和回放的事实源。TDengine 是查询事实源。Kafka 是解耦和重放管道。 + +## 模块边界 + +### 协议接入 APP + +APP: + +- `gb32960-ingest-app` +- `jt808-ingest-app` + +职责: + +- 接收 TCP 连接。 +- 完成帧切分。 +- 执行协议层校验、鉴权、会话绑定和 ACK。 +- 创建 `RawFrameFact`。 +- 通过 `RawArchiveWriter` 持久化原始 bytes。 +- 发布 raw frame fact 到 Kafka。 +- 解析协议帧并发布 decoded facts。 + +接入 APP 不负责历史查询、导出,也不直接写业务查询表。 + +### Ingest Core + +新增模块: + +- `modules/core/ingest-facts` + +核心类型: + +```java +public record RawFrameFact( + String frameId, + ProtocolId protocol, + String vehicleKey, + String vin, + String phone, + int messageId, + int subType, + Instant eventTime, + Instant receivedAt, + String peer, + String rawUri, + String checksum, + long rawSizeBytes, + ParseStatus parseStatus, + String parseError, + Map metadata +) {} +``` + +```java +public sealed interface DecodedFact permits + DecodedFact.Location, + DecodedFact.Realtime, + DecodedFact.Alarm, + DecodedFact.Session, + DecodedFact.Register, + DecodedFact.Extension { + + String factId(); + String frameId(); + ProtocolId protocol(); + String vehicleKey(); + String vin(); + Instant eventTime(); + Instant receivedAt(); + String rawUri(); + Map metadata(); +} +``` + +`vehicleKey` 是稳定分区键: + +- GB32960:VIN。 +- JT808:已解析到 VIN 时使用 VIN,否则使用 `jt808:`。 +- 未知或异常帧:使用 `unknown::`。 + +### Raw Archive + +新增或重构模块: + +- `modules/sinks/raw-archive-store` + +接口: + +```java +public interface RawArchiveWriter { + RawArchiveReceipt write(RawArchiveWriteRequest request) throws IOException; +} +``` + +写入回执: + +```java +public record RawArchiveReceipt( + String rawUri, + String checksum, + long sizeBytes +) {} +``` + +规则: + +- 原始 bytes 一旦写入,不允许原地修改。 +- 写入路径可以按内容寻址,也可以按 `protocol/date/vehicleKey/frameId.bin` 生成。 +- 本地存储需要支持 fsync 能力,接口要预留对象存储实现。 +- `rawUri` 必须能被历史 API 解析,用于单帧回放。 +- 异常帧只要有 bytes,也必须写入,并记录 `parseStatus=FAILED`。 + +### Kafka Topic + +使用稳定的版本化 topic: + +- `vehicle.raw-frame.v1` +- `vehicle.decoded-fact.v1` +- `vehicle.dlq.v1` + +Kafka key: + +```text +: +``` + +`vehicle.raw-frame.v1` 的 payload 是不含 raw bytes 的 `RawFrameFact`,携带 `rawUri`、checksum、byte size、message id、parse status 和 metadata。 + +`vehicle.decoded-fact.v1` 的 payload 是单条 `DecodedFact`,必须引用 `frameId` 和 `rawUri`。 + +`vehicle.dlq.v1` 保存存储失败、发布失败、解析失败和写入失败事件,并携带足够元数据,尽量支持从 raw archive 重放。 + +### History Writer + +新增或重构 APP: + +- `vehicle-history-app` + +职责: + +- 消费 `vehicle.raw-frame.v1`。 +- 消费 `vehicle.decoded-fact.v1`。 +- 批量写 TDengine。 +- 更新实时状态存储。 +- 对非法 payload 或存储失败写 DLQ。 + +Writer 写入需要幂等: + +- raw fact 以 `(frame_id)` 幂等。 +- decoded fact 以 `(fact_id)` 幂等。 + +## TDengine 模型 + +TDengine 是生产历史查询主库。 + +### 原始帧超级表 + +```sql +CREATE STABLE IF NOT EXISTS raw_frames ( + ts TIMESTAMP, + frame_id NCHAR(64), + received_at TIMESTAMP, + message_id INT, + sub_type INT, + event_time TIMESTAMP, + raw_uri NCHAR(512), + checksum NCHAR(128), + raw_size_bytes BIGINT, + parse_status NCHAR(16), + parse_error NCHAR(512), + peer NCHAR(128), + metadata_json NCHAR(4096) +) TAGS ( + protocol NCHAR(16), + vehicle_key NCHAR(128), + vin NCHAR(64), + phone NCHAR(32) +); +``` + +子表命名: + +```text +raw__ +``` + +主要查询模式: + +- 按车辆和时间查询原始帧。 +- 按协议、消息 ID、解析状态和时间查询原始帧。 +- 按 `frame_id` 或 `raw_uri` 查询单帧。 + +### 位置超级表 + +```sql +CREATE STABLE IF NOT EXISTS vehicle_locations ( + ts TIMESTAMP, + fact_id NCHAR(64), + frame_id NCHAR(64), + received_at TIMESTAMP, + longitude DOUBLE, + latitude DOUBLE, + altitude_m DOUBLE, + speed_kmh DOUBLE, + direction_deg DOUBLE, + alarm_flag BIGINT, + status_flag BIGINT, + total_mileage_km DOUBLE, + raw_uri NCHAR(512), + metadata_json NCHAR(4096) +) TAGS ( + protocol NCHAR(16), + vehicle_key NCHAR(128), + vin NCHAR(64), + phone NCHAR(32) +); +``` + +### 实时数据超级表 + +```sql +CREATE STABLE IF NOT EXISTS vehicle_realtime ( + ts TIMESTAMP, + fact_id NCHAR(64), + frame_id NCHAR(64), + received_at TIMESTAMP, + speed_kmh DOUBLE, + total_mileage_km DOUBLE, + battery_soc DOUBLE, + total_voltage_v DOUBLE, + total_current_a DOUBLE, + running_mode NCHAR(64), + vehicle_state NCHAR(64), + charging_state NCHAR(64), + raw_uri NCHAR(512), + fields_json NCHAR(8192), + metadata_json NCHAR(4096) +) TAGS ( + protocol NCHAR(16), + vehicle_key NCHAR(128), + vin NCHAR(64), + phone NCHAR(32) +); +``` + +`fields_json` 存放稀疏全字段值,用于低频排查。高频查询字段应提升为明确列,或放入专用扩展表。 + +### 报警超级表 + +```sql +CREATE STABLE IF NOT EXISTS vehicle_alarms ( + ts TIMESTAMP, + fact_id NCHAR(64), + frame_id NCHAR(64), + received_at TIMESTAMP, + alarm_level NCHAR(32), + alarm_code BIGINT, + alarm_name NCHAR(128), + longitude DOUBLE, + latitude DOUBLE, + raw_uri NCHAR(512), + fields_json NCHAR(8192), + metadata_json NCHAR(4096) +) TAGS ( + protocol NCHAR(16), + vehicle_key NCHAR(128), + vin NCHAR(64), + phone NCHAR(32) +); +``` + +### 会话和注册表 + +```sql +CREATE STABLE IF NOT EXISTS vehicle_sessions ( + ts TIMESTAMP, + fact_id NCHAR(64), + frame_id NCHAR(64), + received_at TIMESTAMP, + session_type NCHAR(32), + result_code INT, + raw_uri NCHAR(512), + fields_json NCHAR(8192), + metadata_json NCHAR(4096) +) TAGS ( + protocol NCHAR(16), + vehicle_key NCHAR(128), + vin NCHAR(64), + phone NCHAR(32) +); +``` + +JT808 注册详情先进入 `vehicle_sessions`,后续可以投影到 MySQL 车辆身份绑定表,用于持久身份解析。 + +## ACK 和持久化规则 + +对需要协议 ACK 的帧: + +1. Netty 接收到 frame bytes。 +2. 原始 bytes 写入 raw archive。 +3. `RawFrameFact` 发布到 Kafka。 +4. 发送协议 ACK。 +5. 解码和 `DecodedFact` 发布可以在 ACK 前或 ACK 后执行,但失败必须发 DLQ,并更新 raw frame 解析状态。 + +对 GB32960 实时上报和补发帧,这条规则保证不会 ACK 一条无法恢复原始 bytes 的帧。 + +对 JT808 注册、鉴权、心跳,也遵循相同 raw archive 和 raw fact 边界。注册 ACK 可以先生成 token,但在 durability boundary 成功前不能 flush 到 channel。 + +如果 raw archive 写入失败,不返回成功 ACK。具体是关闭连接还是保留连接,由协议容忍度决定。 + +如果 Kafka raw fact 发布失败,durable ACK 帧不返回成功 ACK。非 durable 帧可以写本地 DLQ。 + +## 协议解析和扩展 + +使用插件契约: + +```java +public interface ProtocolFrameDecoder { + ProtocolId protocol(); + DecodeResult decode(RawFrameFact rawFact, byte[] rawBytes); +} +``` + +```java +public record DecodeResult( + ParseStatus status, + List facts, + String errorMessage, + Map metadata +) {} +``` + +GB32960: + +- 命令帧映射为会话、实时、位置、报警和扩展事实。 +- 厂商扩展按平台账号、VIN 或配置 profile 选择。 +- 全量 raw replay 仍通过读取 `rawUri` 并运行当前解析器完成。 + +JT808: + +- 消息 ID 映射为注册、鉴权、心跳、位置、批量位置、多媒体元数据和扩展事实。 +- 身份解析会更新 `vehicleKey`。 +- 注册事实携带省、市、厂商、设备 ID、车牌、车牌颜色和 auth token metadata。 + +新增协议消息通常不需要改 TDengine schema,除非它成为高频查询维度。未知字段先作为 extension fact 存入 `fields_json`。 + +## 查询 API + +API 保持明确,不做一个含糊的统一接口。 + +原始帧查询: + +- `GET /api/history/raw-frames` +- 过滤条件:`protocol`、`vin`、`phone`、`vehicleKey`、`messageId`、`parseStatus`、`dateFrom`、`dateTo`、`pageSize`、`cursor`。 +- 返回 frame metadata、raw URI、checksum、parse status,以及可选的紧凑解析摘要。 + +单帧详情: + +- `GET /api/history/raw-frames/{frameId}` +- 返回 raw metadata 和解码详情。 +- 支持 `includeRawHex=true`,但必须有严格大小限制。 + +位置历史: + +- `GET /api/history/locations` +- 过滤条件:`protocol`、`vin`、`phone`、`vehicleKey`、时间范围、cursor。 + +实时历史: + +- `GET /api/history/realtime` +- 查询历史时序数据。 + +实时状态: + +- `GET /api/realtime/vehicles/{vehicleKey}` +- 从 Redis 或内存 latest state 读取,不扫 TDengine。 + +导出: + +- 导出必须异步。 +- 请求创建 export job。 +- Worker 将 TDengine 查询结果流式写入对象存储或本地导出文件。 +- API 返回 job 状态和下载 URI。 + +## 分页 + +高 QPS 历史查询使用 cursor pagination,不使用深 offset。 + +Cursor 字段: + +```text +ts, received_at, fact_id +``` + +倒序查询谓词: + +```sql +WHERE ( + ts < :cursorTs + OR (ts = :cursorTs AND received_at < :cursorReceivedAt) + OR (ts = :cursorTs AND received_at = :cursorReceivedAt AND fact_id < :cursorFactId) +) +``` + +默认 page size 为 100,上限为 1000。 + +## 容量设计 + +假设: + +- 10,000 辆车。 +- 常规上报间隔 5 到 10 秒。 +- 预计写入 1,000 到 2,000 frames/s。 +- 峰值突发系数 3 倍。 +- 历史查询目标 1,000 requests/s。 + +写路径: + +- Netty worker 保持 CPU-light。 +- Raw archive 写入使用有界异步 worker。 +- Kafka producer 开启压缩、批量、幂等,并按 vehicle key 分区。 +- History writer 按超级表和子表批量写 TDengine。 +- TDengine 子表按 vehicle key 分区,使单车查询足够便宜。 + +读路径: + +- 高频查询必须带 vehicle key 和时间范围。 +- 无边界跨车查询仅限管理接口,并强制限流和限制返回量。 +- 实时 latest state 从 Redis 或内存状态读取。 +- Raw 详情只读取单个 raw object 并按需解析。 +- 字典和元数据接口使用内存缓存。 + +背压: + +- Raw archive worker queue 固定上限。 +- Kafka publish future 必须被跟踪。 +- Durability boundary 超时时,durable ACK 失败,channel 关闭或被限流。 +- History writer lag 通过 Kafka consumer lag 和 TDengine batch latency 观察。 + +## 迁移计划 + +阶段 1:在现有代码旁边建立新的 fact model。 + +- 新增 `ingest-facts`。 +- 新增 raw archive receipt contract。 +- 新增 raw frame 和 decoded fact 的 Kafka 序列化。 +- 为 stable id、vehicle key、raw URI 生成规则加测试。 + +阶段 2:强化 raw durability boundary。 + +- 重构 GB32960 和 JT808 handler,让 ACK 等待 raw archive 和 raw fact publish。 +- 保留现有解析器。 +- 向新 Kafka topic 发 decoded facts。 + +阶段 3:构建 TDengine writer。 + +- 创建 TDengine schema manager。 +- 写入 raw frame facts。 +- 写入 location、realtime、alarm、session facts。 +- 按 stable id 实现幂等写入。 + +阶段 4:构建查询 API。 + +- Raw frame 查询。 +- 单帧回放。 +- 位置历史查询。 +- 实时历史查询。 +- 最新实时状态查询。 + +阶段 5:切换并下线旧热路径。 + +- GB32960/JT808 生产热历史不再依赖 DuckDB/Parquet。 +- 历史迁移期间保留 raw replay 兼容。 +- TDengine 覆盖验证完成后,移除协议特定的临时历史存储。 + +## 测试策略 + +单元测试: + +- `RawFrameFact` 校验和 vehicle key 派生。 +- Raw archive 路径生成和 checksum。 +- GB32960 golden frame 到 decoded facts。 +- JT808 sample frame 到 decoded facts。 +- ACK boundary 成功和失败场景。 +- TDengine SQL 生成。 + +集成测试: + +- 使用 mocked Kafka producer 和 fake raw archive 启动 APP。 +- 发送 GB32960 sample frame,验证 raw fact、decoded fact 和 raw archive receipt。 +- 发送 JT808 注册、鉴权、位置帧,验证 registration、session、location facts。 +- 模拟 raw archive 失败,验证不发送成功 ACK。 +- 模拟 Kafka publish 失败,验证 durable ACK 失败。 + +存储测试: + +- TDengine schema 初始化。 +- 批量写入 raw frames 和 decoded facts。 +- 按 vehicle key 和时间范围查询。 +- Cursor 分页正确性。 +- 按 `rawUri` 单帧回放。 + +压测: + +- 持续 30 分钟 2,000 frames/s 写入。 +- 持续 5 分钟 6,000 frames/s 突发写入。 +- 1,000 QPS location/realtime 历史查询,p95 延迟受控。 +- 校验 raw archive receipt、Kafka raw fact、TDengine raw frame row 三者数量无缺口。 + +## 运维要求 + +指标: + +- TCP 连接数。 +- Frame 接收速率。 +- Raw archive 写入延迟和失败数。 +- Kafka 发布延迟和失败数。 +- 按协议和 message id 统计 decoder 成功/失败数。 +- TDengine batch size、延迟和失败数。 +- 各查询 endpoint 的 QPS 和延迟。 +- Kafka consumer lag。 + +日志: + +- 每个失败帧输出一条结构化日志,包含 `frameId`、protocol、vehicle key、message id、raw URI 和 error。 +- 正常日志不输出完整 raw bytes。 +- Raw hex 只允许在显式 debug 工具里输出,并限制大小。 + +告警: + +- Raw archive 写入失败。 +- Durable ACK 超时。 +- Kafka producer 失败。 +- History writer lag。 +- TDengine 写入失败。 +- Raw archive 数量和 TDengine raw frame row 数量不一致。 + +## 验收标准 + +重设计只有满足以下条件才算完成: + +1. 每一条被接收的 GB32960 和 JT808 帧都会创建 raw archive object 和 TDengine `raw_frames` 行。 +2. 有 bytes 的异常帧可以从 raw history 查询到,包含 parse status 和错误原因。 +3. GB32960 和 JT808 位置历史可以按车辆和时间范围 cursor 分页查询。 +4. GB32960 实时历史可以按车辆和时间范围查询。 +5. JT808 注册和会话数据可查询,并且可以更新 identity binding。 +6. 单帧详情可以通过 `rawUri` 重新读取 raw bytes,并使用当前解析器解码。 +7. 需要 ACK 的帧,在 raw archive 和 raw fact publish 成功前,不会发送成功 ACK。 +8. 新协议字段可以通过新增 decoder/projector 扩展;只有成为高频查询字段时才需要 TDengine 表迁移。 +9. 压测证明可以持续 2,000 frames/s 写入,并支持 1,000 QPS 有边界历史查询。 +10. GB32960/JT808 生产查询不再依赖旧 DuckDB/Parquet 热历史路径。 + +## 已固定决策 + +第一版实现固定以下选择: + +- 先使用本地 raw archive,实现接口时预留对象存储。 +- TDengine 是 GB32960/JT808 唯一生产热历史查询库。 +- 使用协议明确的查询 API,不做一个泛化万能查询 API。 +- 高频历史查询统一使用 cursor pagination。 +- Kafka payload 保持轻量,只传 raw 引用,不传 raw bytes。