From f3b4cbbbb24b2f3cc70149b74d0bc5736d00ca13 Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 15:48:14 +0800 Subject: [PATCH] docs: add go data flow diagram --- docs/go-version-data-flow.html | 759 +++++++++++++++++++++++++++++++++ 1 file changed, 759 insertions(+) create mode 100644 docs/go-version-data-flow.html diff --git a/docs/go-version-data-flow.html b/docs/go-version-data-flow.html new file mode 100644 index 00000000..f379fb34 --- /dev/null +++ b/docs/go-version-data-flow.html @@ -0,0 +1,759 @@ + + + + + + Lingniu Vehicle Ingest Go Version Data Flow + + + +
+

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特殊动作
GB32960gateway TCP :32960报文本身通常有 VIN;平台账号用于登录鉴权完整 parsed JSON 写 raw_frames.parsed_json速度、经纬度、SOC、累计里程等按帧类型抽取登录/数据帧需要平台应答,成功发布后 ACK
JT808gateway TCP :808消息头 phone 为主;注册/鉴权/位置帧结合 binding 找 VIN完整 parsed JSON 写 raw_frames.parsed_json位置、速度、方向、状态、报警、GPS 总里程注册/鉴权/通用应答;每日里程按总里程差值统计
YUTONG_MQTTgateway 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写入方内容用途
TDengineraw_frameshistory-writer原始报文、完整 parsed JSON、fields JSON、parse status、source endpointRAW 查询、追溯、排查解析问题
TDengineraw_frame_payload_chunkshistory-writer超长 raw/parsed payload 分片避免 BINARY 长度限制导致完整 JSON 丢失
TDenginevehicle_locationshistory-writer经纬度、速度、方向、状态、报警、总里程等核心字段高频历史位置分页查询
TDenginevehicle_mileage_pointshistory-writer累计里程点、速度、经纬度里程曲线、区间差值分析
Redisvehicle:realtime:{vehicle_key}realtime-api Kafka consumer跨协议合并后的实时字段快照查 VIN 是否在线、查 VIN 实时数据
Redisvehicle:realtime:{vehicle_key}:{protocol}realtime-api Kafka consumer单协议实时快照区分 32960/808/MQTT 的实时状态
Redisvehicle:realtime-raw:{vehicle_key}:{protocol}realtime-api Kafka consumer单协议最新 parsed 全量字段实时 RAW 字段查看,不走 TDengine 历史扫描
Redisvehicle:online:{vehicle_key}vehicle:last_seenrealtime-api Kafka consumer在线状态、最后接收时间、协议列表在线判断和近期车辆列表
MySQLvehicle_identity_binding人工/导入维护,gateway 读取phone/device/plate 到 VIN 的映射808 等缺 VIN 协议反查正确 VIN
MySQLjt808_registrationgateway808 注册、鉴权、首次/最新上报、phone、device、plate、vin 解析结果定位哪些手机号没有映射 VIN,排查注册鉴权
MySQLvehicle_daily_metricstat-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 ECSlingniu-go-gateway公网 115.29.187.205;生产端口 32960808承接外部平台真实连接,发布到 NATS。
Gateway ECSlingniu-go-nats-kafka-bridge连接 NATS 172.17.111.56:4222;Kafka 172.17.111.56:9092当前部署在 gateway ECS,负责 NATS 到 Kafka 桥接。
Gateway ECShistory-writerstat-writerrealtime-apiAPI 115.29.187.205:20200消费 Kafka,写 TDengine/Redis/MySQL,并提供查询。
Kafka ECSKafka + NATS公网 114.55.58.251;内网 172.17.111.56NATS 使用 Docker Compose,端口只绑定内网 4222/8222
TDengine ECSTDengine内网 172.17.111.57:6041历史 RAW、位置、里程点时序存储。
云服务MySQL / RedisRDS 内网、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 查询。 +

+
+
+ + + +