ARTICLE DETAIL

资讯详情

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

Flink集成Spring Boot实现AGV实时监控:Windows本地实践指南

Flink集成Spring Boot实现AGV实时监控:Windows本地实践指南 1. 为什么要把 Flink 和 Spring Boot 绑在一起AGV 监控的真实痛点先聊点实际的。我见过不少做 AGV自动导引车项目的团队早期监控系统大多是一个 Spring Boot 服务定时去查数据库或者前端页面每隔几秒轮询一次接口。小车少的时候还好说一旦上了二三十台车位置、电量、任务状态、告警事件一堆数据同时涌过来轮询方案就开始露馅了——页面刷新慢、告警延迟、数据库压力大更别提要做什么低电量预测、路径拥堵分析这类稍微复杂一点的计算写起来又别扭又难维护。AGV 实时监控的本质是一个典型的流式数据处理场景。小车每秒钟都在上报心跳数据包含坐标、速度、电量、运行状态还可能带一些 IO 信号、避障传感器状态。这些数据明显不是查一次就完事的静态数据而是持续产生、需要连续处理的流。用传统定时任务去循环拉取本质上是在用批处理的思路做流处理的事架构上就别扭。那为什么偏偏是 Flink 加 Spring Boot 这个组合Flink 负责的是流处理的重活持续接收 AGV 心跳流做解析、过滤、窗口聚合、状态管理、规则告警。它是真正的实时计算引擎能保证毫秒级处理延迟还自带 exactly-once 之类的状态一致性保证。而 Spring Boot 负责的是传统后端服务该干的活提供 REST API 给前端看板、通过 WebSocket 主动推送告警、做用户权限管理、把转换后的数据落库。两个框架各有专长拼在一起正好互补。还有一个很现实的原因现在很多 AGV 项目现场上位机就是 Windows调度系统、监控大屏都跑在 Windows 上。开发阶段更不用说大部分人手里就是一台 Windows 笔记本。在 Windows 本地把整套链路跑通意味着你可以在不依赖 Linux 服务器的情况下把所有逻辑调试验证完毕再平滑切换到生产环境。这也是我写这篇文章的初衷——在 Windows 上把 Flink 集成进 Spring Boot 项目踩过的坑、走过的弯路一次说清楚。适合谁来参考一是做 AGV、AMR 这类移动机器人项目的后端或上位机开发人员二是想在自己的 Spring Boot 工程里引入流处理能力、又不想一上来就上大数据集群的同学三是对 Flink 感兴趣但被 Hadoop 生态吓住、想先找个 Windows 环境练手的人。这篇文章不假设你有大数据集群所有东西都跑在本地 Windows 上。2. 开工前的环境准备Windows 本地跑 Flink 的隐性门槛2.1 版本选型JDK 和 Flink 的搭配讲究很多人第一步就栽在版本搭配上。Flink 对 JDK 版本有明确要求拿 Flink 1.17 来说官方支持 JDK 8、11、17但如果你用的是 JDK 17会遇到一些模块化相关的权限问题比如反射访问受限。我个人的建议是本地开发统一用 JDK 8 或 JDK 11别追求最新。Spring Boot 这边也是一样。Spring Boot 2.7.x 搭配 JDK 8/11 非常成熟而 Spring Boot 3.x 强制要求 JDK 17。如果你跟我一样选择了 JDK 11 作为基础那 Spring Boot 就用 2.7.x别硬上 3.x不然会遇到各种兼容性麻烦。这套组合在 Flink 1.17 Spring Boot 2.7 JDK 11 下是我实测最稳的搭配。Flink 本身下载很简单到官网下载页面选对应 Hadoop 版本的包本地模式其实不需要 Hadoop选一个通用的即可。解压后目录结构里有个 bin 文件夹里面是启动脚本。2.2 启动 Flink 的连环坑Windows 下启动 Flink 本地集群传统方式是双击start-cluster.bat默认启动一个 JobManager 和一个 TaskManager。但这里有几个极容易踩的坑第一环境变量 JAVA_HOME 必须配好。这个不用多说但我要补充一个细节JAVA_HOME 路径里不能有空格。很多人把 JDK 装在C:\Program Files\Java\jdk1.8.0_xxxFlink 的 bat 脚本解析这个路径时会因为空格拆分出错导致脚本闪退。解决方案是手动把 JAVA_HOME 指到一个无空格路径比如C:\Java\jdk8或者修改 Flink 的conf\flink-conf.yaml里的env.java.home指定到 JDK 安装路径。第二闪退问题的另一个常见原因是 bat 脚本执行策略或杀毒软件拦截。Windows 的 PowerShell 默认执行策略可能会阻止脚本运行这时候用 CMD 而不是 PowerShell 启动会更省事。杀毒软件把java.exe当作可疑程序拦截的情况我也遇到过解决办法就是把 Flink 目录加入杀毒软件的白名单。第三端口占用。Flink 默认的 Web UI 端口是 8081如果你本机已经有什么服务占用了 8081启动会直接失败JobManager 日志里会报Address already in use。用 CMD 跑一下netstat -ano | findstr 8081看看是谁占着然后在flink-conf.yaml里改rest.port即可。启动成功的标志很简单浏览器打开http://localhost:8081能看到 Flink 的 Web UI 页面显示一个空闲的 TaskManager就说明集群已经就绪了。这个页面非常有用后面任务跑起来之后你可以在这里看到每个任务的处理速率、反压状态和错误日志。2.3 创建 Spring Boot 工程的正确姿势Spring Boot 工程推荐用 IDEA 社区版加 Spring Initializr 插件或者直接到start.spring.io网站生成。这里我建议你手动加依赖而不是全选因为后面要用的几个依赖各有讲究spring-boot-starter-web提供 REST APIspring-boot-starter-websocket实时推送看板数据flink-streaming-javaFlink 流处理核心库注意要指定版本号flink-clients本地提交任务需要flink-connector-kafka或flink-connector-mqtt根据你的数据通道选后面细说Flink 相关依赖的 scope 有一个大坑本地测试时要把它们设为provided或显式排除掉内置的冲突依赖。因为 Flink 集群本身带了一份 Flink 库如果 Spring Boot 的可执行 Jar 里也打包了一份提交任务时就会类冲突。本地跑没问题但你一旦部署到远程集群就会报出各种NoClassDefFoundError。我习惯在开发环境把 Flink 依赖 scope 设置为provided这样打包时不会冗余本地运行时因为直接依赖 IDE 的 classpath 也没问题。这里附一个 pom.xml 里 Flink 依赖的参考写法dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.17.2/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version1.17.2/version scopeprovided/scope /dependency3. 打通第一层Flink 任务怎么把数据交给 Spring Boot3.1 先想清楚数据链路再写代码这是整个项目里最关键的设计决策。我见过太多人上来就写代码结果数据通路没理清后面全部返工。AGV 实时监控的数据流大致分成这么几段AGV 车载终端/调度系统 - 消息中间件 - Flink解析、聚合、告警 - 存储/接口 - Spring Boot - WebSocket - 前端看板这里有两处过墙的地方AGV 数据怎么进 FlinkFlink 处理完怎么交给 Spring Boot。这两处决定了整个架构的形态。AGV 数据进 Flink 这一侧现实中通常有两种方式。一种是 AGV 调度系统本身就集成了 MQTT 或 Kafka心跳数据已经进入了消息队列Flink 只需要做一个 source 消费者。另一种是调度系统没有消息队列直接把数据写到数据库或者通过 HTTP 上报。这种情况下Flink 需要拿着 JDBC 连接器去轮询数据库或者通过 HTTP source 去拉取。性能差异很大强烈建议接入了消息中间件之后再交给 Flink。Flink 处理完往 Spring Boot 传递这一侧有三种常见方案我对比一下方案链路优点缺点适用场景A. Flink - Kafka - Spring Boot消息队列中转解耦彻底Spring Boot 天然消费本地要多跑一个 Kafka吃内存生产环境推荐B. Flink - Redis - Spring Boot状态共享实时读取最新状态方便Spring Boot 只管读Redis 不是消息队列不承载大量事件流状态看板、位置跟踪C. Flink HTTP Sink - Spring Boot 接口直接调用最简单本地调试方便耦合紧Flink 重试逻辑要自己处理本地开发、小规模验证本地开发我的建议是先用方案 C 把全链路打通逻辑验证没问题之后再换方案 A。原因很简单本地再跑一个 Kafka内存和配置成本都上来了而且 Kafka 在 Windows 上的体验不算友好Zookeeper 和 Kafka 两个进程都得手动启动还要处理 Windows 下的路径和权限问题折腾一轮下来半天没了。先用 HTTP Sink 验证业务逻辑等你要模拟高并发或者上生产了再切 Kafka 不迟。3.2 Flink 里的处理完到底是什么状态很多人对 Flink 输出给 Spring Boot 的数据内容没有概念。AGV 心跳原始数据是高频且冗余的假设每台车每 500ms 上报一条几十台车就是每秒上百条消息。如果把这些全都推给 Spring Boot服务将被淹没。Flink 的价值恰恰在这一层体现过滤去掉无效、重复、乱序的数据富化把 AGV ID 关联到具体线体、区域、班组等业务维度窗口聚合把每秒的明细聚合成每秒的指标比如当前在线车辆数、各区域车辆密度、平均速度模式识别识别异常状态比如连续 3 秒无心跳判定为通信中断电量低于阈值触发充电告警只有经过这些处理之后才把该让业务层知道的事发给 Spring Boot。这就是我在标题里说的打通任督二脉的意思——Flink 是任脉持续不断的流处理Spring Boot 是督脉对外提供服务和展示两脉通了整个监控体系才真正活起来。4. 核心链路逐段拆解从 AGV 心跳到前端看板4.1 模拟 AGV 数据源没有真车怎么开发开发阶段不一定有真实 AGV 可以接入所以第一步是写一个数据模拟器。我这边用一个 Java 线程池定时生成 JSON 心跳数据通过 HTTP POST 发送到 Flink 的 source或者直接写入一个 Kafka topic如果你切到了方案 A。模拟器里的关键逻辑是让每台车的运动轨迹有一定规律比如沿直线或绕圈移动电量按时间递减这样后面做窗口聚合和告警才看得出效果。模拟的心跳 JSON 长这样{ agvId: AGV-001, x: 23.4, y: 15.8, speed: 1.2, battery: 87.5, status: RUNNING, timestamp: 1700000000000 }核心字段就六个车辆 ID、坐标、速度、电量、运行状态、时间戳。模拟器每 500ms 生成一条几十辆车同时跑带宽压力不大但足够暴露问题。4.2 Flink 端的处理逻辑DataStream API 还是 Flink SQL在 Flink 里处理这份数据有两种路径DataStream API 或者 Flink SQL。我推荐本地项目用 DataStream API因为 AGV 监控里需要做比较细的状态管理和窗口计算DataStream API 的灵活性更好而且排查问题的时候能直观地在代码里打断点。核心处理流程包含四个算子第一个算子是 source负责订阅数据。如果走 HTTP Sink 方案source 就是 Flink 自带的一个 HTTP 连接器轮询模拟器接口如果走消息队列就是 Kafka Source。第二个算子是解析和过滤。收到的原始字符串转成 Java POJO校验必填字段过滤掉 timestamp 距今超过 5 秒的乱序数据。这里提醒一点Flink 的序列化对 POJO 要求很严格必须是公有无参构造、字段有公有 getter/setter否则运行时直接报错说无法序列化这个类型。第三个算子是窗口聚合。举个例子按每个 AGV 的 ID 做 keyBy然后用滚动窗口每 5 秒统计一次当前位置、平均速度、电量变化率。Flink 的窗口 API 写起来非常顺手DataStreamAgvHeartbeat stream ...; stream .keyBy(AgvHeartbeat::getAgvId) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .aggregate(new AgvMetricsAggregate()) .map(new MetricsToJson());第四个算子是告警检测。用 Flink 的状态编程做连续 N 条数据异常才触发告警避免单条数据抖动导致误报。这里我用的是 ValueState 加计数器每条数据进来先更新状态再判断是否达到触发条件。这些模式识别逻辑用 FLink 做的好处是状态是托管在内存里的而且支持 Checkpoint任务挂了能从最近的状态恢复不会丢状态。这在 Spring Boot 里你要自己写的话可费劲了光一个超时判断的状态存储就够你头疼。4.3 Spring Boot 端接收 Flink 的数据并推送前端当 Flink 把聚合后的指标通过 HTTP Sink POST 到 Spring Boot 接口时Spring Boot 端要做三件事第一接收并解析这批数据把最新的车辆位置、电量、状态存到本地缓存Caffeine 或 Redis。这里不建议直接写数据库原因很简单高频写入数据库扛不住而且看板只需要最新状态。真正的历史记录可以单独落库那是另一个更低频的路径。第二通过 WebSocket 把数据推送给前端看板。前端页面已经订阅了/topic/agv-status这个频道后端有了新数据就主动推过去前端不用再轮询了。跟传统轮询相比WebSocket 是服务器推延迟低、性能好实时性完全不在一个量级。第三把告警事件单独建立一张表方便事后追溯。Flink 检测到的告警Spring Boot 收到后先落库再推送。这样即使前端当时没人在线告警也不会丢。WebSocket 配置在 Spring Boot 里不算复杂重写一下WebSocketHandler或者用 STOMP 协议都行。我自己习惯用 STOMP因为前端用sockjs-client加stomp.js就能直接订阅后端写起来也很干净。4.4 从模拟器到看板的完整验证路线整个链路跑通后最好的验证方式不是看日志而是直接看效果。我在本地跑通了这么一套完整演示打开 Flink Web UI能看到一个任务在跑每个算子的记录吞吐量都显示正常没有反压打开前端看板页面地图上能看到每一台 AGV 的小车图标在移动位置更新延迟肉眼不可见几乎实时手动在模拟器里把某台车的电量调到 15% 以下最多 1 秒钟前端页面弹出低电量告警数据库里也多了一条告警记录把模拟器停掉其中一台车的发送3 秒后前端标记这台车为通信中断同时触发告警这几个现象能出现就说明 Flink 的实时处理和 Spring Boot 的服务输出整个链路是通的。我见过很多项目开发到一半就卡在数据能收到但前端不刷新告警总是延迟几十秒这类问题上其实都是链路设计的问题不是某个组件单独的问题。5. 实测中的坑Flink 集成 Spring Boot 的常见障碍5.1 Flink JDBC 连接器异常一个必然遇到的坎在做 Flink 集成 MyBatis 或直接写 MySQL 的时候几乎每个人都会碰上一次 JDBC 连接器异常。最常见的是任务启动时报这个错误Caused by: java.lang.ClassNotFoundException: org.postgresql.Driver或者更隐蔽的一种MySQL 驱动类找不到报Cannot load driver class: com.mysql.cj.jdbc.Driver。原因不是你没加依赖而是 Flink 的类加载机制。Flink 里用户代码和框架代码是隔离的JDBC 驱动这类第三方 Jar 需要显式放到lib目录或者通过--classpath参数指定。你在 IDE 里跑没问题一提交到集群就翻车就是因为这个隔离机制。解决方案有两个一是把数据库驱动 Jar 复制到 Flink 的lib目录然后重启 Flink二是用flink run -C file:///path/to/mysql-connector.jar指定额外的 classpath。我建议本地开发直接用-C参数不用动 Flink 的 lib 目录干净省事。5.2 反压和背压本地看不出来但必须懂Flink Web UI 里有一个很直观的概念叫反压。简单理解就是上游算子处理速度大于下游算子消费速度导致数据在下游堆积这会拖慢整个任务的实时性。在 AGV 监控场景如果 Spring Boot 的 HTTP 接口响应变慢Flink Sink 的写入速度会下降最终反向传导到 source导致 AGV 心跳数据积压。本地调试数据量小反压基本不会出现但你可以人为压测一下把模拟器并发数调高到 200 台车然后看 Flink Web UI 上的背压指标。如果某个算子显示高背压通常意味着需要增加并行度或者优化下游处理逻辑。一个很常见的优化点是 Sink 改成批量写入不要一条一条调 HTTP 接口。5.3 事件时间和时区窗口统计结果偏了 8 小时Flink 处理事件时间数据时需要指定时间戳提取器和水位线Watermark。我第一次跑完 5 秒窗口聚合发现结果总是有点不对劲数据像是被分到了错误的窗口。排查了半天发现是在时间戳上踩了坑Java 里System.currentTimeMillis()是毫秒时间戳而 JSON 里如果用秒时间戳两者混在一起用Flink 解析出来的时间就乱了。还有一个本地 Windows 特有的坑系统时区如果不是 UTCFlink 默认会按系统时区解释时间导致用 Event Time 窗口聚合时统计的结果偏了 8 个小时。解决方法是显式在 Flink 环境里设置env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);然后时间戳提取器统一用毫秒并指定TimeZone。这类问题出现在数据看起来对但窗口边界总是怪怪的这种症状上不做数据验证很容易漏过去。5.4 Windows 特有的资源问题内存不够用Flink 本地模式默认会占用不少内存尤其是你同时跑着 IDEA、浏览器和模拟器的时候一不小心就 OOM 了。Flink 的 TaskManager 默认堆内存是 1GB外加堆外内存、网络缓冲一个任务轻松吃掉 1.5GB 以上。如果你电脑总共只有 8GB 内存建议在flink-conf.yaml里把内存调小taskmanager.memory.process.size: 1024m taskmanager.memory.jvm-overhead.fraction: 0.15本地调试追求的是逻辑正确不是性能吞吐。把内存压到 1GB任务照样能跑还不会把整台电脑拖垮。5.5 端口冲突和防火墙拦截Flink 的 Web UI 默认 8081Spring Boot 开发默认 8080加上你要给前端接口开一个 WebSocket 端口三个端口一块跑是常态。有一个很容易踩的坑是 Windows 防火墙突然拦截了 Java 进程的网络访问症状是Flink 能把数据发给本机 Spring Boot但从另外一台机器访问你的 Spring Boot 接口时连不通。如果你需要局域网内的前端看板访问记得在防火墙里放行对应端口或者至少是在开发阶段临时允许 Java 进程通过。端口是否被占用排查命令是netstat -ano | findstr 8080找到占用进程的 PID 后再taskkill /F /PID pid就能清理掉。Windows 下这个操作我几乎每周都要用一次。6. 从 Demo 到生产这套架构还能怎么演进6.1 消息中间件的升级路线本地 Demo 跑通了到了现场环境大概率要面对两个变化数据量变大以及稳定性要求变高。这时候 HTTP Sink 直连 Spring Boot 的方案就要让位给消息队列方案了。AGV 项目里最常碰到的消息中间件是 Kafka 和 MQTT两者选择标准很简单如果你的 AGV 终端本身通过 MQTT 上报那就 Flink 直接消费 MQTT如果调度中心统一收数据再分发那多半是 Kafka。切换的关键是 Flink Source 端换一个连接器业务逻辑算子根本不用动这也是把 Flink 放在中间层的架构优势。6.2 引入 Flink SQL 简化报表逻辑随着监控指标越加越多你会发现 DataStream API 里写聚合逻辑开始变得繁琐——每个新指标都要写一个自定义聚合函数。这个阶段可以考虑把一部分指标计算切换到 Flink SQL直接用 SQL 做窗口聚合和维度关联代码量能减少一半以上。比如统计每小时每个区域的 AGV 平均运行时长SELECT region, TUMBLE_START(proc_time, INTERVAL 1 HOUR) AS window_start, AVG(running_duration) AS avg_running_duration FROM agv_metrics GROUP BY region, TUMBLE(proc_time, INTERVAL 1 HOUR)Flink 的 Table API 和 DataStream API 可以在同一个任务里混用不用二选一。生产环境的常见做法是复杂的状态管理和告警规则用 DataStream常规报表聚合用 SQL各取所长。6.3 对接 Spring Boot Admin 做服务治理Spring Boot 这边的演进方向一个是引入 Spring Boot Admin 做服务健康监控。Flink 任务本身有自己的 Web UI但 Spring Boot 服务的健康状态、内存占用、接口响应时间Flink 管不到。Spring Boot Admin 的集成很简单在 Spring Boot 工程里加个spring-boot-admin-starter-client依赖再起一个 Admin Server就能在页面上一眼看到哪些客户端在线、哪些接口有告警。跟 Flink Web UI 搭配使用一个管流处理任务一个管业务服务整个系统就有了两个清晰的监控视角。6.4 历史轨迹和实时定位的分层存储AGV 项目到后期一定会遇到一个需求既要看实时位置也要看历史轨迹回放。这两类需求的数据特征完全不同。实时位置要求毫秒级查询延迟数据量不大历史轨迹要按时间段查询、轨迹线渲染数据量大但查询频率低。我目前的做法是实时状态存 Redis历史轨迹落 ClickHouse 或者 InfluxDB中间用 Flink 定期把 Redis 里的明细转存到时序数据库。这样 Flink 就承担了数据冷热分层转换的角色Spring Boot 只管从不同的存储里取数职责非常清晰。说句实在话这套架构的核心价值不在于某个单独组件而在于把流处理和传统后端服务合理分工。Flink 到底该管多宽、Spring Boot 该管多窄这个边界需要根据实际业务体量来拿捏。我个人的体会是本地调试期尽量少引入组件把链路跑通是第一位等业务逻辑稳定了再逐步引入 Kafka、时序数据库这些重型装备。技术栈越简单排查问题越容易这个原则在 AGV 这类现场环境尤其重要。最后分享一个我自己的实操习惯每到一个新环境我都会先花十几分钟把 Flink 的 Web UI 打开确认 TaskManager 的槽位和内存配置符合预期再开始提交任务。很多现场莫名其妙的性能问题最后都是出在默认配置和实际资源不匹配上。你在 Windows 本地练手的时候尽早养成看 Flink Web UI 的习惯后面到了真实项目里会省掉无数个深夜排查问题的时间。
返回列表