Files
lingniu-vehicle-ingest/README.md
2026-07-01 18:32:09 +08:00

6.3 KiB
Raw Permalink Blame History

lingniu-vehicle-ingest

羚牛车辆数据接入平台 v2

设计目标

  • 协议接入统一抽象GB/T 32960、JT/T 808、Yutong MQTT 是默认生产接入JT/T 1078、JSATL12 作为显式 profile 的可选能力;信达 Push 已废弃并删除源码
  • 原子能力化:每个协议一个独立 Maven 模块 + 独立 AutoConfiguration + 独立配置开关
  • 业务 / 实时彻底解耦:协议接入应用只负责收、解析、校验、规整和投递 Kafka历史与统计由独立 Kafka 消费应用落库
  • 高并发低延迟Netty + Disruptor + Java 25 虚拟线程;目标单节点 ≥ 5 万 msg/sP99 < 50 ms
  • 可观测、可回放、可灰度:统一 traceId、Envelope、幂等键、DLQ、原始报文冷存

技术栈

类别 选型
语言 Java 25
框架 Spring Boot 3.5.3
网络 Netty 4.2.9.Final
并发 Disruptor 4 + 虚拟线程
事件流 Kafka唯一生产消息通道vin 分区,保证单车有序)
线上消息格式 Protobuf,调试走 JSON
MQTT 客户端 Eclipse Paho 1.2.5
会话索引 Redis
熔断/限流 Resilience4j 2
可观测 Micrometer + Prometheus
冷存 本地文件系统
构建 Maven

模块划分

lingniu-vehicle-ingest/
├── modules/
│   ├── core/
│   │   ├── ingest-api/              SPI + sealed 领域事件 + 注解
│   │   ├── ingest-codec-common/     公共编解码工具BCD/CRC/BCC/bit utils
│   │   ├── ingest-core/             Pipeline / Dispatcher / Disruptor / Session 桥
│   │   ├── session-core/            设备会话 + 鉴权 + Token + Redis SessionStore
│   │   ├── vehicle-identity/        跨协议车辆身份解析 + MySQL 外部标识绑定
│   │   └── observability/           metrics / health
│   ├── protocols/
│   │   ├── protocol-gb32960/        GB/T 32960
│   │   ├── protocol-jt808/          JT/T 808统一身份映射 + 事件 metadata 内部 VIN + 注册/鉴权/心跳/注销/位置/批量位置/参数/属性/媒体/透传/未知上行和坏帧兜底/断链清理会话/下行分包)
│   │   ├── protocol-jt1078/         JT/T 1078optional-command-gateway profile808 信令按需桥接 + 常用下行信令编码 + TCP/UDP RTP 媒体流分段归档 + 事件 metadata 内部 VIN + 坏 RTP 统一 RawArchive/Passthrough + 归档失败兜底 + SIM 身份映射)
│   │   └── protocol-jsatl12/        苏标主动安全报警附件optional-attachments profile
│   ├── inbound/
│   │   └── inbound-mqtt/            MQTT 接入endpoint 生命周期 + profile 注册扩展 + 统一身份映射 + PEM 双向 TLS + 未知 profile/解析失败/profile异常/连接订阅失败兜底 + 统一 UNKNOWN 身份 metadata
│   ├── sinks/
│   │   ├── sink-kafka/              Kafka producer + Protobuf Envelope
│   │   └── sink-archive/            原始报文冷存
│   ├── services/
│   │   ├── event-history-service/   Kafka 全字段事件消费 + 历史查询
│   │   ├── vehicle-state-service/   Kafka 全字段事件消费 + Redis 热状态查询optional-latest-state profile
│   │   └── vehicle-stat-service/    Kafka 全字段事件消费 + 可配置日统计
│   ├── testing/
│   │   └── vehicle-identity-test-support/  协议测试共享身份解析夹具
│   └── apps/
│       ├── command-gateway/         可选 HTTP → 设备下行命令optional-command-gateway profile
│       ├── gb32960-ingest-app/      GB32960 TCP 接入 + Kafka 投递
│       ├── jt808-ingest-app/        JT808 TCP 接入 + Kafka 投递
│       ├── yutong-mqtt-app/         宇通 MQTT 接入 + Kafka 投递
│       ├── vehicle-history-app/     TDengine 历史查询 + RAW JSON
│       └── vehicle-analytics-app/   JT808 每日里程指标消费
├── docs/                           架构文档、模块图、实施计划
└── reference/                      参考资料

快速开始

# 要求JDK 25, Maven 3.9+
mvn -v
mvn -pl :gb32960-ingest-app,:jt808-ingest-app,:yutong-mqtt-app,:vehicle-history-app,:vehicle-analytics-app -am package -Dmaven.test.skip=true

生产服务只在 ECS/Portainer 上运行本机只用于源码检查、Maven 构建和单元测试。部署、健康检查和真实流量验证步骤见 docs/operations/gb32960-service-split-runbook.mddocs/operations/vehicle-ingest-tdengine-verification.md

信达 Push 已废弃并删除源码,不参与 Maven reactor、CI、Portainer 部署和历史消费链路。

JSATL12 附件上传不在默认 Maven reactor 中;需要时使用 -Poptional-attachments 构建。

Command Gateway/JT1078 下行与音视频信令不在默认 Maven reactor 中;需要时使用 -Poptional-command-gateway 构建。

Redis 最新状态查询不在默认 Maven reactor 中;需要时使用 -Poptional-latest-state 构建 vehicle-state-service

迁移说明

本项目是 lingniu-vehicle-data-reception 的 v2 重构,采用 strangler fig 渐进式迁移,旧项目保留只读参考。迁移路径与决策参见 ../REFRACTOR_PLAN.md

架构文档

  • 目标架构:docs/target-architecture.md
  • 内部字段模型:docs/vehicle-telemetry-internal-fields.md
  • 模块与数据流:docs/module-data-flow.html

核心原则

  1. 接入与业务落库分离:协议接入应用只负责收、解析、校验、规整和投递 Kafkavehicle-history-app 写入 TDengine 历史与 RAW JSONvehicle-analytics-app 将 JT808 daily_mileage_km 写入 MySQL vehicle_stat_metric
  2. Kafka 唯一消息通道:生产链路只支持 Kafka接入、历史、统计统一通过 Kafka Envelope 解耦
  3. 协议即插拔:每个 protocol-* / inbound-* 都可独立开关(lingniu.ingest.<name>.enabled
  4. 顺序保证:同一 VIN 严格有序Disruptor hash + Kafka 分区 key
  5. 幂等消费Envelope 带 eventId,下游去重
  6. 原始可回放:协议接入会把 rawArchiveKey/rawArchiveUri 追加到事件 metadataKafka Envelope、TDengine raw_frames 和导出都使用 archive://... 逻辑 URI 追溯原始 bytes实际文件由 sink-archive 管理