104 lines
6.4 KiB
Markdown
104 lines
6.4 KiB
Markdown
# 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/s,P99 < 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 1078(optional-command-gateway profile;808 信令按需桥接 + 常用下行信令编码 + 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)
|
||
│ │ └── inbound-xinda-push/ 信达 Push 接入(废弃兼容模块,仅 -Plegacy-xinda 显式构建)
|
||
│ ├── 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 全字段事件消费 + 可配置日统计
|
||
│ └── 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/ 参考资料
|
||
```
|
||
|
||
## 快速开始
|
||
|
||
```bash
|
||
# 要求: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.md` 与 `docs/operations/vehicle-ingest-tdengine-verification.md`。
|
||
|
||
信达 Push 已废弃,不参与默认 Maven reactor、CI、Portainer 部署和历史消费链路;如需兼容排查,可显式使用 `-Plegacy-xinda` 构建。
|
||
|
||
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. **接入与业务落库分离**:协议接入应用只负责收、解析、校验、规整和投递 Kafka;`vehicle-history-app` 写入 TDengine 历史与 RAW JSON,`vehicle-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` 追加到事件 metadata;Kafka Envelope、TDengine `raw_frames` 和导出都使用 `archive://...` 逻辑 URI 追溯原始 bytes,实际文件由 `sink-archive` 管理
|