AllenVB/Vehicle-Tracking-Simulation
GitHub: AllenVB/Vehicle-Tracking-Simulation
基于Kafka和Spring Boot构建的事件驱动车队遥测仿真平台,模拟车辆沿真实道路网络行驶,提供实时地图追踪、违规检测与驾驶员评分功能。
Stars: 0 | Forks: 0
# 车辆追踪系统
事件驱动的车队远程信息处理平台。来自模拟车辆设备的遥测数据沿着 **ingestion → Kafka → 处理/分析 → 通知 → API 网关** 的管线流动,并在**真实地图上实时**显示。
- **设计目标:** 1000 辆车,约 1000 条消息/秒,每天约 8600 万行。
- **开发 目标:** 覆盖土耳其的 100 辆车,1 秒 tick。
- **原则:** 规模扩展仅源自配置;架构从一开始就正确构建。
技术栈:**Java 21 · Spring Boot 3.3 · Apache Kafka (KRaft) · Kafka Streams · TimescaleDB + PostGIS · Redis · Leaflet · 多模块 Maven monorepo。**
## 屏幕截图
单页、单服务 (`:8080`),**12 列网格:2 列车队栏 · 5 列实时地图 · 5 列操作员地图**。

- **左侧 (2/12) — 车队栏:** 车辆列表、驾驶员评分、实时违规、选择和控制框。
- **中间 (5/12) — 实时地图** (OpenStreetMap):105 辆车行驶在**真实道路**上(OSRM 路线)。
每辆车都根据其**类型显示彩色图标** —— 汽车(蓝色)· 卡车(黄色)· 摩托车(白色)· 直升机(紫色) —— 并根据其行驶方向旋转;违规时变为红色。地图上还标记了**加油站** (⛽)。选中车辆时,其**即将行驶的路线**会以流动的虚线绘制,并显示**到达最近加油站的距离**。
- **右侧 (5/12) — 操作员地图** (CartoDB):选择车辆,**双击**新位置。
地面车辆只能移动到**道路上**;如果在道路外点击,会出现 *"无法到达此地点"* 的警告,车辆**保持在原地**。直升机可以降落在任何地方。移动的车辆**不会停下**:它会从放置点以自己的速度出发**前往新的目的地**。更改会通过真实的遥测管线,在 **~0.1 秒内反映到左侧地图上**。
- **违规**会列出其土耳其语名称和**以 TL 计价的罚款**;总罚款显示在栏上方。由于大多数车辆遵守限制,因此数据流比较稀疏(在几秒钟内不会像洪水般涌入)。
### 操作员:路线生成与车辆警告

- **路线生成:** 选中车辆后,操作员地图右下角会出现一个按钮;选择**目标省份**,车辆将被引导沿新路线(地面车辆走 OSRM 道路,直升机直线飞行)前往 —— 更改会立即反映在左侧地图上。
- **车辆警告:** 操作员可以向车辆发送**文本警告**(🔥 易燃物品,📦 易碎物品,❄️ 冷链,⚠️ 注意速度…)。警告是持久性的(点击车辆时可见),并通过 WebSocket 作为**即时通知** (toast) 发送给所有操作员。
两张地图都由**单个 WebSocket 订阅**提供数据 —— 没有 polling。
## 架构
数据单向流动。**操作员地图**闭环了反馈回路:通过 gateway 覆盖模拟器中的位置,该位置通过正常管线并返回到地图。
整个 UI 都托管在单个服务 (gateway) 中 —— 没有第二个前端服务。
```
flowchart TB
subgraph KAYNAK["1 - Kaynak"]
SIM["vts-simulator :8085
100 araç · 81 il
OSRM gerçek yol rotaları
Virtual Threads
(konumların TEK kaynağı)"] end subgraph GIRIS["2 - Giriş"] ING["vts-ingestion :8081
imei→vehicle lookup
stateless · DLQ"] end K{{"Apache Kafka · 24 partition
key = vehicleId"}} subgraph ISLEME["3 - İşleme ve Analitik"] PROC["vts-processing :8082
JDBC batch insert
durumsuz kurallar
+ ihlal cooldown"] STR["vts-stream-analytics :8083
Kafka Streams (RocksDB)
durumlu kurallar · trip · geofence"] end subgraph DEPO["4 - Depolama"] TS[("TimescaleDB + PostGIS
hypertable · continuous aggregate")] RD[("Redis
cache · cooldown")] end NOT["vts-notification :8084
Strategy sender · quiet hours"] SCH["vts-scheduler :8086
ShedLock · outbox publisher"] subgraph SUNUM["5 - Sunum (tek sayfa, tek origin)"] GW["vts-api-gateway :8080
JWT · REST · STOMP WebSocket
1 sn delta · viewport filtresi
operatör kontrolünü simülatöre proxy'ler"] UI["Tek sayfa UI
2 filo barı · 5 canlı harita · 5 operatör haritası
Leaflet (OSM + CartoDB)"] end UI -.->|"araç konumunu ez
(vehicleId)"| GW GW -.->|"proxy: /api/control
(imei index'e çevirir)"| SIM SIM -->|"POST /telemetry/batch"| ING ING -->|"vehicle.telemetry.raw"| K K --> PROC K --> STR PROC --> TS PROC --> RD PROC -->|"vehicle.violation + outbox"| K STR -->|"violation · geofence · trip"| K K --> NOT NOT -->|"vehicle.notification"| K K --> GW GW <--> TS SCH --> TS SCH --> K GW -->|"STOMP /topic/fleet/live
(tek abonelik, iki harita)"| UI ``` ### 关键流程:从操作员到地图 模拟器是车队的**唯一位置来源**。因此,从操作员地图进行的移动不是虚假的“地图作弊”,而是通过真实管线传输的真实遥测 —— 并且会在覆盖时立即发布(无需等待下一个 tick): ``` Sağ haritada çift tık → Gateway (proxy) → Simülatör (override + anında yayın) ↓ Ingestion → Kafka ↓ Sol harita ← STOMP delta ← Gateway ← Processing ~0.1 sn ``` ## 模块 | 模块 | 端口 | 职责 | |---|---|---| | `vts-common` | — | 事件模型、topic 常量、枚举、TenantContext、通用 Kafka 消费者支持(反序列化 + retry/DLQ 策略) | | `vts-simulator` | 8085 | 车队模拟器 (Virtual Threads)、OSRM 道路路线、位置覆盖 API(不提供 UI) | | `vts-ingestion-service` | 8081 | 无状态 HTTP 入口;imei→vehicle 查找 (Caffeine→Redis→DB)、Kafka 发布、DLQ | | `vts-processing-service` | 8082 | 批量消费者;JDBC 批量插入、无状态规则 **+ 违规冷却**、outbox | | `vts-stream-analytics` | 8083 | Kafka Streams;有状态规则(急刹车、持续超速、怠速、geofence、trip) | | `vts-notification-service` | 8084 | 策略发送器、冷却、免打扰时段 | | `vts-api-gateway` | 8080 | JWT 安全、REST、STOMP WebSocket、schema 所有者 (Flyway) **+ 单页 UI(两张地图)** + 操作员控制代理 | | `vts-scheduler-service` | 8086 | ShedLock 任务:离线检测、评分、维护、outbox 发布者 | ## 车队模型 ### 覆盖土耳其的分布 100 辆地面车辆根据人口加权分布到**所有 81 个省份**(+ 5 架直升机 = 105): | 省份分组 | 车辆 | 示例 | |---|---|---| | 最大的 3 个大都市 | 各 3 辆 | 伊斯坦布尔、安卡拉、伊兹密尔 | | 接下来的 13 个大都市 | 各 2 辆 | 布尔萨、安塔利亚、阿达纳、科尼亚、加济安泰普… | | 其余 65 个省份 | 各 1 辆 | 锡诺普、阿尔达汉、亚洛瓦… | ### 类别与车辆类型 车队使用**分类法**进行建模(`vehicle_category` + `vehicle_type` 表): | 类别 | 类型 | 数量 | |---|---|---| | **地面车辆** (`LAND`) | 汽车、卡车、摩托车 | 50 / 30 / 20 | | **空中车辆** (`AIR`) | 直升机 | 5 | | **海上车辆** (`SEA`) | — | 0 | `SEA` 故意留空:分类法旨在即使还没有一艘船,也能表明可以表示一个海上舰队。`GET /api/v1/vehicle-types` 返回结果包括类别为空的项。 **类别**是决定运动和规则适用性的核心轴:地面车辆受道路网络限制,而空中车辆则不受限。代码检查的是类别而不是类型 —— “这是直升机吗”这个问题,在两个服务中产生了两个不同的豁免列表。它还包含车牌类型:`VTS-001-Otomobil`、`VTS-027-Tır`、`VTS-063-Motor`。 ### 哪条规则适用于哪种类型?(`rule_assignment` → `VEHICLE_TYPE`) 违规**仅适用于地面车辆**;例外是燃料/电池规则,它们属于机器而不是道路。适用性和阈值数据都可以在无需更改代码的情况下进行调整: | 规则 | 汽车 | 卡车 | 摩托车 | 直升机 | |---|---|---|---|---| | `SPEED_LIMIT` / `SUSTAINED_SPEEDING` | 110 | 80 | 90 | **无** | | `HARSH_BRAKING` (减速度差值) | −40 | **−30** | **−50** | **无** | | `IDLING`, `GEOFENCE_ENTER/EXIT` | ✓ | ✓ | ✓ | **无** | | `LOW_FUEL` | 15 | 20 | 10 | **25** | | `LOW_BATTERY` | 20 | 20 | 20 | 20 | 对于急刹车,*较小*的 magnitude 是更严格的设置:满载的卡车如果在一次测量中损失 30 km/s,这属于剧烈刹车,而摩托车通常都会达到 50。通用的 −40 阈值会导致对卡车发现过晚,而对摩托车频繁标记。 如果表中没有某种类型的行,则应用规则的默认值;该表仅记录**差异**。如果 `enabled = false`,则该规则根本不适用于该类型。 启动时,地面车辆会被分配 **0–120 km/s 的随机基准巡航速度**;超过限制的车辆会产生违规(受下文提到的冷却限制)。 ### 直升机(车牌 101–105) 5 架直升机因为是飞行的,所以行为与地面车辆不同: - **直线飞行、高速** (180–260 km/s);路线不是 OSRM,而是直线航线 —— 直升机的真实路线本来就是直线的。 - **不适用任何基于道路的规则**:超速限制、持续超速、急刹车、geofence 和怠速。这些规则描述的都是车辆与道路或地面区域的关系;直升机没有这种关系。豁免现在集中在一个地方 —— `rule_assignment` 的 `VEHICLE_TYPE` 行 —— 且两个规则引擎都读取相同的行。 - **燃料和电池规则有效**(燃料阈值为 25,会更早发出警告):它们描述的是机器而不是道路,而燃料耗尽的直升机是最紧急的情况。 - **Trip 规则有效** —— 一次飞行也是一个 trip,而 trip 不是违规。 - 操作员可以将直升机放置在**任何地方**(海上、建筑物顶部);而地面车辆只能移动到道路上,点击道路外会被拒绝。 ### 行程:目的地、真实路径、剩余公里数 车辆不会漫无目的地乱转;它们进行**真实的行程**: 1. 从附近的省份选择一个**目的地**。 2. 从 OSRM 获取通往该目的地的**真实驾驶路线**(道路几何形状)。 3. 车辆沿路线行驶;**剩余公里数**根据真实道路长度计算(不是直线)。 4. 到达时 **速度降为 0 → 停靠 (5.5–12 分钟)** → trip 关闭 → 选择新的目的地。 这是一个唤醒系统另一半的设计决定:如果车辆不停下,**trip 就永远不会关闭**,那么 `trip`、`trip_point`、`stop_event`、驾驶员评分和维护数据就会*永久为空*。车辆在路线上以**随机进度**开始 —— 否则首次到达(和首次 trip)要在几个小时后才会出现。 ### 路线引擎是本地的 (OSRM) `docker compose` 会启动一个**本地 OSRM**(`osrm` 服务,包含土耳其的 OSM 数据)。首次启动非常耗时 —— 需要下载约 400MB 数据并进行几分钟的预处理 —— 但结果会保存在一个 volume 中,随后的启动会瞬间完成。系统因此独立于互联网运行。 模拟器以前连接到 OSRM 的**公共演示服务器**,当它变慢/无响应时,地面车辆会被赋予 **A 到 B 的直线**:车辆依然在行驶,依然会到达,trip 依然会关闭 —— 但是穿越了土耳其的山脉、湖泊和农田。因此,车辆穿过地形的原因不是错误,而是有意的“继续移动”的选择。 现在,**地面车辆如果没有道路路线就不会移动**。如果无法获取路线,车辆将停靠,分配器会每隔几秒钟重试一次。停止的车辆是路线引擎损坏的可见且诚实的标志;而不是飞越湖泊的车辆。生成合成路线的代码已被完全删除 —— 现在已经没有任何代码路径可以为地面车辆提供偏离道路的几何形状了。 ### 驾驶员评分 `driver_score_daily`:填充 30 天的历史数据(带有每位驾驶员持久的“性格”偏好 —— 否则每个人的平均分都一样,排名就失去了意义)。日常处理则根据**真实数据**计算评分:trip 距离 + 所有违规类型,罚款**按每 100 公里进行标准化**。如果不进行标准化,行驶时间长的驾驶员会自动显得“最差” —— 这是一个破坏评分表可靠性的经典错误。 ## 规模限制(从一开始就正确确立的决定) 1. **遥测数据绝不通过单独的 `save()` 写入** —— 批量 Kafka 消费者 + `JdbcTemplate.batchUpdate()` + `ON CONFLICT DO NOTHING`;没有用于遥测的 JPA entity;使用 `reWriteBatchedInserts=true`。 2. **没有每个事件一次的 Redis round-trip** —— 有状态的状态使用 Kafka Streams state store (RocksDB);批量 Redis 操作使用 pipeline。 3. **没有每个事件一次的 WebSocket 消息** —— gateway 在内存中保存状态,并使用 `@Scheduled(1s)` 仅发布更改的内容(delta);client 发送视口 (bbox)。 4. **违规不会针对每次读取产生** —— 违数量应与**离散事件**成正比,而不是与遥测速率成正比。在三个地方进行了 debounce 处理: - *无状态规则* (`SPEED_LIMIT`, `LOW_BATTERY`, `LOW_FUEL`):在 `RuleEngine` 中采用基于车辆+规则的冷却机制 —— 持续超速的车辆会在规则的 `cooldown_seconds` 时间窗口内产生一次违规。 - *急刹车*:在 `HarshBrakingRule` 中采用基于车辆的 120 秒冷却 (RocksDB state store)。 - *持续超速*:5 分钟的 **tumbling** 窗口。(以前是 `advanceBy(1 分钟)` 的 hopping 窗口;任何持续超速的车辆都会**每分钟**产生一次违规,仅此一项就占了数据洪流的 **83%**。) 测得的整体影响:**~20 违规/秒 → ~0.2 违规/秒**。 5. **内存不设无限** —— 如果没有容器限制,JVM 会根据*宿主机* RAM(%25)选择堆上限(heap ceiling);8 个服务相乘会导致 RAM 随着时间推移而膨胀。通过每个服务设置 `mem_limit` + `MaxRAMPercentage` 来固定堆。此外,由于 Kafka Streams 中的 **5 个 store × 24 个 partition ≈ 120 个 RocksDB 实例**各自会开启自己的缓存,因此它们全部共享一个有限且统一的 LRU 缓存 + write-buffer manager (`BoundedRocksDBConfig`)。Java 服务总计:**~2.4 GB,硬性上限 3.7 GB**。 6. **Trip 距离不是原始 GPS 的总和** —— 连续读取之间超过 2 公里的步进会被视为*重新定位*(操作员传送、重启、GPS 跳跃),不计入距离。如果计入,trip 长度以及随之而来的所有按公里标准化的驾驶员评分都会悄无声息地作废。 7. **Kafka partition = 24**(与 profile 无关) —— 以后增加会破坏按车辆的排序 (per-vehicle ordering) 和 Streams state store。 8. **遥测 = TimescaleDB hypertable** —— dashboard 查询指向 continuous aggregate 而不是原始表。 9. 每个表从一开始就带有 `tenant_id` + Outbox Pattern。 10. **实时流不能匿名监听** —— 在 STOMP `CONNECT` frame 中验证 JWT。由于无法在 SockJS 握手期间设置 `Authorization` header,因此握手保持公开;身份在 CONNECT 时检查。否则,任何人都可以在没有 token 的情况下观看整个车队。 ## 数据模型 Flyway `V1`–`V15` 定义了 **38 张业务表**。主要包括: - `telemetry` **hypertable**:`by_range(ts)` + `by_hash(vehicle_id, 8)`,PK `(vehicle_id, ts)`,无 FK(保证批量插入速度)。 - `violation` **hypertable**。 - `vehicle_driver_assignment`:用于将违规准确归咎于当时驾驶员的时间记录。 - `rule` + `rule_assignment`:阈值绝不在代码中;在 TENANT/GROUP 范围内被覆盖(基于类型的限速即来源于此)。 - Continuous aggregate:`telemetry_1min`、`telemetry_hourly`、`violation_daily_summary`。 - 压缩 + retention 策略,GIST/BRIN/partial index。 ## 运行 整个系统只需一条命令即可启动(基础设施 + 8 个服务): ``` docker compose up -d --build ``` 然后在浏览器中访问: | 界面 | 地址 | 登录 | |---|---|---| | **VTS — 单页(车队栏 + 两张地图)** | http://localhost:8080 | `admin` / `password` | | Swagger UI | http://localhost:8080/swagger-ui.html | JWT | | Kafka UI | http://localhost:8090 | — | | Prometheus | http://localhost:9090 | — | | Grafana | http://localhost:3000 | `admin` / `admin` | 其他端口:ingestion 8081,processing 8082,stream-analytics 8083,notification 8084,scheduler 8086,**Postgres 5433**,Redis 6379。 负载 profile(1000 辆车 / 1 秒,3 个 broker 覆盖): ``` docker compose -f docker-compose.yml -f docker-compose.load.yml up -d ``` ### API 示例 ``` # 登录(开发用户:admin / password) TOKEN=$(curl -s -X POST localhost:8080/api/v1/auth/login \ -H 'Content-Type: application/json' \ -d '{"username":"admin","password":"password"}' | jq -r .token) curl localhost:8080/api/v1/vehicles -H "Authorization: Bearer $TOKEN" curl localhost:8080/api/v1/live/positions -H "Authorization: Bearer $TOKEN" curl localhost:8080/api/v1/dashboard/summary -H "Authorization: Bearer $TOKEN" curl "localhost:8080/api/v1/violations?limit=20" -H "Authorization: Bearer $TOKEN" ``` 实时地图 WebSocket (STOMP):`ws://localhost:8080/ws` → `/topic/fleet/live`、`/topic/violations`、`/user/queue/notifications`; 针对视口的 `/app/viewport`。 ### 操作员控制 API (gateway 代理 — 通过 `vehicleId`) ``` # 工具的目标/路线状态 curl localhost:8080/api/v1/control/state -H "Authorization: Bearer $TOKEN" # 移动工具(使用 vehicleId — 而非车牌号) curl -X POST localhost:8080/api/v1/control/49/position -H "Authorization: Bearer $TOKEN" \ -H 'Content-Type: application/json' -d '{"lat":39.92,"lon":32.85}' ``` 移动是一种**锁定,而不是传送**:车辆降落在放置点,并从那里以自己的速度出发**前往新的目的地**。目的地是根据车辆的*新*位置来选择的 —— 保留旧目的地会把移动变成一个自我取消的操作:一辆从哈塔伊移动到安卡拉的卡车仍然需要前往开塞利,并且会立即掉头向东南方向返回。 响应通过 `moved` 字段告知结果 —— 点击道路外会返回 200 状态码但 `moved:false`,因为请求已被理解,**拒绝即是回答本身**: ``` {"found":true,"flying":false,"moved":false,"reason":"OFF_ROAD","offRoadMeters":10576} ``` 地面车辆要将点击视为“在道路上”,最近的道路必须在 `vts.simulator.road-click-tolerance-meters`(默认为 **50 m**)范围内。这个阈值是对“点击在道路上”的定义:如果太窄,视觉上对准道路的点击会被拒绝(点击包含了缩放误差、道路绘制粗细和 OSM 几何误差),如果太宽,车辆会悄无声息地停在操作员未指示的地方。 Gateway 通过 imei 将 `vehicleId` 转换为模拟器的设备索引;模拟器的 `:8085/api/**` 端点保留在内网中,UI 永远不会直接访问它。 ## Profiles | Profile | 车辆 | Tick | Chunk | Retention | |---|---|---|---|---| | `dev` (默认) | 100 (分布到 81 个省份) | 1 秒 | 1 天 | 30 天 | | `load` | 1000 | 1 秒 | 1 小时 | 7 天 | | `prod` | 外部配置 | — | — | — | ## 测试 - **Kafka Streams:** `TopologyTestDriver`(急刹车、怠速、geofence 进出、trip)。 - **单元测试:** 规则引擎 **+ 违规冷却窗口**、ingestion 路由、notification 冷却/免打扰时段、JWT、实时地图 delta+viewport、模拟器运动和速度模型。 - **Schema 验证:** JPA entity 针对运行中的 TimescaleDB 进行 `ddl-auto=validate` (Testcontainers)。 - **端到端:** 在真实容器中验证了 simulator → Kafka → DB 数据流(遥测、最后位置、违规、基于类型的阈值覆盖、驾驶员归属、操作员覆盖在主地图上的反映)。 ``` mvn test ``` ## 可观测性 Micrometer + Prometheus + Grafana。指标:`telemetry.ingested`、`telemetry.persisted`、`violation.produced`、`notification.sent`、consumer lag、DLQ 比率。Grafana 会自动加载一个预置的 **"VTS — Fleet Telematics Overview"** dashboard。每个事件都带有 `correlationId` 的结构化 JSON 日志。
- **路线生成:** 选中车辆后,操作员地图右下角会出现一个按钮;选择**目标省份**,车辆将被引导沿新路线(地面车辆走 OSRM 道路,直升机直线飞行)前往 —— 更改会立即反映在左侧地图上。
- **车辆警告:** 操作员可以向车辆发送**文本警告**(🔥 易燃物品,📦 易碎物品,❄️ 冷链,⚠️ 注意速度…)。警告是持久性的(点击车辆时可见),并通过 WebSocket 作为**即时通知** (toast) 发送给所有操作员。
两张地图都由**单个 WebSocket 订阅**提供数据 —— 没有 polling。
## 架构
数据单向流动。**操作员地图**闭环了反馈回路:通过 gateway 覆盖模拟器中的位置,该位置通过正常管线并返回到地图。
整个 UI 都托管在单个服务 (gateway) 中 —— 没有第二个前端服务。
```
flowchart TB
subgraph KAYNAK["1 - Kaynak"]
SIM["vts-simulator :8085100 araç · 81 il
OSRM gerçek yol rotaları
Virtual Threads
(konumların TEK kaynağı)"] end subgraph GIRIS["2 - Giriş"] ING["vts-ingestion :8081
imei→vehicle lookup
stateless · DLQ"] end K{{"Apache Kafka · 24 partition
key = vehicleId"}} subgraph ISLEME["3 - İşleme ve Analitik"] PROC["vts-processing :8082
JDBC batch insert
durumsuz kurallar
+ ihlal cooldown"] STR["vts-stream-analytics :8083
Kafka Streams (RocksDB)
durumlu kurallar · trip · geofence"] end subgraph DEPO["4 - Depolama"] TS[("TimescaleDB + PostGIS
hypertable · continuous aggregate")] RD[("Redis
cache · cooldown")] end NOT["vts-notification :8084
Strategy sender · quiet hours"] SCH["vts-scheduler :8086
ShedLock · outbox publisher"] subgraph SUNUM["5 - Sunum (tek sayfa, tek origin)"] GW["vts-api-gateway :8080
JWT · REST · STOMP WebSocket
1 sn delta · viewport filtresi
operatör kontrolünü simülatöre proxy'ler"] UI["Tek sayfa UI
2 filo barı · 5 canlı harita · 5 operatör haritası
Leaflet (OSM + CartoDB)"] end UI -.->|"araç konumunu ez
(vehicleId)"| GW GW -.->|"proxy: /api/control
(imei index'e çevirir)"| SIM SIM -->|"POST /telemetry/batch"| ING ING -->|"vehicle.telemetry.raw"| K K --> PROC K --> STR PROC --> TS PROC --> RD PROC -->|"vehicle.violation + outbox"| K STR -->|"violation · geofence · trip"| K K --> NOT NOT -->|"vehicle.notification"| K K --> GW GW <--> TS SCH --> TS SCH --> K GW -->|"STOMP /topic/fleet/live
(tek abonelik, iki harita)"| UI ``` ### 关键流程:从操作员到地图 模拟器是车队的**唯一位置来源**。因此,从操作员地图进行的移动不是虚假的“地图作弊”,而是通过真实管线传输的真实遥测 —— 并且会在覆盖时立即发布(无需等待下一个 tick): ``` Sağ haritada çift tık → Gateway (proxy) → Simülatör (override + anında yayın) ↓ Ingestion → Kafka ↓ Sol harita ← STOMP delta ← Gateway ← Processing ~0.1 sn ``` ## 模块 | 模块 | 端口 | 职责 | |---|---|---| | `vts-common` | — | 事件模型、topic 常量、枚举、TenantContext、通用 Kafka 消费者支持(反序列化 + retry/DLQ 策略) | | `vts-simulator` | 8085 | 车队模拟器 (Virtual Threads)、OSRM 道路路线、位置覆盖 API(不提供 UI) | | `vts-ingestion-service` | 8081 | 无状态 HTTP 入口;imei→vehicle 查找 (Caffeine→Redis→DB)、Kafka 发布、DLQ | | `vts-processing-service` | 8082 | 批量消费者;JDBC 批量插入、无状态规则 **+ 违规冷却**、outbox | | `vts-stream-analytics` | 8083 | Kafka Streams;有状态规则(急刹车、持续超速、怠速、geofence、trip) | | `vts-notification-service` | 8084 | 策略发送器、冷却、免打扰时段 | | `vts-api-gateway` | 8080 | JWT 安全、REST、STOMP WebSocket、schema 所有者 (Flyway) **+ 单页 UI(两张地图)** + 操作员控制代理 | | `vts-scheduler-service` | 8086 | ShedLock 任务:离线检测、评分、维护、outbox 发布者 | ## 车队模型 ### 覆盖土耳其的分布 100 辆地面车辆根据人口加权分布到**所有 81 个省份**(+ 5 架直升机 = 105): | 省份分组 | 车辆 | 示例 | |---|---|---| | 最大的 3 个大都市 | 各 3 辆 | 伊斯坦布尔、安卡拉、伊兹密尔 | | 接下来的 13 个大都市 | 各 2 辆 | 布尔萨、安塔利亚、阿达纳、科尼亚、加济安泰普… | | 其余 65 个省份 | 各 1 辆 | 锡诺普、阿尔达汉、亚洛瓦… | ### 类别与车辆类型 车队使用**分类法**进行建模(`vehicle_category` + `vehicle_type` 表): | 类别 | 类型 | 数量 | |---|---|---| | **地面车辆** (`LAND`) | 汽车、卡车、摩托车 | 50 / 30 / 20 | | **空中车辆** (`AIR`) | 直升机 | 5 | | **海上车辆** (`SEA`) | — | 0 | `SEA` 故意留空:分类法旨在即使还没有一艘船,也能表明可以表示一个海上舰队。`GET /api/v1/vehicle-types` 返回结果包括类别为空的项。 **类别**是决定运动和规则适用性的核心轴:地面车辆受道路网络限制,而空中车辆则不受限。代码检查的是类别而不是类型 —— “这是直升机吗”这个问题,在两个服务中产生了两个不同的豁免列表。它还包含车牌类型:`VTS-001-Otomobil`、`VTS-027-Tır`、`VTS-063-Motor`。 ### 哪条规则适用于哪种类型?(`rule_assignment` → `VEHICLE_TYPE`) 违规**仅适用于地面车辆**;例外是燃料/电池规则,它们属于机器而不是道路。适用性和阈值数据都可以在无需更改代码的情况下进行调整: | 规则 | 汽车 | 卡车 | 摩托车 | 直升机 | |---|---|---|---|---| | `SPEED_LIMIT` / `SUSTAINED_SPEEDING` | 110 | 80 | 90 | **无** | | `HARSH_BRAKING` (减速度差值) | −40 | **−30** | **−50** | **无** | | `IDLING`, `GEOFENCE_ENTER/EXIT` | ✓ | ✓ | ✓ | **无** | | `LOW_FUEL` | 15 | 20 | 10 | **25** | | `LOW_BATTERY` | 20 | 20 | 20 | 20 | 对于急刹车,*较小*的 magnitude 是更严格的设置:满载的卡车如果在一次测量中损失 30 km/s,这属于剧烈刹车,而摩托车通常都会达到 50。通用的 −40 阈值会导致对卡车发现过晚,而对摩托车频繁标记。 如果表中没有某种类型的行,则应用规则的默认值;该表仅记录**差异**。如果 `enabled = false`,则该规则根本不适用于该类型。 启动时,地面车辆会被分配 **0–120 km/s 的随机基准巡航速度**;超过限制的车辆会产生违规(受下文提到的冷却限制)。 ### 直升机(车牌 101–105) 5 架直升机因为是飞行的,所以行为与地面车辆不同: - **直线飞行、高速** (180–260 km/s);路线不是 OSRM,而是直线航线 —— 直升机的真实路线本来就是直线的。 - **不适用任何基于道路的规则**:超速限制、持续超速、急刹车、geofence 和怠速。这些规则描述的都是车辆与道路或地面区域的关系;直升机没有这种关系。豁免现在集中在一个地方 —— `rule_assignment` 的 `VEHICLE_TYPE` 行 —— 且两个规则引擎都读取相同的行。 - **燃料和电池规则有效**(燃料阈值为 25,会更早发出警告):它们描述的是机器而不是道路,而燃料耗尽的直升机是最紧急的情况。 - **Trip 规则有效** —— 一次飞行也是一个 trip,而 trip 不是违规。 - 操作员可以将直升机放置在**任何地方**(海上、建筑物顶部);而地面车辆只能移动到道路上,点击道路外会被拒绝。 ### 行程:目的地、真实路径、剩余公里数 车辆不会漫无目的地乱转;它们进行**真实的行程**: 1. 从附近的省份选择一个**目的地**。 2. 从 OSRM 获取通往该目的地的**真实驾驶路线**(道路几何形状)。 3. 车辆沿路线行驶;**剩余公里数**根据真实道路长度计算(不是直线)。 4. 到达时 **速度降为 0 → 停靠 (5.5–12 分钟)** → trip 关闭 → 选择新的目的地。 这是一个唤醒系统另一半的设计决定:如果车辆不停下,**trip 就永远不会关闭**,那么 `trip`、`trip_point`、`stop_event`、驾驶员评分和维护数据就会*永久为空*。车辆在路线上以**随机进度**开始 —— 否则首次到达(和首次 trip)要在几个小时后才会出现。 ### 路线引擎是本地的 (OSRM) `docker compose` 会启动一个**本地 OSRM**(`osrm` 服务,包含土耳其的 OSM 数据)。首次启动非常耗时 —— 需要下载约 400MB 数据并进行几分钟的预处理 —— 但结果会保存在一个 volume 中,随后的启动会瞬间完成。系统因此独立于互联网运行。 模拟器以前连接到 OSRM 的**公共演示服务器**,当它变慢/无响应时,地面车辆会被赋予 **A 到 B 的直线**:车辆依然在行驶,依然会到达,trip 依然会关闭 —— 但是穿越了土耳其的山脉、湖泊和农田。因此,车辆穿过地形的原因不是错误,而是有意的“继续移动”的选择。 现在,**地面车辆如果没有道路路线就不会移动**。如果无法获取路线,车辆将停靠,分配器会每隔几秒钟重试一次。停止的车辆是路线引擎损坏的可见且诚实的标志;而不是飞越湖泊的车辆。生成合成路线的代码已被完全删除 —— 现在已经没有任何代码路径可以为地面车辆提供偏离道路的几何形状了。 ### 驾驶员评分 `driver_score_daily`:填充 30 天的历史数据(带有每位驾驶员持久的“性格”偏好 —— 否则每个人的平均分都一样,排名就失去了意义)。日常处理则根据**真实数据**计算评分:trip 距离 + 所有违规类型,罚款**按每 100 公里进行标准化**。如果不进行标准化,行驶时间长的驾驶员会自动显得“最差” —— 这是一个破坏评分表可靠性的经典错误。 ## 规模限制(从一开始就正确确立的决定) 1. **遥测数据绝不通过单独的 `save()` 写入** —— 批量 Kafka 消费者 + `JdbcTemplate.batchUpdate()` + `ON CONFLICT DO NOTHING`;没有用于遥测的 JPA entity;使用 `reWriteBatchedInserts=true`。 2. **没有每个事件一次的 Redis round-trip** —— 有状态的状态使用 Kafka Streams state store (RocksDB);批量 Redis 操作使用 pipeline。 3. **没有每个事件一次的 WebSocket 消息** —— gateway 在内存中保存状态,并使用 `@Scheduled(1s)` 仅发布更改的内容(delta);client 发送视口 (bbox)。 4. **违规不会针对每次读取产生** —— 违数量应与**离散事件**成正比,而不是与遥测速率成正比。在三个地方进行了 debounce 处理: - *无状态规则* (`SPEED_LIMIT`, `LOW_BATTERY`, `LOW_FUEL`):在 `RuleEngine` 中采用基于车辆+规则的冷却机制 —— 持续超速的车辆会在规则的 `cooldown_seconds` 时间窗口内产生一次违规。 - *急刹车*:在 `HarshBrakingRule` 中采用基于车辆的 120 秒冷却 (RocksDB state store)。 - *持续超速*:5 分钟的 **tumbling** 窗口。(以前是 `advanceBy(1 分钟)` 的 hopping 窗口;任何持续超速的车辆都会**每分钟**产生一次违规,仅此一项就占了数据洪流的 **83%**。) 测得的整体影响:**~20 违规/秒 → ~0.2 违规/秒**。 5. **内存不设无限** —— 如果没有容器限制,JVM 会根据*宿主机* RAM(%25)选择堆上限(heap ceiling);8 个服务相乘会导致 RAM 随着时间推移而膨胀。通过每个服务设置 `mem_limit` + `MaxRAMPercentage` 来固定堆。此外,由于 Kafka Streams 中的 **5 个 store × 24 个 partition ≈ 120 个 RocksDB 实例**各自会开启自己的缓存,因此它们全部共享一个有限且统一的 LRU 缓存 + write-buffer manager (`BoundedRocksDBConfig`)。Java 服务总计:**~2.4 GB,硬性上限 3.7 GB**。 6. **Trip 距离不是原始 GPS 的总和** —— 连续读取之间超过 2 公里的步进会被视为*重新定位*(操作员传送、重启、GPS 跳跃),不计入距离。如果计入,trip 长度以及随之而来的所有按公里标准化的驾驶员评分都会悄无声息地作废。 7. **Kafka partition = 24**(与 profile 无关) —— 以后增加会破坏按车辆的排序 (per-vehicle ordering) 和 Streams state store。 8. **遥测 = TimescaleDB hypertable** —— dashboard 查询指向 continuous aggregate 而不是原始表。 9. 每个表从一开始就带有 `tenant_id` + Outbox Pattern。 10. **实时流不能匿名监听** —— 在 STOMP `CONNECT` frame 中验证 JWT。由于无法在 SockJS 握手期间设置 `Authorization` header,因此握手保持公开;身份在 CONNECT 时检查。否则,任何人都可以在没有 token 的情况下观看整个车队。 ## 数据模型 Flyway `V1`–`V15` 定义了 **38 张业务表**。主要包括: - `telemetry` **hypertable**:`by_range(ts)` + `by_hash(vehicle_id, 8)`,PK `(vehicle_id, ts)`,无 FK(保证批量插入速度)。 - `violation` **hypertable**。 - `vehicle_driver_assignment`:用于将违规准确归咎于当时驾驶员的时间记录。 - `rule` + `rule_assignment`:阈值绝不在代码中;在 TENANT/GROUP 范围内被覆盖(基于类型的限速即来源于此)。 - Continuous aggregate:`telemetry_1min`、`telemetry_hourly`、`violation_daily_summary`。 - 压缩 + retention 策略,GIST/BRIN/partial index。 ## 运行 整个系统只需一条命令即可启动(基础设施 + 8 个服务): ``` docker compose up -d --build ``` 然后在浏览器中访问: | 界面 | 地址 | 登录 | |---|---|---| | **VTS — 单页(车队栏 + 两张地图)** | http://localhost:8080 | `admin` / `password` | | Swagger UI | http://localhost:8080/swagger-ui.html | JWT | | Kafka UI | http://localhost:8090 | — | | Prometheus | http://localhost:9090 | — | | Grafana | http://localhost:3000 | `admin` / `admin` | 其他端口:ingestion 8081,processing 8082,stream-analytics 8083,notification 8084,scheduler 8086,**Postgres 5433**,Redis 6379。 负载 profile(1000 辆车 / 1 秒,3 个 broker 覆盖): ``` docker compose -f docker-compose.yml -f docker-compose.load.yml up -d ``` ### API 示例 ``` # 登录(开发用户:admin / password) TOKEN=$(curl -s -X POST localhost:8080/api/v1/auth/login \ -H 'Content-Type: application/json' \ -d '{"username":"admin","password":"password"}' | jq -r .token) curl localhost:8080/api/v1/vehicles -H "Authorization: Bearer $TOKEN" curl localhost:8080/api/v1/live/positions -H "Authorization: Bearer $TOKEN" curl localhost:8080/api/v1/dashboard/summary -H "Authorization: Bearer $TOKEN" curl "localhost:8080/api/v1/violations?limit=20" -H "Authorization: Bearer $TOKEN" ``` 实时地图 WebSocket (STOMP):`ws://localhost:8080/ws` → `/topic/fleet/live`、`/topic/violations`、`/user/queue/notifications`; 针对视口的 `/app/viewport`。 ### 操作员控制 API (gateway 代理 — 通过 `vehicleId`) ``` # 工具的目标/路线状态 curl localhost:8080/api/v1/control/state -H "Authorization: Bearer $TOKEN" # 移动工具(使用 vehicleId — 而非车牌号) curl -X POST localhost:8080/api/v1/control/49/position -H "Authorization: Bearer $TOKEN" \ -H 'Content-Type: application/json' -d '{"lat":39.92,"lon":32.85}' ``` 移动是一种**锁定,而不是传送**:车辆降落在放置点,并从那里以自己的速度出发**前往新的目的地**。目的地是根据车辆的*新*位置来选择的 —— 保留旧目的地会把移动变成一个自我取消的操作:一辆从哈塔伊移动到安卡拉的卡车仍然需要前往开塞利,并且会立即掉头向东南方向返回。 响应通过 `moved` 字段告知结果 —— 点击道路外会返回 200 状态码但 `moved:false`,因为请求已被理解,**拒绝即是回答本身**: ``` {"found":true,"flying":false,"moved":false,"reason":"OFF_ROAD","offRoadMeters":10576} ``` 地面车辆要将点击视为“在道路上”,最近的道路必须在 `vts.simulator.road-click-tolerance-meters`(默认为 **50 m**)范围内。这个阈值是对“点击在道路上”的定义:如果太窄,视觉上对准道路的点击会被拒绝(点击包含了缩放误差、道路绘制粗细和 OSM 几何误差),如果太宽,车辆会悄无声息地停在操作员未指示的地方。 Gateway 通过 imei 将 `vehicleId` 转换为模拟器的设备索引;模拟器的 `:8085/api/**` 端点保留在内网中,UI 永远不会直接访问它。 ## Profiles | Profile | 车辆 | Tick | Chunk | Retention | |---|---|---|---|---| | `dev` (默认) | 100 (分布到 81 个省份) | 1 秒 | 1 天 | 30 天 | | `load` | 1000 | 1 秒 | 1 小时 | 7 天 | | `prod` | 外部配置 | — | — | — | ## 测试 - **Kafka Streams:** `TopologyTestDriver`(急刹车、怠速、geofence 进出、trip)。 - **单元测试:** 规则引擎 **+ 违规冷却窗口**、ingestion 路由、notification 冷却/免打扰时段、JWT、实时地图 delta+viewport、模拟器运动和速度模型。 - **Schema 验证:** JPA entity 针对运行中的 TimescaleDB 进行 `ddl-auto=validate` (Testcontainers)。 - **端到端:** 在真实容器中验证了 simulator → Kafka → DB 数据流(遥测、最后位置、违规、基于类型的阈值覆盖、驾驶员归属、操作员覆盖在主地图上的反映)。 ``` mvn test ``` ## 可观测性 Micrometer + Prometheus + Grafana。指标:`telemetry.ingested`、`telemetry.persisted`、`violation.produced`、`notification.sent`、consumer lag、DLQ 比率。Grafana 会自动加载一个预置的 **"VTS — Fleet Telematics Overview"** dashboard。每个事件都带有 `correlationId` 的结构化 JSON 日志。
标签:Kafka, SonarQube插件, Spring Boot, TimescaleDB, 事件驱动架构, 地理信息系统, 域名枚举, 搜索引擎查询, 请求拦截, 车联网, 车队管理