Files
lingniu-vehicle-ingest/docs/architecture/production-data-plane-inventory.md

151 lines
16 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 生产数据面清单
审计时间2026-07-02
本文记录当前 ECS 上 Go 版本车辆接入链路的真实数据面用来约束后续重构新增表、topic、key 之前先对照这里,避免把已经删除的历史设计重新带回来。
## 运行服务
生产接入 ECS`115.29.187.205`
| 服务 | systemd 单元 | 端口 | 职责 |
| --- | --- | --- | --- |
| Gateway | `lingniu-go-gateway.service` | `0.0.0.0:32960``0.0.0.0:808``127.0.0.1:20211` | GB32960、JT808、宇通 MQTT 接入、协议解析、即时鉴权响应和内存身份解析 |
| NATS fast writer | `lingniu-go-nats-fast-writer.service` | `127.0.0.1:20215` | NATS raw 消费,写 Redis 实时投影;生产默认关闭 TDengine stage |
| NATS Kafka bridge | `lingniu-go-nats-kafka-bridge.service` | `127.0.0.1:20214` | NATS canonical raw 到 Kafka raw/fields 的可靠桥接和字段投影 |
| History writer | `lingniu-go-history-writer.service` | `127.0.0.1:20212` | 多 worker 消费 Kafka raw写 TDengine |
| Stat writer | `lingniu-go-stat-writer.service` | `127.0.0.1:20213` | 多 worker 消费 Kafka fields写每日里程 |
| Realtime writer | `lingniu-go-realtime-writer.service` | `127.0.0.1:20216` | 多 worker 消费 Kafka raw写 MySQL 实时表;生产关闭 Redis projector |
| Identity writer | `lingniu-go-identity-writer.service` | `127.0.0.1:20217` | 多 worker 消费 Kafka JT808 raw幂等写注册、鉴权和低频在线触达事实 |
| Realtime API | `lingniu-go-realtime-api.service` | `0.0.0.0:20200` | 读取 Redis/MySQL/TDengine 并提供查询 API不消费 Kafka |
## 总线
### NATS JetStream
NATS 部署在 Kafka ECS 内网 `172.17.111.56:4222`
| 项 | 当前值 |
| --- | --- |
| Stream | `VEHICLE_INGEST` |
| Subjects | Gateway 只新增三类 `vehicle.raw.go.*`;三类 `vehicle.fields.go.*` 暂留在 Stream 合同中兼容切换前消息,不再由 Gateway 直接写入 |
| Durable consumer | `vehicle-kafka-bridge` |
| 语义 | Gateway 每帧只发布一份 canonical rawbridge 复用其中已计算的 `parsed_fields` 生成 fieldsKafka raw 与 fields 都成功后才 ACK 同一条 NATS raw |
| 保留策略 | `NATS_STREAM_MAX_BYTES=21474836480``NATS_STREAM_MAX_AGE_HOURS=24``NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60`NATS 仅做入口缓冲 |
协议 parser 和字段扁平化都只在 Gateway 接入线程执行一次。TDengine、Redis、MySQL、bridge 和统计链路只能读取 envelope 中预计算的 `parsed_fields`;实时帧缺失该字段时 RAW 仍保留,但实时投影会以 `skipped_missing_fields` 跳过,禁止在存储层重新解析或误标在线。
Gateway 到 JetStream 的生产接受边界是本地分段 WALcanonical RAW 先进入 1ms/256 条分组提交并完成 `fsync`,然后才异步发布 JetStream只有收到 PubAck 才把记录标记完成,提交失败或 ACK 超时会释放给有界重放。进程崩溃后从 WAL 重建待发布记录,使用稳定 `Nats-Msg-Id=kind:subject:event_id` 让崩溃窗口内的重复重放由 JetStream 去重。生产基线为 16MiB/5s 分段、100000 append queue、10000 PubAck inflight禁止和旧逐条 JSON spool 共用目录。Redis 快路径使用 `FAST_WRITER_WORKERS=16``FAST_WRITER_BATCH_SIZE=100``FAST_WRITER_FETCH_WAIT_MS=5`,仍在 Redis 成功后手工 ACK NATS。
### Kafka
当前 Go 链路正在使用的 topic
| Topic | 分区 | 写入方 | 主要消费方 |
| --- | --- | --- | --- |
| `vehicle.raw.go.gb32960.v1` | 12 | NATS Kafka bridge | history、realtime |
| `vehicle.raw.go.jt808.v1` | 12 | NATS Kafka bridge | history、realtime、identity |
| `vehicle.raw.go.yutong-mqtt.v1` | 12 | NATS Kafka bridge | history、realtime |
| `vehicle.fields.go.gb32960.v1` | 12 | NATS Kafka bridge 从 canonical raw 投影 | stat |
| `vehicle.fields.go.jt808.v1` | 12 | NATS Kafka bridge 从 canonical raw 投影 | stat |
| `vehicle.fields.go.yutong-mqtt.v1` | 12 | NATS Kafka bridge 从 canonical raw 投影 | stat |
| `vehicle.event.go.unified.v1` | 12 | 显式打开 `PUBLISH_UNIFIED_ENABLED=true` 后才写入 | 兼容旧统一事件消费;当前最小链路不依赖 |
审计时 Kafka 上仍存在旧 Java/Xinda/telemetry topic例如 `vehicle.raw.gb32960.v1``vehicle.raw.jt808.v1``vehicle.event.xinda.v1``vehicle.raw.telemetry-input.v1`。这些不属于当前 Go 最小链路,后续如果确认没有消费者依赖,应单独做 Kafka topic 下线清单,不在业务代码里继续引用。
history、stat、realtime-writer、identity-writer 四个消费服务会以 `vehicle_kafka_consumer_info{service,group,topic}` 暴露当前运行时 Kafka 订阅配置;容量巡检将实时 writer 采集为 `realtime`,将身份 writer 采集为 `identity`,将查询服务采集为 `realtime-api`,以此区分“消费线程实际挂载 topic”和“API 可用性”。
history、stat、realtime-writer、identity-writer 的 Kafka 批消费采用 failure-closed 语义存储失败后不抓取下一批history/realtime/stat 只提交各分区连续成功前缀并保留失败后缀identity 的 MySQL 批事务整体重试commit 失败时四个 writer 都只重试 commit不重复写存储。`vehicle_{history,realtime,stat}_retry_pending_messages``vehicle_identity_writer_retry_pending_messages` 表示当前卡在失败重试中的消息数,生产健康值必须为 `0`
history、realtime、stat 和 identity writer 默认都在一个进程内启动 3 个同 group Kafka consumer由 Kafka 分配 topic 分区。bridge 使用车辆键作为 Kafka key因此同一车辆、同一协议的消息仍在单一分区内有序处理不同分区分别并行写 TDengine、MySQL 实时投影、MySQL 统计和 JT808 身份事实。`HISTORY_WORKERS``REALTIME_WORKERS``STATS_WORKERS``IDENTITY_WRITER_WORKERS` 只能在 topic 分区数和后端连接容量允许的范围内调整,不能靠并发打乱单车状态、累计里程或注册鉴权顺序。
生产环境中 `nats-fast-writer` 只负责 Redis 快速当前态,`FAST_WRITER_TDENGINE_ENABLED=false`TDengine RAW/位置历史由 `history-writer` 通过 Kafka raw topic 单路写入,避免 NATS fast path 与 Kafka history path 双写同一证据层。
生产环境中 Redis 当前态只允许 `nats-fast-writer` 写入。`realtime-writer``REALTIME_REDIS_PROJECTOR_ENABLED=false`,只保留 MySQL 当前态投影;`realtime-api` 只负责查询。指标 `vehicle_realtime_redis_projector_enabled` 必须为 `0`,否则 `capacity-check` 判为 degraded避免 Kafka 慢链路覆盖 NATS 快链路的实时结果。
## TDengine
数据库:`lingniu_vehicle_ts`
| Stable | 主列 | Tags | 职责 |
| --- | --- | --- | --- |
| `raw_frames` | `ts``frame_id``event_id``message_id``event_time``received_at``raw_size_bytes``raw_hex``raw_text``parsed_json``parse_status``parse_error``source_endpoint` | `protocol``vehicle_key``vin``phone``device_id` | RAW 证据层,保存完整解析 JSON |
| `raw_frame_payload_chunks` | 分片 payload 列 | 协议和车辆标识 tags | 保存超长 raw/parsed payload |
| `vehicle_locations` | `ts``event_id``frame_id``received_at``longitude``latitude``altitude_m``speed_kmh``direction_deg``alarm_flag``status_flag``total_mileage_km``soc_percent` | `protocol``vin` | 高频历史位置、SOC 和总里程查询;只写 VIN 非空车辆。writer 启动时幂等执行 `ADD COLUMN soc_percent`,兼容已有 stable |
不应重新出现:
- `vehicle_mileage_points`
- `raw_frames.fields_json`
## MySQL
数据库:`lingniu_vehicle_data`
| 表 | 核心字段 | 写入方 | 职责 |
| --- | --- | --- | --- |
| `vehicle` | `vin``plate``oem``enabled` | 导入/人工维护Gateway 只读 | 车辆主实体,承载 VIN 到默认车牌/厂商的轻量索引 |
| `vehicle_identifier` | `protocol``source_code``identifier_type``identifier_value``vin``plate``oem` | 映射文件导入/人工维护Gateway 只读 | 多平台、多标识反查 VIN当前主要承载 JT808 手机号/车牌映射 |
| `vehicle_identity_binding` | `vin``plate``phone``oem` | 导入/人工维护Gateway 兼容只读 | 旧版 VIN/车牌/手机号事实表;新映射导入会用它反查 VIN运行期解析优先使用 `vehicle_identifier` |
| `jt808_registration` | `phone``device_id``plate``vin``manufacturer``auth_token``source_endpoint`、首次/最新注册鉴权时间 | Identity writer | JT808 注册、鉴权和 VIN 匹配状态;可从 Kafka JT808 raw 重放重建 |
| `vehicle_realtime_snapshot` | `protocol``vin``plate``platform_name``peer``parsed_json``event_time``received_at``event_id``access_first_seen_at/access_previous_received_at/access_latest_received_at/access_report_interval_ms/access_sample_count/access_latest_event_id/access_first_seen_source` | Realtime writer | 每协议每 VIN 最新合并扁平字段快照;同一条原子 upsert 还维护与设备事件时间解耦的接收证据。重复 event ID、相同或回退接收时间不推进连续间隔旧行只做 `snapshot_backfill` 上线基线,不宣称历史首次接入 |
| `vehicle_realtime_location` | `protocol``vin``plate`、经纬度、速度、总里程、SOC、事件时间 | Realtime writer | 每协议每 VIN 最新位置业务缓存;由 Kafka raw 写入 MySQL不参与 Redis 当前态 |
| `vehicle_data_source` | `protocol``source_ip``source_code``platform_name``trust_priority``enabled` | Stat writer / 映射导入 / 人工维护 | 车辆数据来源管理;`source_code` 是机器可用的稳定来源编码,`platform_name` 是展示和人工维护名称 |
| `vehicle_daily_mileage_source` | `vin``stat_date``protocol``source_key``source_ip`、首末总里程、样本数、质量状态、是否选中 | Stat writer | 每协议、每来源的每日里程事实层;高频更新,保留多源候选 |
| `vehicle_daily_mileage` | `vin``stat_date``protocol``source_id``daily_mileage_km``latest_total_mileage_km` | Stat writer | 对外查询的每日里程结果层;由 `vehicle_daily_mileage_source` 按来源优先级/质量选举投影 |
身份表使用业务主键:`vehicle``vin` 为主键,`vehicle_identifier``(protocol, source_code, identifier_type, identifier_value)` 为主键,`vehicle_identity_binding``vin` 为主键,`jt808_registration``phone` 为主键;这些核心身份表不保留代理自增主键和 `created_at``vehicle_identity_binding` 是外部维护的兼容事实表,服务运行时不回写。
Gateway 默认每 60 秒把 `vehicle_identity_binding``vehicle_identifier``jt808_registration``vehicle_data_source` 原子刷新为本机只读身份快照。808 每帧只查内存,注册帧仍会立即更新当前 Gateway 的 phone 会话并返回鉴权码;持久化由 `identity-writer` 消费 Kafka JT808 raw 完成Gateway 必须设置 `JT808_REGISTRATION_GATEWAY_WRITES_ENABLED=false`。注册、鉴权帧不节流,普通位置帧默认每个 phone 每 10 分钟触达一次MySQL 写失败时不提交 Kafka offset恢复后继续重放。快照刷新失败继续使用上一版MySQL 不可用不会阻断 TCP 接入。
NATS→Kafka Bridge 只消费 canonical RAW并从 RAW 中已有的 `parsed_fields` 派生 fields envelope不重新解析协议。单批次按 Kafka topic 分组并发写入,默认 `BRIDGE_KAFKA_WRITE_CONCURRENCY=6`;同一个 RAW 对应的 RAW/fields topic 全部成功后才 ACK NATS失败时保留源消息重放。因此并行化只缩短独立 topic 的网络确认等待,不改变 RAW 证据和统计字段的一致性边界。
日里程事实表 `vehicle_daily_mileage_source` 使用 `(vin, stat_date, protocol, source_key)` 作为业务主键,保留首末总里程和样本数用于解释差值计算;对外结果表 `vehicle_daily_mileage` 使用 `(vin, stat_date, protocol)` 作为业务主键,不保留自增 `id``created_at` 和候选来源明细。
Stat writer 每条有效里程样本都会更新 `vehicle_daily_mileage_source`,但 `vehicle_daily_mileage` 的选举投影默认按 `STATS_PROJECT_INTERVAL_SECONDS=15` 节流;新车辆/新协议/新日期/新来源会立即投影。这样保留每日里程事实的实时性,同时避免高频帧对 MySQL 结果表做多 SQL 放大。
每日里程只采用同来源累计里程边界差值:`当日最新累计总里程 - 最近历史日最后累计总里程`。优先取前一自然日;前一日没有数据时继续向更早日期查找,历史完全为空才以当天第一条为基线。不同协议、平台和终端来源分别计算候选,再由来源质量和优先级选举最终结果,不混用两个来源的累计总里程。
车辆数据中台按用户启用的协议顺序查询“单车、单日”结果,但只在有有效增量的来源之间应用优先级:首选协议为 `0 km`、其他已启用协议存在正向里程时,查询层降级到有增量的最高优先级来源;全部来源均为 `0 km` 时仍严格按原优先级返回 `0`。该规则只选择一个完整协议事实,不拼接不同协议的首末累计里程。
`vehicle_data_source``(protocol, source_ip)` 唯一识别来源。直连终端可能产生很多来源 IP这些记录可以保留但默认不要求平台名转发平台来源通过 `vehicle_identifier.source_code``jt808_registration.source_endpoint` 推断后补充 `source_code/platform_name`,后续统计优先使用 `source_id/trust_priority/enabled` 做可信源选择。
统计分日优先使用协议事件时间;如果事件时间明显晚于接收时间超过 10 分钟则认为设备时间异常使用接收时间归属统计日。RAW 历史仍保留原始事件时间,便于追溯。
两张实时当前态表均使用 `(protocol, vin)` 作为业务主键,不保留自增 `id``created_at`;最新更新时间使用 `updated_at`
不应重新出现:
- `vehicle_daily_metric`
- `vehicle_identity_binding_registration`
- `vehicle_identity_bindings`
## Redis
Redis 使用 DB 50定位为实时缓存不作为历史事实来源。
| Key 族 | 用途 |
| --- | --- |
| `vehicle:latest:{vin}` | 跨协议合并后的最新核心字段快照,不重复保存完整 parsed只写 VIN 非空车辆 |
| `vehicle:latest:{vin}:{protocol}` | 单协议最新轻量快照,保留核心 fields 和时间信息,不重复保存完整 parsed |
| `vehicle:realtime-raw:{protocol}:{vin}` | 单协议最新完整 parsed 状态,是实时完整协议字段的唯一 Redis 副本 |
| `vehicle:rt-kv:{protocol}:{vin}:values` | 单协议扁平化实时字段值,字段名遵循协议字段映射 |
| `vehicle:rt-kv:{protocol}:{vin}:types` | `values` 中每个字段的值类型 |
| `vehicle:rt-kv:{protocol}:{vin}:times` | `values` 中每个字段最后写入的归一事件时间;快路径用它防止乱序旧帧覆盖新字段 |
| `vehicle:rt-kv:{protocol}:{vin}:meta` | 单协议 KV 投影的最新事件、接收时间和字段映射版本 |
| `vehicle:online:{protocol}:{vin}` | 在线状态和 TTL |
| `vehicle:online-state:{protocol}:{vin}` | 在线状态的 Hash 副本,便于分页查询 |
| `vehicle:protocols:{vin}` | 当前车辆最近出现过的协议集合 |
| `vehicle:last_seen` | 最近活跃车辆排序集合 |
审计时活跃 key 族主要是 `vehicle:latest:*``vehicle:realtime-raw:*``vehicle:rt-kv:*``vehicle:online:*``vehicle:online-state:*``vehicle:protocols:*``vehicle:last_seen`。旧文档中的 `vehicle:realtime:*``vehicle:merged:*` 不是当前活跃 key 族。
## 后续优化约束
1. 接入层每帧只解析和扁平化一次并写一份 NATS canonical rawKafka 是持久回放层fields 必须由可重放消费者从 raw 中已有的 `parsed_fields` 投影,不能由 Gateway 双写或由存储层重算。
2. TDengine 只放高写入时序数据RAW 证据和位置历史。
3. MySQL 只放低基数业务状态:身份、实时轻量快照、每日里程。
4. Redis 只放当前态,所有 key 都必须允许 TTL 过期后从 Kafka/TDengine/MySQL 重建。
5. 协议新增字段默认进入 `raw_frames.parsed_json` 物理列中的扁平 `parsed_fields` JSON并同步进入 Redis KV只有稳定查询需求出现后才提升为 TDengine/MySQL 独立列。
6. Gateway 不直接写业务数据库JT808 即时会话留在内存,注册鉴权事实由 Kafka identity-writer 单路投影。