ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java公交车实时监控系统:GPS接入、Kafka缓冲与WebSocket推送实战

Java公交车实时监控系统:GPS接入、Kafka缓冲与WebSocket推送实战 简介这是一套基于Java打造的公交车实时监控系统设计源码面向具有一定Java基础、希望学习企业级实时监控系统开发的开发者可支撑课程设计、项目实训或毕业设计。系统采用微服务架构基于Spring Boot封装各业务模块并通过RESTful API与笑园实时公交API对接实现车辆位置、运行轨迹、车辆状态等数据的实时获取与展示。资源包共43个文件约61KB其中包含37个Java源代码文件作为核心业务逻辑辅以2个XML配置文件用于组件灵活配置1个YAML文件简化配置流程1个SQL文件支撑监控数据持久化存储另有gitignore及说明文档。当前已有279人学习下载。从源码中可梳理出数据接口调用的标准化流程、后端模块划分与配置管理思路有助于理解实时监控系统的工程化实现适合作为智能交通方向的项目参考。1. 用 Java 把公交车的实时位置搬到浏览器里这套源码解决的不只是画点调度室里的大屏要实时显示每一辆公交车的当前位置这个需求听起来简单真正动手做的时候会发现它牵扯到 GPS 数据接入、消息缓冲、实时推送、前端地图渲染四块内容。我拆的这套基于 Java 实现的公交车实时监控系统设计源码就是把“上报 → 缓冲 → 推送 → 展示”这条链路完整打通并且每一层都预留了替换的余地。适合做 Java 课程设计、毕业设计或者公司内部需要一套可演示的调度原型系统的开发者。动手之前先记住一个结论这个项目最耗时间的不是地图上画点而是数据对齐和时间戳处理源码的价值恰恰是把这些坑提前踩平。2. 系统架构与数据模型先把车辆位置从报文到浏览器走一遍2.1 五个环节组成的链路每一环都能单独拆开替换拿到这套源码第一件事不是看代码而是看它的数据流向。整个系统从模拟器到前端地图一共经过五个环节理解这一条链路之后后面改代码才有明确的下手点。车辆 GPS 模拟器生成经纬度、速度、方向角封装成一条 JSON 消息上报给接入接口接入接口负责验签和基本格式校验随后把消息投递到 Kafka下游的位置消费服务从 Kafka 拉取消息一方面把最新位置写入 Redis 做缓存另一方面通过 WebSocket 推送给浏览器浏览器端的地图组件收到 JSON 后解析经纬度更新车辆图标。全程没有数据库参与实时链路MySQL 只用来存车辆基础信息和历史轨迹这很关键。中间的每个环节都可以替换模拟器可以换成真实车载终端接入接口可以直接对接 TCP 长连接Kafka 在单机演示时可以用内存队列替代。源码里用配置开关控制这条数据链路我一般拿到项目先搜KafkaListener和ServerEndpoint这两个注解就能快速定位核心处理逻辑。各层职责按模块拆开看会更清楚下面这张表可以直接对照源码包结构去翻。环节核心组件职责替换场景数据源GpsSimulator生成并上报车辆位置换成真实终端报文接入层GpsController校验消息、投递到消息队列增加鉴权或限流缓冲层Kafka Topicbus-gps削峰填谷、解耦上下游单机换内存队列消费推送PositionConsumer / WebSocket更新缓存、推送到浏览器增加多路推送策略展示层网页地图组件渲染车辆标点、运行状态换成小程序或客户端2.2 车辆与位置表设计字段冗余少一点查询才不至于绕数据库在这套系统里承担的是辅助角色但表结构设计直接影响后面做历史轨迹回放时的工作量。源码里核心就是两张表车辆基础信息表和位置记录表。车辆表保存车牌号、线路编号、车辆状态位置表保存每一次上报的经纬度、速度、方向角和 GPS 设备时间。位置表的设计有几个值得注意的决策。经纬度字段用DECIMAL(10,6)这个精度大约是 0.1 米足够公交定位使用速度用SMALLINT保存整数 km/h没有用浮点方向角用SMALLINT存 0 到 359 的整数。GPS 时间单独用一个DATETIME字段而不是直接用数据库的CURRENT_TIMESTAMP这个细节在实际运行中很重要——报文到达时间和设备采集时间之间存在网络延迟后续做轨迹分析要以设备时间为准。CREATE TABLE vehicle ( id BIGINT PRIMARY KEY AUTO_INCREMENT, plate_no VARCHAR(20) NOT NULL UNIQUE COMMENT 车牌号, line_no VARCHAR(10) NOT NULL COMMENT 线路编号, status TINYINT DEFAULT 1 COMMENT 1运营 0停运, create_time DATETIME DEFAULT CURRENT_TIMESTAMP ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT公交车基础信息; CREATE TABLE bus_position ( id BIGINT PRIMARY KEY AUTO_INCREMENT, plate_no VARCHAR(20) NOT NULL COMMENT 车牌号, lng DECIMAL(10,6) NOT NULL COMMENT 经度, lat DECIMAL(10,6) NOT NULL COMMENT 纬度, speed SMALLINT DEFAULT 0 COMMENT 速度km/h, direction SMALLINT DEFAULT 0 COMMENT 方位角0-359, gps_time DATETIME NOT NULL COMMENT GPS设备时间, KEY idx_plate_time (plate_no, gps_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT车辆位置上报记录;注意idx_plate_time这个联合索引它是给历史轨迹查询用的。如果你直接按plate_no查全部记录车辆跑一天就会产生几万条数据没有索引的查询会直接拖垮页面有了联合索引之后按时间和车牌过滤基本在毫秒级完成。我一般建议把这两张表的数据保留策略做得简单一些实时表只保留最近三天历史轨迹归档到另一张分区表即可。2.3 为什么用 Kafka 做缓冲而不是直接写数据库这是整套系统里最值得在面试或答辩时展开的设计点。一台公交车平均 3 到 5 秒上报一条位置看似频率不高但一辆运行 14 小时的公交车一天就要产生 1 万条左右的位置记录一个城市按 200 辆车计算一天就是 200 万条写入。如果每次上报直接 INSERT 到 MySQL数据库会持续处于高负载状态而且前端推送和数据库写入互相争抢连接资源时实时性会被拉垮。Kafka 在这里承担的是削峰和异步化的角色。接入接口收到位置消息后只做校验和序列化投递到 Kafka 就立即返回下游的 PositionConsumer 按照自己的节奏消费消息批量更新 Redis、批量推送、按批次落库。这样数据库不再直接面对高频上报而是接收集中批量写入压力完全可控。这套系统里 Kafka 主题只保存最近的消费进度不需要长时间保留数据所以主题的retention.ms我建议直接设成3600000也就是一小时。如果你跑单机演示觉得 Kafka 太重源码里有开关可以直接切换到内存队列模式这就是链路解耦带来的优势——每一层替换起来都不影响其他模块运行。3. 数据采集与模拟手写一个 GPS 模拟器的正确姿势3.1 模拟器关注三个参数路线、车速、上报节奏没有真实车载终端时模拟器就需要生成足够可信的位置消息让整个系统转起来。源码里的模拟器核心就三件事按预设路线移动车辆、随机调整速度和方向、按上报间隔发送消息。初次看代码的时候我建议先把这三个参数找出来别急着跑工程。路线本质上是一组经纬度坐标点数组模拟器在相邻两个坐标点之间做线性插值。车速决定点位前进的步长比如 30 km/h 比 50 km/h 每一步移动的距离更短。上报节奏决定消息产生的频率源码默认 1000 毫秒上报一次如果单机演示想看得更连贯调整到 500 毫秒也没问题但要注意下游消费端的处理速度。很多人写模拟器容易犯一个错误把路线当成直线来跑。真实公交是在站点之间走折线的直线插值会让车从楼顶上穿过去。源码的路线表里每个点就是一个站台或拐点站点之间用线性插值拐点处切换目标点这样跑出来的轨迹在地图上才像一条正常的公交线路。3.2 模拟器核心代码插值移动和随机扰动下面这段代码是从模拟器里抽出来的核心逻辑它的作用就是让一辆车沿着预设路线持续移动并叠加随机扰动让数据看着更真实。public class GpsSimulator implements Runnable { private final String plateNo; private final ListPoint route; private final Random random new Random(); private int segmentIndex 0; private double progress 0.0; Override public void run() { Point start route.get(segmentIndex); Point next route.get(segmentIndex 1); // 每次移动 2% 到 3% 的区间速度由这个步长体现 progress 0.02 random.nextDouble() * 0.01; if (progress 1.0) { segmentIndex (segmentIndex 1) % (route.size() - 1); progress 0.0; } double lng start.lng (next.lng - start.lng) * progress; double lat start.lat (next.lat - start.lat) * progress; GpsMessage msg new GpsMessage(); msg.setPlateNo(plateNo); msg.setLng(lng); msg.setLat(lat); msg.setSpeed(20 random.nextInt(25)); // 20-45km/h 之间随机 msg.setDirection(computeDirection(start, next)); // 按线段方向计算方位角 msg.setGpsTime(System.currentTimeMillis()); sendToKafka(msg); // 实际生产环境通过 KafkaTemplate 发送 } }这里的核心逻辑是progress这个变量它代表车辆在当前线段中走过的比例。每次run()增加 2% 到 3%步长随机浮动车速因此有了变化。计算新坐标时用start (next - start) * progress做线性插值这就是最简单也最不容易出错的轨迹生成方式。方向角通过当前线段两个点的经纬度差值计算用Math.atan2转成角度即可。上报时间直接用System.currentTimeMillis()是模拟器可以接受的因为模拟器没有真实设备时钟。但在真实终端场景里这里必须换成设备采集时间我在第 5 章会详细说明原因。3.3 让 200 台车同时动起来调度与线程池模拟器类本身只负责生成一条位置消息要让多辆车同时运行需要一个调度层来控制节奏。源码里用一个线程池来管理所有模拟器任务每个任务代表一辆车。public class SimulatorScheduler { private final ExecutorService executor Executors.newFixedThreadPool(30); public void start(ListString plateNos, ListListPoint routes) { for (int i 0; i plateNos.size(); i) { GpsSimulator simulator new GpsSimulator(plateNos.get(i), routes.get(i)); // 每辆车按自己的频率上报不需要全局同步 executor.submit(simulator); } } }线程池设了 30 个线程而不是车辆数。原因很简单模拟器的任务并不计算密集每个任务只是循环执行插值和消息发送大部分时间在等待网络 I/O30 个线程处理 200 辆车完全足够。submit()之后模拟器会在内部用while(running)循环持续上报每轮循环之间通过Thread.sleep(1000)控制节奏。有个细节值得注意每辆车都在自己的循环里上报所以不会出现所有车辆在同一时刻集中发送消息的情况。这样更接近真实场景消息到达接入层的时间是错开的Kafka 的写入压力也是平稳的。如果你想提高模拟强度把 sleep 间隔缩短到 200 毫秒就能模拟高频上报场景顺便测试下游的吞吐上限。4. 实时推送链路WebSocket 与 Redis 的配合4.1 为什么实时推送选 WebSocket 而不是轮询浏览器要实时显示车辆位置最简单的做法是前端每 5 秒调一次接口拉最新位置这种轮询方式在数据量小的时候还算凑合但一旦车辆超过 100 辆每秒要重复拉取大量不变的数据浪费带宽也浪费数据库连接。这套系统选 WebSocket 的原因很直接位置信息是服务端主动产生的数据流客户端只需要在连接建立后被动接收实时性可以压到 1 秒以内而且连接数量可控。轮询的延迟是固定间隔决定的间隔设得越短越费资源WebSocket 是事件驱动数据一到就推没有冗余请求。另外一个隐藏优势是 WebSocket 链接可以复用不用每次建立新的 HTTP 连接这对持续推送的场景非常重要。4.2 WebSocket 服务端实现连接管理、广播、心跳Component ServerEndpoint(/ws/bus) public class BusWebSocket { // 线程安全的连接集合避免并发遍历问题 private static final CopyOnWriteArraySetSession SESSIONS new CopyOnWriteArraySet(); OnOpen public void onOpen(Session session) { SESSIONS.add(session); } OnClose public void onClose(Session session) { SESSIONS.remove(session); } OnError public void onError(Session session, Throwable error) { SESSIONS.remove(session); } /** * 广播位置消息所有订阅了车辆位置的页面都会收到 */ public static void broadcast(String json) { for (Session session : SESSIONS) { if (session.isOpen()) { session.getBasicRemote().sendText(json); } } } }这段代码最关键的两个点连接集合用了CopyOnWriteArraySet它保证并发遍历时不抛ConcurrentModificationException代价是写操作开销高但连接建立和断开本身不频繁这个取舍划算broadcast是静态方法因为消费者线程从 Kafka 拿到消息时没有 HttpSession 上下文必须通过静态方法直接调用这也是很多教程里容易省略的细节。代码里对每个 session 直接sendText是同步调用。200 辆车每辆每秒一条消息每秒就有 200 次广播每个连接都要写一次如果页面连接数到了 200 个单线程循环阻塞就会明显拖慢消费速度。所以真正的高并发场景里需要把同步 send 改成异步发送或者引入多线程推送但那是后话先用同步方案把链路跑通最重要。4.3 Redis 在链路里的角色车辆信息预热与最新位置缓存Redis 在这套系统里做两件事缓存最新位置供页面回显和查询使用缓存车辆基础信息避免下游每次消费消息都去查 MySQL。车辆信息变化极少启动时一次性预热到 Redis 是常规做法。Component public class CacheService { private final StringRedisTemplate redisTemplate; public void preloadVehicles() { // 从数据库查出所有车辆基础信息后写入 Redis Hash ListVehicle vehicles vehicleMapper.selectList(null); for (Vehicle v : vehicles) { redisTemplate.opsForHash().put(vehicle:info, v.getPlateNo(), JSON.toJSONString(v)); } } public void updatePosition(GpsMessage msg) { // 用 String 结构存最新位置过期时间 10 分钟 String key bus:pos: msg.getPlateNo(); redisTemplate.opsForValue().set(key, JSON.toJSONString(msg), 10, TimeUnit.MINUTES); } }preloadVehicles在系统启动后执行一次后续车辆基本信息都从 Redis 的 Hash 结构里取查询复杂度是 O(1)数据库压力几乎为零。这部分的另一个用途是给新接入的页面做初始化——当页面刷新时WebSocket 只推送新消息初始的位置数据就得靠 Redis 兜底否则页面打开的瞬间地图上全是空点。这个场景很常见我一般建议在浏览器页面加载完成后先调用一次 REST 接口从 Redis 拉全量最新位置再建立 WebSocket 订阅增量更新这样前端体验最好。5. 避坑与排查把这套系统真正跑起来要过的五关5.1 车辆漂移坐标系混用是头号元凶现象车辆轨迹在地图上正常移动但整体位置偏移车开到河里或者道路几十米开外的地方。原因模拟器生成的是 WGS-84 坐标系经纬度前端地图用的是 GCJ-02 坐标系两种坐标系在城市区域会有几十米的偏差看起来就是“漂移”。很多刚入手的人以为是硬件问题折腾半天发现是坐标没转换。解决统一在后端做坐标转换模拟器输出原坐标接入层或消费层拿到消息后调一次坐标转换工具转成 GCJ-02 再推给前端。转换代码用标准算法就行上百行网上一搜一大片。重点是把转换放在后端前端不再做任何坐标处理。5.2 Kafka 消费积压广播阻塞了整个消费线程现象模拟器正常运行Kafka 控制台显示生产者消息数量在涨但消费者只消费了一小部分积压越来越严重页面上的车辆位置最新更新时间明显滞后。原因消费者线程同步调broadcast()而broadcast()内部逐条 session 同步发送当 WebSocket 连接数较多或者某个浏览器页面卡顿时发送会阻塞消费者被卡住无法继续拉取消息积压由此产生。解决把 WebSocket 推送从 Kafka 消费线程中拆出去。消费消息后先更新 Redis再交给一个异步推送执行器去广播消费线程本身立刻返回继续拉取下一批。我在源码里把异步推送改成ExecutorService配合submit()消费者不再直接调broadcast()。实际上手时可以先看 Kafka Lag 指标确认是不是这个原因再决定优化方向。5.3 sendText 抛异常导致消费者线程挂起现象系统运行几分钟后页面上所有车辆全部卡死不再更新重启后恢复正常过几分钟又卡死。原因某个浏览器页面被强制关闭WebSocket 连接断开了但服务器端的session还没有从连接集合里移除。下一次循环广播时对失效 session 调用sendText抛出IOException异常没有拦截导致整个线程中断后续所有车辆的推送全部停止。解决给sendText调用包一层 try-catch捕获到IOException就立即把 session 从SESSIONS集合中移除同时关闭连接释放资源。这个异常处理看起来很简单但它直接决定系统的可用性我见过不少项目忽略这个细节导致推送服务不稳定。5.4 时间戳用错用后端时间给 GPS 报文盖章现象车辆在隧道里行驶时WebSocket 推送的消息到达时间正常但前端用时间排序后车辆位置看起来在倒退轨迹线来回穿插行驶里程统计偏大。原因模拟器或接入层在数据落库时用了System.currentTimeMillis()作为最终时间而这个时间是服务端接收调度的时间不是 GPS 设备采集的时间。当消息因为网络波动延迟到达时服务端时间戳会比设备时间晚前端按时间排序就把顺序打乱了。解决消息体里始终保留设备上报的gpsTime所有排序和轨迹回放都用它。服务端时间只在写入日志或计算处理延迟时使用绝不覆盖设备时间。这个教训之后我一直把“时间戳只能由数据源产生”当成铁律。5.5 MySQL 批量插入没提速先查数据包大小现象把单条 INSERT 改成批量 INSERT 后写入速度不仅没提升反而频繁报错日志里出现PacketTooBigException。原因MySQL 默认的max_allowed_packet是 4MB位置记录表字段较多一次批量插入 2000 条的消息体很容易超过这个限制导致整批失败。解决后端批量插入按 500 条一批切割或者执行前先算一下 JSON 报文大小。我通常会顺手把 MySQL 的max_allowed_packet调整到 16MB给后续扩展留余量。两个方案都做效果最好这个参数在配置文件中改一行就能生效。6. 端到端验证与压测一条命令确认链路还活着6.1 写一个最小 WebSocket 客户端做冒烟测试拿到源码跑起来之后我习惯先不打开浏览器而是用一个最小客户端直接连 WebSocket 端点确认消息真的能推出来。这个方法比打开地图页面更快更可靠也方便在服务器上排查问题。public class WsSmokeTest { public static void main(String[] args) throws Exception { // 本地起一个 WebSocket 连接不依赖浏览器 WebSocketClient client new WebSocketClient(new URI(ws://localhost:8080/ws/bus)) { Override public void onMessage(String message) { // 解析经纬度字段打印结构确认数据完整 GpsMessage msg JSON.parseObject(message, GpsMessage.class); System.out.printf(车辆 %s 位置 %.6f,%.6f 速度 %d km/h%n, msg.getPlateNo(), msg.getLng(), msg.getLat(), msg.getSpeed()); } }; client.connect(); } }这个测试类不用集成任何测试框架main 方法就能跑。拿到消息后打印车牌、经纬度和速度确认三件事消息能持续到达、JSON 字段完整、推送间隔正常。如果客户端连不上优先检查 WebSocket 端点路径和端口别先怀疑模拟器。6.2 用“消费积压量”判断系统健康度冒烟测试通过之后还有一个更长期的验证指标值得关注Kafka 消费者积压量Lag。这个数值代表消费者还没处理的消息条数理想情况下应该接近零。如果积压量持续增加说明消费者处理速度跟不上生产者速度这个指标能提前暴露问题比等到页面卡死再排查省事得多。我经常在测试环境用一段脚本每 5 秒打印一次 Lag 值观察它是否稳定。启动模拟器后积压量先飙升再下降说明缓存接收正常如果积压量一直不降就得按第 5 章的顺序排查广播阻塞、异常捕获和时间戳问题。用这种方式压了一个晚上积压量始终在个位数徘徊这套系统的可靠性基本就有底了。从那以后我每次做实时系统都会强制先跑一遍 WebSocket 冒烟测试和 Lag 监控确认链路通畅再考虑功能扩展。处理完这五类坑这套源码基本就能稳定运行在单机环境里了换数据源、加地图适配、调整推送频率这些扩展点也都在链路里预留好了接口位置。希望这个拆解过程能帮到你少走我踩过的弯路。本文还有配套的精品资源点击获取
返回列表