一、系统定位
当前项目不是业务库写入服务。它的核心职责是把 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_metric:metric_key=daily_mileage_km、calculation_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.v1、vehicle.event.gb32960.v1。 |
生产接入32960 |
| modules/apps/jt808-ingest-app | 生产 JT808 TCP 接入应用,只负责 808 注册/鉴权/位置等上行解析、冷存引用和 Kafka 投递。 | JT/T 808 TCP 808。 | vehicle.raw.jt808.v1、vehicle.event.jt808.v1。 |
生产接入808 |
| modules/apps/yutong-mqtt-app | 生产宇通 MQTT 接入应用,按 endpoint profile 解析 MQTT 报文并投递 Kafka。 | 宇通 MQTT topic。 | vehicle.raw.mqtt-yutong.v1、vehicle.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。 |
每日里程指标表 |
六、氢能业务开发落点
| 业务主题 | 当前建议落点 | 说明 |
|---|---|---|
| 车辆状态、速度、位置 | RealtimePayload、LocationPayload |
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、车企私有协议时不会重写统计口径。