1. 总览图:从车辆平台到 API
TCP 连接到 ECS :32960,包含登录、登出、实时信息、补发等帧。
TCP 连接到 ECS :808,包含注册、鉴权、位置上报 0x0200 等。
gateway 作为 MQTT client 订阅 /ytforward/shln/+。
读取连接、拆包、处理半包/粘包、按协议路由到解析器。
gb32960、jt808、yutongmqtt 解析 RAW,抽取核心 fields。
通过 MySQL vehicle_identity_binding 将 phone/device/plate 尽量映射到 VIN。
32960/808 需要应答时,在成功发布后写 ACK。
统一封装:protocol、message_id、vin、phone、device_id、source_endpoint、raw、parsed、fields。
基于协议、消息、车辆键、序号、时间、raw 生成稳定 ID,用于去重和追踪。
gateway 发布到 VEHICLE_INGEST stream,文件存储,24 小时保留。
vehicle.raw.go.gb32960.v1vehicle.raw.go.jt808.v1vehicle.raw.go.yutong-mqtt.v1vehicle.event.go.unified.v1
gateway 内部异步队列发布 NATS,降低协议连接被 Kafka 慢写拖住的风险。
durable pull consumer 批量拉 NATS,写 Kafka 成功后才 ACK NATS。
与 NATS subject 同名:raw 三个协议 + unified 一个实时主题。
Kafka 写失败不 ACK,NATS 保留消息等待重试;未知 subject 不 ACK,防止静默丢数据。
消费 RAW topic,写 TDengine:raw_frames、payload_chunks、locations、mileage_points。
消费 RAW topic,根据总里程差值写 MySQL 每日指标。
消费 unified topic,写 Redis 在线状态、实时合并快照、各协议 realtime-raw。
realtime-api 暴露实时、历史、统计查询接口,端口 :20200。
2. 单条报文的完整生命周期
3. FrameEnvelope:Go 版本内部统一数据格式
| 字段 | 含义 | 来源 |
|---|---|---|
event_id | 稳定事件 ID | 已有则保留,否则按协议/消息/车辆/时间/raw 计算 |
protocol | 协议 | GB32960、JT808、YUTONG_MQTT |
message_id | 协议消息类型 | 如 808 0x0200,32960 命令标识等 |
vin | VIN | 协议原始字段或 MySQL binding 反查 |
phone | 808 终端手机号 | 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 | 全量解析 JSON | TDengine parsed_json;Redis realtime-raw |
fields | 最小核心字段 | TDengine 位置/里程点;Redis 合并快照;统计计算 |
event_time_ms | 车端事件时间 | TDengine 主时间;Redis field time |
received_at_ms | 平台接收时间 | 延迟分析、在线状态 last_seen |
parse_status | 解析状态 | OK、PARTIAL、BAD_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 的职责分工
部署在 Kafka ECS,Docker Compose 管理;监听内网 172.17.111.56:4222,监控 172.17.111.56:8222。gateway 只需要把事件快速、可靠地写入 NATS。
使用 durable pull consumer;批量写 Kafka;只有 Kafka 写入成功才 ACK NATS。Kafka 短暂异常时,NATS 消息保留并重试。
保持现有 topic 和消费者模型,history/realtime/stat 不直接依赖 gateway,也不直接依赖 NATS API。
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 key | TTL 内有数据即在线,默认 TTL 600 秒。 |
| VIN 实时数据 | Redis merged snapshot | 跨协议合并 fields,并保留各协议 ProtocolData。 |
| 单协议实时 RAW | Redis 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;生产端口 32960、808 |
承接外部平台真实连接,发布到 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-writer、stat-writer、realtime-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. 关键可靠性设计
gateway 先写 NATS,Kafka 慢写由 bridge 消化,降低 32960/808 连接积压风险。
bridge 写 Kafka 失败时不 ACK,消息仍在 JetStream,恢复后继续拉取。
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 查询。