Go 版本车辆数据处理与流转全图

当前 Go 架构以 gateway 做高性能协议接入,NATS JetStream 做入口缓冲和解耦,nats-kafka-bridge 保持现有 Kafka 下游兼容,最终进入 TDengine、Redis、MySQL,并由 realtime-api 对外提供查询能力。
外部数据源 Go 接入与解析 NATS / Kafka 解耦 存储 查询 API 可靠性边界

1. 总览图:从车辆平台到 API

外部输入
GB/T 32960 平台

TCP 连接到 ECS :32960,包含登录、登出、实时信息、补发等帧。

JT/T 808 平台

TCP 连接到 ECS :808,包含注册、鉴权、位置上报 0x0200 等。

宇通 MQTT

gateway 作为 MQTT client 订阅 /ytforward/shln/+

gateway 接入层
TCPServer / MQTTClient

读取连接、拆包、处理半包/粘包、按协议路由到解析器。

协议解析器

gb32960jt808yutongmqtt 解析 RAW,抽取核心 fields。

身份解析

通过 MySQL vehicle_identity_binding 将 phone/device/plate 尽量映射到 VIN。

协议响应

32960/808 需要应答时,在成功发布后写 ACK。

统一消息
FrameEnvelope

统一封装:protocol、message_id、vin、phone、device_id、source_endpoint、raw、parsed、fields。

StableEventID

基于协议、消息、车辆键、序号、时间、raw 生成稳定 ID,用于去重和追踪。

入口解耦
NATS JetStream

gateway 发布到 VEHICLE_INGEST stream,文件存储,24 小时保留。

Subjects

vehicle.raw.go.gb32960.v1
vehicle.raw.go.jt808.v1
vehicle.raw.go.yutong-mqtt.v1
vehicle.event.go.unified.v1

Async + Retry

gateway 内部异步队列发布 NATS,降低协议连接被 Kafka 慢写拖住的风险。

Kafka 兼容层
nats-kafka-bridge

durable pull consumer 批量拉 NATS,写 Kafka 成功后才 ACK NATS。

Kafka Topics

与 NATS subject 同名:raw 三个协议 + unified 一个实时主题。

失败语义

Kafka 写失败不 ACK,NATS 保留消息等待重试;未知 subject 不 ACK,防止静默丢数据。

落库与查询
history-writer

消费 RAW topic,写 TDengine:raw_frames、payload_chunks、locations、mileage_points。

stat-writer

消费 RAW topic,根据总里程差值写 MySQL 每日指标。

realtime-api 消费器

消费 unified topic,写 Redis 在线状态、实时合并快照、各协议 realtime-raw。

HTTP API

realtime-api 暴露实时、历史、统计查询接口,端口 :20200

2. 单条报文的完整生命周期

1. 到达32960/808 TCP 报文或 MQTT 消息到 gateway。
2. 拆包TCP 按协议提取完整帧;MQTT 直接按消息处理。
3. 解析生成 parsed 全量结构化数据和 fields 核心字段。
4. 绑定 VIN根据 VIN/phone/device_id/plate 生成 vehicle_key。
5. 发 NATS写 RAW subject;同时写 unified subject。
6. 桥接 Kafkabridge 写 Kafka 成功后 ACK NATS。
7. 多路消费历史、实时、统计各自消费 Kafka,互不阻塞。

3. FrameEnvelope:Go 版本内部统一数据格式

字段含义来源
event_id稳定事件 ID已有则保留,否则按协议/消息/车辆/时间/raw 计算
protocol协议GB32960JT808YUTONG_MQTT
message_id协议消息类型如 808 0x0200,32960 命令标识等
vinVIN协议原始字段或 MySQL binding 反查
phone808 终端手机号808 消息头 BCD 解析后去前导 0
device_id设备标识协议内设备号、MQTT 设备字段等
plate车牌注册帧或 binding 表
source_endpoint来源地址TCP remote ip:port 或 MQTT endpoint/topic
字段含义去向
raw_hex / raw_text原始报文TDengine raw_frames,超长进入 chunks
parsed全量解析 JSONTDengine parsed_json;Redis realtime-raw
fields最小核心字段TDengine 位置/里程点;Redis 合并快照;统计计算
event_time_ms车端事件时间TDengine 主时间;Redis field time
received_at_ms平台接收时间延迟分析、在线状态 last_seen
parse_status解析状态OKPARTIALBAD_FRAME
parse_error解析错误RAW 查询排查使用
vehicle_key车辆主键优先 VIN;否则 PROTOCOL:phone/device

4. 三个协议的处理差异

协议 接入方式 身份主线 RAW 数据 核心 fields 特殊动作
GB32960 gateway TCP :32960 报文本身通常有 VIN;平台账号用于登录鉴权 完整 parsed JSON 写 raw_frames.parsed_json 速度、经纬度、SOC、累计里程等按帧类型抽取 登录/数据帧需要平台应答,成功发布后 ACK
JT808 gateway TCP :808 消息头 phone 为主;注册/鉴权/位置帧结合 binding 找 VIN 完整 parsed JSON 写 raw_frames.parsed_json 位置、速度、方向、状态、报警、GPS 总里程 注册/鉴权/通用应答;每日里程按总里程差值统计
YUTONG_MQTT gateway MQTT client 订阅 MQTT payload 中 VIN/设备字段;source endpoint 为 MQTT endpoint/topic 完整 parsed JSON 写 raw_frames.parsed_json 速度、经纬度、总里程、SOC 等 不需要 TCP ACK;实时多帧在 Redis 按协议合并

5. NATS + Kafka 的职责分工

NATS JetStream:入口缓冲层

部署在 Kafka ECS,Docker Compose 管理;监听内网 172.17.111.56:4222,监控 172.17.111.56:8222。gateway 只需要把事件快速、可靠地写入 NATS。

nats-kafka-bridge:可靠桥接层

使用 durable pull consumer;批量写 Kafka;只有 Kafka 写入成功才 ACK NATS。Kafka 短暂异常时,NATS 消息保留并重试。

Kafka:业务消费层

保持现有 topic 和消费者模型,history/realtime/stat 不直接依赖 gateway,也不直接依赖 NATS API。

当前阶段没有让下游直接消费 NATS,这是为了最小化切换风险:入口换成 NATS 后,下游仍沿用 Kafka。后续如果要进一步实时化,可以让实时缓存或控制命令直接使用 NATS subject。

6. 存储模型:哪些数据进入哪里

存储 表/Key 写入方 内容 用途
TDengine raw_frames history-writer 原始报文、完整 parsed JSON、fields JSON、parse status、source endpoint RAW 查询、追溯、排查解析问题
TDengine raw_frame_payload_chunks history-writer 超长 raw/parsed payload 分片 避免 BINARY 长度限制导致完整 JSON 丢失
TDengine vehicle_locations history-writer 经纬度、速度、方向、状态、报警、总里程等核心字段 高频历史位置分页查询
TDengine vehicle_mileage_points history-writer 累计里程点、速度、经纬度 里程曲线、区间差值分析
Redis vehicle:realtime:{vehicle_key} realtime-api Kafka consumer 跨协议合并后的实时字段快照 查 VIN 是否在线、查 VIN 实时数据
Redis vehicle:realtime:{vehicle_key}:{protocol} realtime-api Kafka consumer 单协议实时快照 区分 32960/808/MQTT 的实时状态
Redis vehicle:realtime-raw:{vehicle_key}:{protocol} realtime-api Kafka consumer 单协议最新 parsed 全量字段 实时 RAW 字段查看,不走 TDengine 历史扫描
Redis vehicle:online:{vehicle_key}vehicle:last_seen realtime-api Kafka consumer 在线状态、最后接收时间、协议列表 在线判断和近期车辆列表
MySQL vehicle_identity_binding 人工/导入维护,gateway 读取 phone/device/plate 到 VIN 的映射 808 等缺 VIN 协议反查正确 VIN
MySQL jt808_registration gateway 808 注册、鉴权、首次/最新上报、phone、device、plate、vin 解析结果 定位哪些手机号没有映射 VIN,排查注册鉴权
MySQL vehicle_daily_metric stat-writer 每日里程等指标,按总里程差值计算 统计查询

7. 查询入口

实时查询

接口类型数据源说明
VIN 是否在线Redis online keyTTL 内有数据即在线,默认 TTL 600 秒。
VIN 实时数据Redis merged snapshot跨协议合并 fields,并保留各协议 ProtocolData。
单协议实时 RAWRedis realtime-raw查看某 VIN/phone 在某协议下最新 parsed 全量字段。

历史查询

接口类型数据源说明
RAW 帧查询TDengine raw_frames + chunks按协议、VIN/phone、时间、消息类型分页。
位置历史TDengine vehicle_locations高频位置分页查询,避免每次扫完整 JSON。
里程点TDengine vehicle_mileage_points区间里程、曲线、异常总里程排查。
每日指标MySQL vehicle_daily_metric按日期、协议、指标查询统计结果。

8. 生产部署视图

节点组件地址/端口说明
Gateway ECS lingniu-go-gateway 公网 115.29.187.205;生产端口 32960808 承接外部平台真实连接,发布到 NATS。
Gateway ECS lingniu-go-nats-kafka-bridge 连接 NATS 172.17.111.56:4222;Kafka 172.17.111.56:9092 当前部署在 gateway ECS,负责 NATS 到 Kafka 桥接。
Gateway ECS history-writerstat-writerrealtime-api API 115.29.187.205:20200 消费 Kafka,写 TDengine/Redis/MySQL,并提供查询。
Kafka ECS Kafka + NATS 公网 114.55.58.251;内网 172.17.111.56 NATS 使用 Docker Compose,端口只绑定内网 4222/8222
TDengine ECS TDengine 内网 172.17.111.57:6041 历史 RAW、位置、里程点时序存储。
云服务 MySQL / Redis RDS 内网、Redis 内网 身份绑定、808 注册、每日指标、实时缓存。

9. 关键可靠性设计

接入不直接等待 Kafka
gateway 先写 NATS,Kafka 慢写由 bridge 消化,降低 32960/808 连接积压风险。
NATS ACK 在 Kafka 成功之后
bridge 写 Kafka 失败时不 ACK,消息仍在 JetStream,恢复后继续拉取。
RAW 完整 JSON 保留
parsed 全量写 RAW,核心表只存查询常用字段,避免位置查询每次扫完整 JSON。
实时多协议合并
Redis 同时保留 merged snapshot、protocol snapshot、realtime-raw,既能看统一实时,也能追单协议原始字段。
身份解析降级
VIN 优先,其次 phone、vehicle_key_hint、device_id,保证无 VIN 时仍可临时归档和查询。
数据分层查询
高频查询走 TDengine 核心字段表和 Redis;完整追溯走 RAW JSON,性能和完整性分开处理。

10. 当前运行链路一句话

外部 32960/808/MQTT 数据进入 Go gateway 后,被解析成 FrameEnvelope,先写入 NATS JetStream; nats-kafka-bridge 把 NATS 消息可靠转写到 Kafka;Kafka 再分发给 history-writer 写 TDengine、 stat-writer 写 MySQL 每日指标、realtime-api 写 Redis 实时缓存并提供 HTTP 查询。