lingniu-vehicle-ingest 模块与数据流

本图基于当前 Maven 多模块项目整理,用于后续讨论协议接入、内部字段、氢能车辆运营统计、安全告警和下游消费边界。 当前架构将接入层、历史明细、Redis 热状态和日统计拆成独立模块,通过统一 Kafka Envelope 解耦。

一、系统定位

当前项目不是业务库写入服务。它的核心职责是把 GB/T 32960、JT/T 808、Yutong MQTT 等默认生产来源的车辆数据统一转换为 VehicleEvent 和全字段 TelemetrySnapshot; JT/T 1078、JSATL12 作为显式 profile 的可选能力保留, 再通过 Kafka、原始报文归档、TDengine 历史索引、Redis 热状态和独立统计模块交给下游业务。

接入方式 Netty TCP、MQTT、后续可扩展更多 Inbound Adapter;信达 Push 源码已删除。
统一模型 所有协议归一到 VehicleEvent 和全字段内部 TelemetrySnapshot
输出方式 Kafka Protobuf Envelope + 原始 bytes 冷存 + TDengine raw_frames/locations + Redis 热状态。
业务边界 日统计由 vehicle-stat-service 消费 Kafka 实现,告警工单、资产业务继续由下游消费。

二、模块分层图

flowchart TB
  subgraph external["外部来源"]
    vehicle["车辆终端
GB/T 32960 / JT808 / JT1078"] mqttSource["车企或平台 MQTT"] commandClient["业务系统 / 运维平台
HTTP 下行命令"] end subgraph protocolModules["协议模块 modules/protocols"] gb["protocol-gb32960
32960 Netty 接入、鉴权、ACK、解析"] jt808["protocol-jt808
808 Netty 接入、会话、位置/批量位置、媒体、透传、下行"] jt1078["protocol-jt1078
1078 信令 + TCP/UDP RTP 媒体流归档"] jsatl12["protocol-jsatl12
主动安全附件流接入"] end subgraph inbound["入口适配 modules/inbound"] mqtt["inbound-mqtt
MQTT endpoint 生命周期、PEM TLS"] mqttProfile["MqttProfileRegistry
按 endpoint profile 解析/映射"] end subgraph protocolBase["核心与公共能力 modules/core"] codec["ingest-codec-common
BCD / CRC / BCC / bit 工具"] session["session-core
设备会话、命令分发接口、会话存储"] identity["vehicle-identity
跨协议身份解析、MySQL 绑定表"] obs["observability
metrics / health"] api["ingest-api
ProtocolId / RawFrame / VehicleEvent / 注解 SPI"] registry["HandlerRegistry
按协议、命令、infoType 路由"] dispatcher["Dispatcher
RawFrame 到 Handler 到 EventBus"] bus["DisruptorEventBus
高吞吐事件发布"] end subgraph sink["输出与明细存储 modules/sinks"] kafka["sink-kafka
VehicleEvent 到 Protobuf Envelope 到 Kafka"] archive["sink-archive
RawArchive 到本地文件系统冷存"] tdengineStore["tdengine-history-store
TDengine raw_frames + locations"] end subgraph services["消费服务 modules/services"] history["event-history-service
Kafka event/raw 到 TDengine 历史查询"] state["vehicle-state-service
可选 Redis 热状态
optional-latest-state"] stat["vehicle-stat-service
Kafka 全字段事件到指标表"] end subgraph apps["应用入口 modules/apps"] gbApp["gb32960-ingest-app
GB32960 TCP 接入"] jtApp["jt808-ingest-app
JT808 TCP 接入"] yutongApp["yutong-mqtt-app
宇通 MQTT 接入"] historyApp["vehicle-history-app
历史查询与 TDengine 写入"] analyticsApp["vehicle-analytics-app
808 日里程指标消费"] gateway["command-gateway
可选 HTTP 设备命令入口"] end vehicle --> gbApp vehicle --> jtApp vehicle --> jt1078 vehicle --> jsatl12 mqttSource --> yutongApp gbApp --> gb jtApp --> jt808 yutongApp --> mqtt mqtt --> mqttProfile mqttProfile --> api commandClient --> gateway gb --> codec jt808 --> codec jt1078 --> codec jsatl12 --> codec gb --> api jt808 --> api jsatl12 --> archive jsatl12 --> jt808 api --> registry registry --> dispatcher dispatcher --> bus bus --> kafka bus --> archive kafka --> historyApp kafka --> analyticsApp historyApp --> history analyticsApp --> stat kafka -.可选热状态.-> state history --> tdengineStore gateway --> session session --> jt808 session --> gb identity --> jt808 identity --> mqttProfile obs -.监控.-> dispatcher obs -.监控.-> bus obs -.监控.-> kafka classDef entry fill:#eff6ff,stroke:#2563eb,color:#1e3a8a; classDef core fill:#ecfdf3,stroke:#0f766e,color:#134e48; classDef sink fill:#fffbeb,stroke:#b45309,color:#7c2d12; classDef command fill:#f5f3ff,stroke:#6d28d9,color:#4c1d95; classDef support fill:#f8fafc,stroke:#64748b,color:#334155; class gb,jt808,jt1078,jsatl12,mqtt,mqttProfile,gbApp,jtApp,yutongApp,historyApp,analyticsApp entry; class api,registry,dispatcher,bus,identity core; class kafka,archive,tdengineStore sink; class gateway,session command; class codec,obs support;
入口/协议模块 核心管线 输出层 下行命令/会话

三、上行数据流转

sequenceDiagram
  autonumber
  participant Device as 车辆/平台
  participant Inbound as Inbound Adapter
Netty/MQTT participant Decoder as FrameDecoder
MessageDecoder participant Dispatcher as Dispatcher participant Handler as Protocol Handler
Mapper participant EventBus as DisruptorEventBus participant Kafka as sink-kafka
Kafka Envelope participant Archive as sink-archive
ArchiveStore participant History as vehicle-history-app participant TDengine as tdengine-history-store
raw_frames + locations participant Consumer as 下游业务消费者 Device->>Inbound: 原始报文 bytes / MQTT message Inbound->>Decoder: 粘包拆帧、校验、协议体解析 Inbound->>Inbound: MQTT 按 endpoint profile 选择厂商解析器 Decoder->>Dispatcher: RawFrame(protocol, command, infoType, payload, rawBytes) Dispatcher->>EventBus: 先发布 RawArchive 事件 EventBus->>Archive: 保存原始 bytes,生成可回放材料 Dispatcher->>Handler: 根据 protocol + command + infoType 路由 Handler->>Handler: 协议字段映射为内部字段 Dispatcher->>Handler: 绑定 rawArchiveKey / rawArchiveUri Handler->>EventBus: VehicleEvent(Realtime / Location / Alarm / Login ...) EventBus->>Kafka: 投递规整后的业务事件 Kafka->>Kafka: EnvelopeMapper 转 Protobuf Kafka->>History: event/raw topic,key=vin,单车有序 History->>TDengine: 写 raw_frames、locations 和 raw parsed JSON TDengine->>Consumer: 分页历史查询、原始帧追溯、位置查询 Kafka->>Consumer: Redis 热状态、指标统计、告警工单等下游消费
关键边界: 协议字段只在协议模块内解释,出了 Mapper 以后都应该使用内部字段。统计每日里程、每日用电量、每日用氢量、储氢安全和氢泄露告警时,应消费全字段 TelemetrySnapshot,不要直接依赖某个协议的原始字段名。
原始追溯: Dispatcher 为所有带 rawBytes 的 RawFrame 生成统一的 rawArchiveKey,业务事件 metadata 使用 rawArchiveUri=archive://... 指向同一份原始 bytes。归档文件的真实本地路径只属于 ArchiveStore,查询、导出和 Kafka Envelope 使用逻辑 URI 解耦存储实现。

四、32960 细化链路

flowchart LR
  terminal["氢能车 TBOX
GB/T 32960"] --> netty["Gb32960NettyServer
TCP / TLS / Idle 检测"] netty --> access["Gb32960AccessService
VIN 白名单 / 平台登录鉴权"] netty --> frame["Gb32960FrameDecoder
拆帧、转义、校验"] frame --> decoder["Gb32960MessageDecoder
Header + Body 解析"] decoder --> diag["Gb32960FrameDiagnostics
首帧、Raw 块、异常诊断"] decoder --> handler["Gb32960ChannelHandler
ACK / NACK / dispatch"] handler --> ack["Gb32960AckService
登录应答、失败断开"] handler --> dispatcher["Dispatcher"] dispatcher --> rtHandler["Gb32960RealtimeHandler"] rtHandler --> mapper["Gb32960EventMapper
协议字段到内部字段"] mapper --> realtime["RealtimePayload
速度、里程、电量、氢量、压力、温度"] mapper --> alarm["AlarmPayload
安全分类、氢泄露、告警等级"] mapper --> login["Login / Logout / Heartbeat"] realtime --> bus["DisruptorEventBus"] alarm --> bus login --> bus bus --> kafka["Kafka Envelope"] bus --> archive["RawArchive 冷存"] classDef safety fill:#fff1f3,stroke:#b42318,color:#7a271a; classDef core fill:#ecfdf3,stroke:#0f766e,color:#134e48; classDef entry fill:#eff6ff,stroke:#2563eb,color:#1e3a8a; classDef sink fill:#fffbeb,stroke:#b45309,color:#7c2d12; class terminal,netty,frame,decoder,handler entry; class access,diag,ack,dispatcher,rtHandler,mapper,bus core; class alarm safety; class kafka,archive sink;

五、模块职责表

模块 职责 主要输入 主要输出 业务开发关注点
modules/core/ingest-api 定义系统边界:协议 ID、RawFrame、IngestContext、VehicleEvent、Payload、Handler 注解、Sink SPI;ProtocolId 提供 UNKNOWN 隔离桶,供消费侧兼容未来协议或脏数据,不作为正式接入协议使用;consumer 包定义 EnvelopeIngestor、EnvelopeConsumerProcessor 和 EnvelopeDeadLetterSink,统一 Kafka 消费结果与死信边界。 无直接外部输入。 全项目共享的接口和内部事件模型。 内部字段定义新增事件类型字段兼容性
modules/core/ingest-codec-common 公共编解码工具,承载 BCD、CRC、BCC、bit 操作等协议底层能力。 协议模块传入的 byte、bit、校验材料。 解析辅助结果。 协议基础能力
modules/core/ingest-core 核心管线。扫描 Handler,按协议和命令路由,执行拦截器,发布 VehicleEvent 到 Disruptor。 RawFrame。 VehicleEvent、RawArchive。 吞吐路由原始可回放
modules/core/session-core 设备会话、命令下发抽象、SessionStore、CommandDispatcher 默认实现。SessionStore 仅保留 Redis 后端;协议进程本地只保存真实 Netty Channel,Redis 保存 sessionId、VIN 和 phone 三索引及 TTL,协议模块只依赖 SPI,不感知存储细节。 设备连接、命令请求。 会话状态、命令分发结果、Redis 会话索引。 在线状态下行控制多实例会话索引
modules/core/vehicle-identity 跨协议车辆身份解析,维护 phone、deviceId、plate 到 VIN 的绑定关系;生产运行使用 MySQL 绑定表,避免重启丢失绑定并保证多实例一致。 协议模块传入的外部标识。 稳定 VIN、解析来源、是否已命中绑定、MySQL 绑定记录。 统计主键多协议归一持久化绑定
modules/core/observability 监控、指标、健康检查等横切能力。 核心管线和 Sink 运行状态。 Micrometer 指标、健康信息。 可用性延迟监控
modules/protocols/protocol-gb32960 GB/T 32960 接入、鉴权、ACK、报文解析、实时数据和告警映射。 32960 TCP 报文。 Realtime、Location、Alarm、Login、Logout、Heartbeat、RawFrame。 氢泄露储罐安全电量/氢量里程
modules/protocols/protocol-jt808 JT/T 808 接入、注册/鉴权/心跳/注销、位置/位置附加项、批量位置、参数/属性、多媒体、透传、未知上行兜底透传、分帧边界异常和协议解析异常坏帧兜底、下行命令和超长下行分包能力;0x0200 位置附加项 0x01 总里程按 0.1km 解码为内部字段 total_mileage_km,后续每日里程统计不依赖 808 原始字段;0x0102 鉴权帧解析会保留完整 token,并尽量提取 IMEI 和软件版本;注册和鉴权会话均使用共享身份解析后的内部 VIN,避免 deviceId/IMEI 污染后续会话和下行边界;终端连接断开时同步解绑 Channel 并清理 SessionStore 会话,避免 command-gateway 误判离线车辆仍可下行;正常上行 RawArchive 和业务事件 metadata 均使用共享身份解析后的内部 vin,并保留 phone、identityResolved、identitySource 便于冷存和业务事件按同一车辆查询;无法解析终端身份的坏帧 RawArchive 和 Passthrough 均按 vin=unknown、identityResolved=false、identitySource=UNKNOWN 标记;0x0900 透传事件按原始消息 ID 归类,透传类型写入 passthroughType metadata。 808 TCP 报文。 Location、Login、Logout、Heartbeat、MediaMeta、Passthrough、未知上行 Raw 兜底、坏帧 Passthrough、会话状态、下行响应。 GPS 位置在线状态媒体证据命令链路
modules/protocols/protocol-jt1078 JT/T 1078 音视频能力。信令在 JT808 mapper 存在时桥接复用 JT808 连接和包头;0x1005 乘客流量保留原始 body,同时结构化 channelId、startTime、endTime、passengerGetOn、passengerGetOff metadata 便于查询/导出;0x1205 文件列表保留原始 body,同时结构化 responseSerialNo、fileCount 摘要;下行信令编码由本模块提供,覆盖 0x1003 音视频属性查询、0x9101 实时预览、0x9102 实时控制、0x9201 历史回放、0x9202 回放控制、0x9205 资源列表查询、0x9206 文件上传和 0x9207 文件上传控制,命令发送仍由 command-gateway 经 CommandDispatcher 复用 JT808 在线通道;TCP/UDP RTP 媒体流独立端口接入,生产默认端口按旧接收服务对齐为 11078,按 VIN/通道/时间分段写 ArchiveStore,Kafka 只发送 MediaMeta 引用,事件 metadata 暴露内部 vin、sim、channelId、dataType、packetType、segment、sequence、segmentSizeBytes、archiveKey 和 archiveRef,MediaMeta.sizeBytes 同步当前片段大小,MediaMeta 按 archiveKey 去重,避免同一 SIM 运行中从 fallback 身份切换到内部 VIN 时漏发新归档引用,便于文件明细库追踪、展示和导出;超长/短包/坏魔数等 RTP 入口异常和媒体归档失败都会兜底产出 Passthrough 错误事件,坏 RTP 统一走 Dispatcher 产出 RawArchive 和 Passthrough,统一标记 parseError 和 parseErrorMessage,并保留内部 vin、原始 bytes、peer、长度、原因、RawArchive 引用、segment 和 archiveKey 便于排障;坏 RTP 头部可读时会提取 SIM 并通过共享身份服务映射 VIN。 1078 信令报文、TCP/UDP RTP 媒体流。 MediaMeta、乘客流量 Passthrough、文件列表 Passthrough、坏 RTP/归档失败 Passthrough、归档媒体分段。 视频扩展信令事件化大流冷存
modules/protocols/protocol-jsatl12 苏标主动安全报警附件接入(仅 optional-attachments profile 显式构建,不属于默认生产面)。依赖 ArchiveStore 和 JT808 decoder;端口按旧接收服务对齐为 7612;附件 DataPacket 按旧服务真实布局 01cd + 50B文件名 + offset + length + data 分帧解析,不再把文件名前 4 字节误判为帧长;附件数据写 ArchiveStore 后通过 Dispatcher 产出 MediaMeta 引用事件,事件 metadata 保留 fileName、fileOffset、declaredChunkSizeBytes、archiveKey 和 archiveRef,便于附件明细查询、导出和补传诊断;内置附件状态机处理 T1210 文件清单、DataPacket 分块区间和 T1212 上传完成消息,数据分块会按 T1210 文件清单中的文件名继承手机号并通过共享身份解析器得到内部 VIN,入口 RawFrame、ArchivedChunk 和 MediaMeta 均保留 phone、vin、identityResolved、identitySource,只有找不到文件归属时才标记 UNKNOWN;按文件名合并无身份数据分块并返回 0x9212 完成/补传应答,补传应答包含缺失 offset/length 区间;附件归档失败会产出 JSATL12 Passthrough 错误事件并保留原始分块;JT_MESSAGE 信令复用 JT808 decoder 后进入 Dispatcher,桥接 RawFrame 使用共享身份解析后的内部 VIN并保留 phone、identityResolved、identitySource,坏 JT 信令入口 RawFrame 同样标记 UNKNOWN 身份;超长/畸形附件帧、未闭合超长 JT 信令和坏 JT 信令都以 JSATL12 Passthrough 兜底,附件帧异常同时标记 frameError 与统一 parseError/parseErrorMessage,保留原始字节、peer、长度和错误元数据,帧长和 worker 线程支持部署配置。 主动安全附件流、JT808 信令帧。 归档文件、MediaMeta、JT808 统一事件、坏帧/坏信令/归档失败 Passthrough。 附件归档报警证据信令复用错误隔离
modules/inbound/inbound-mqtt MQTT 多 endpoint 生命周期管理,支持 PEM CA、客户端证书/私钥双向 TLS,按 endpoint profile 路由到厂商解析/映射实现,默认提供宇通 profile;生产默认接入口径参考旧接收服务:ssl://cpxlm.axxc.cn:38883、topic /ytforward/shln/+、QoS 2、cleanSession=false、keepAlive=20s、connectTimeout=10s,用户名和密码仍通过环境变量注入;入站 RawFrame 和事件映射均接入统一车辆身份解析,解析器保留 externalVin、phone、mqttDeviceId、plateNo,优先按设备号/终端 ID/IMEI 绑定解析内部 VIN,未绑定且字段形态为 VIN 时再作为显式 VIN 兜底,冷存和事件 metadata 均保留这些外部标识、identityResolved、identitySource;未知 profile、解析失败、profile 运行时异常、endpoint 初始化/连接/订阅失败均兜底为 Passthrough,解析失败诊断和 profile 异常诊断都会同步写入 RawArchive metadata,并在最终 Passthrough 事件保留 parseErrorMessage 或 operational reason,无法解析设备身份的兜底事件和运行类归档都按 vin=unknown、identityResolved=false、identitySource=UNKNOWN 标记,保留原始 payload 或 operational 错误信息便于排障和归档追溯。 MQTT topic/message、endpoint profile、TLS PEM 配置。 统一 RawFrame、Realtime、Location、Passthrough。 车企平台接入profile 扩展身份映射双向 TLS错误隔离
modules/sinks/sink-kafka Kafka Sink。把 VehicleEvent 转成 Protobuf Envelope,按 VIN 分区投递;业务事件携带全字段 TelemetrySnapshot 和 rawArchiveUri,显式 RawArchive Envelope 会填充 archive:// 逻辑 URI 与 size,默认 Kafka Sink 仍不发送原始 bytes 本体;同时提供 KafkaEnvelopeConsumerFactory、KafkaEnvelopeConsumerRunner、KafkaEnvelopeConsumerWorker 和 KafkaEnvelopeDeadLetterSink,统一服务侧 Kafka 消费启动、独立 group 绑定、死信发布和 topic 到 EnvelopeConsumerProcessor 的分发。 VehicleEvent。 Kafka Protobuf 消息、消费侧 DLQ 记录。 单车有序Schema 演进安全字段下发给消费者
modules/sinks/sink-archive 原始报文冷存,当前为本地文件系统实现;写入键与业务事件 metadata 中的 rawArchiveKey 保持一致。 RawArchive 事件、附件流。 可回放原始文件、archive://... 逻辑引用。 问题排查合规留痕
modules/sinks/tdengine-history-store 生产历史库边界,负责 TDengine schema、raw_frames/location 行映射、批量写入和分页查询语句;raw_frames 保留完整 payloadJson.parsed,位置表保持轻量字段并通过 rawUri 关联原始帧。 Kafka Protobuf Envelope、Raw Envelope。 TDengine raw_frames 和位置表。 生产历史分页查询raw 追溯
modules/services/event-history-service 消费 32960、808、宇通 MQTT 的 Kafka event/raw topic,生产默认写入 TDengine raw_frames 和位置表,并提供 raw 帧、位置历史等分页查询边界;Raw 查询返回完整 parsed JSON,位置历史只保留核心字段并通过 rawUri 关联原始帧;生产消费可使用 tryIngest 和自动装配的 EnvelopeConsumerProcessor 将坏 protobuf、缺快照和存储异常收敛为结构化结果,并通过 EnvelopeDeadLetterSink 发布死信,避免阻塞 Kafka 分区。 Kafka Protobuf Envelope。 TDengine raw_frames、位置分页结果、raw archive 逻辑引用。 历史明细分页查询raw JSON
modules/services/vehicle-state-service 消费 Kafka 全字段事件,更新车辆最新状态、位置、安全和最后事件到 Redis;仅通过 optional-latest-state profile 显式构建,不属于默认生产 reactor;生产消费可使用 tryIngest 和自动装配的 EnvelopeConsumerProcessor 隔离坏 protobuf 和缺快照 Envelope,并把失败记录交给死信出口。 Kafka Protobuf Envelope。 vehicle:state:{vin}vehicle:location:{vin}vehicle:safety:{vin} 毫秒级热查询氢泄露安全
modules/services/vehicle-stat-service 消费 Kafka 全字段事件,808 每日里程只使用 0x0200 附加项 0x01 上报的总里程做当日首末差值;生产消费可使用 tryIngest 和自动装配的 EnvelopeConsumerProcessor 隔离坏 protobuf,缺少 telemetry_snapshot 的消息会标记为 SKIPPED 并进入死信出口。 Kafka Protobuf Envelope、808 total_mileage_km vehicle_stat_metricmetric_key=daily_mileage_kmcalculation_method=JT808_TOTAL_MILEAGE_DIFF 每日里程指标表
modules/apps/command-gateway 可选 HTTP 到设备下行命令入口。支持会话查询、808 位置查询、参数查询/设置、终端控制、平台通用应答、终端属性查询 0x8107、区域删除 0x8601、人工报警确认 0x8203,以及 1078 音视频属性查询、实时预览/控制、历史回放/控制、资源列表查询、文件上传/控制;通过 CommandDispatcher 复用 protocol-jt808 的在线通道、流水号和同步应答等待能力。 业务系统命令请求。 CommandDispatcher 调用、设备应答摘要。 远程控制参数设置属性查询区域删除报警确认视频控制
modules/apps/gb32960-ingest-app 生产 GB32960 TCP 接入应用,只负责接收、鉴权、解析、冷存引用和 Kafka 投递。 GB/T 32960 TCP 32960。 vehicle.raw.gb32960.v1vehicle.event.gb32960.v1 生产接入32960
modules/apps/jt808-ingest-app 生产 JT808 TCP 接入应用,只负责 808 注册/鉴权/位置等上行解析、冷存引用和 Kafka 投递。 JT/T 808 TCP 808。 vehicle.raw.jt808.v1vehicle.event.jt808.v1 生产接入808
modules/apps/yutong-mqtt-app 生产宇通 MQTT 接入应用,按 endpoint profile 解析 MQTT 报文并投递 Kafka。 宇通 MQTT topic。 vehicle.raw.mqtt-yutong.v1vehicle.event.mqtt-yutong.v1 生产接入MQTT
modules/apps/vehicle-history-app 生产历史应用,消费 32960、808、宇通 MQTT 的 raw/event Kafka Envelope,写 TDengine 并提供历史查询。 Kafka Protobuf Envelope。 TDengine raw_frames、位置表和历史查询 API。 历史查询raw JSON
modules/apps/vehicle-analytics-app 生产指标应用,当前只消费 JT808 事件并用 GPS 总里程首末差值写每日里程指标。 vehicle.event.jt808.v1 MySQL vehicle_stat_metric 每日里程指标表

六、氢能业务开发落点

业务主题 当前建议落点 说明
车辆状态、速度、位置 RealtimePayloadLocationPayload 32960 和 808 都可能提供位置。统计侧需要明确优先级,例如优先 32960,缺失时用 808 补点。
每日里程 Kafka 下游统计服务 接入层只投递总里程和实时点位;日增里程建议下游按 VIN、自然日、事件时间聚合,处理回补和乱序。
每日用电量 内部电池字段 + 下游统计 接入层保留 SOC、电压、电流等瞬时字段;若有累计电耗字段,可加入内部字段模型后统一投递。
每日用氢量 hydrogenRemainingKg + 下游统计 优先使用累计氢耗字段;没有累计值时用氢余量差值估算,并在统计结果里标记估算口径。
储氢安全 RealtimePayload 压力/温度 + AlarmPayload.safetyCategory 高压、低压、温度、异常 bit 都应统一归入储罐安全域,下游可做趋势、阈值和连续异常判断。
氢气泄露 AlarmPayload.hydrogenLeakDetected 当前已作为高优先级安全事件。32960 中出现 HYDROGEN_LEAK 时强制映射为 CRITICAL
原始报文追溯 RawArchive + rawArchiveUri + sink-archive 业务统计出现争议时,用事件 metadata 中的 archive://... 逻辑 URI 回查原始 bytes,结合 eventId、traceId、VIN、时间复现解析链路。
明细展示和导出 vehicle-history-app + tdengine-history-store 生产查询通过 TDengine raw_frames、位置表和 rawUri 关联完成;不再维护文件型事件索引旁路。
后续开发建议: 当前 Kafka Protobuf 已加入全字段 TelemetrySnapshot。后续业务统计服务继续只依赖内部字段,不依赖 32960 原始字段名,这样接入 808、MQTT、车企私有协议时不会重写统计口径。