ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Flink实时交通监控平台实战:从架构设计到踩坑全记录

2026/9/17 14:46:43 拓冰建站 浏览量
Flink实时交通监控平台实战:从架构设计到踩坑全记录 今年手上正好做了一个城市交通实时监控平台从需求梳理到上线维护前后折腾了快三个月。这个项目最大的感受是它不是单纯把数据扔给Flink跑几个SQL就完事而是从数据接入、窗口计算、状态管理、结果落地到可视化展示一整条链路都要打通每一环都有坑等着你。这个平台最终要做的事情简单说就是三件实时接入城市主要路口的卡口数据、浮动车GPS数据计算各路段的平均车速和拥堵指数再把结果推给大屏和管理后台。数据量不大峰值也就每秒两三千条但胜在计算逻辑有点意思——要处理迟到数据、要维护路段状态、要跟静态的GIS路网数据做关联各种边界条件让人头大。如果你正在做类似的实时监控、实时大屏、实时数仓的项目或者准备用Flink做点实战练手这篇文章应该能帮你少走一些弯路。1. 整体设计与技术选型思路1.1 项目需求拆解监控平台到底要监控什么需求方一开始提得很笼统领导要看城市交通的实时态势。这句话翻译成技术语言需要拆出四块具体能力车辆实时位置与轨迹通过浮动车GPS数据在地图上实时描点看车辆走位。路段平均速度与拥堵指数这是核心中的核心。按路段link维度统计过去5分钟内通过车辆的平均速度再换算成拥堵等级畅通、缓行、拥堵、严重拥堵。异常事件告警比如某路段速度骤降、某路口车流长时间停滞需要秒级触发告警。历史回放与趋势对比虽然叫实时平台但值班人员和领导免不了要回看过去一小时、昨天同一时段的路况。拆完需求就清楚了这是一个典型的实时流式计算场景数据有GPS坐标、有卡口过车记录还要叠加上路段静态信息做维表关联。Flink在这个场景里几乎是标准答案后面细说。1.2 技术栈选型为什么是Flink而不是Spark Streaming或Kafka Streams选型时也做过对比主要纠结过Spark Streaming和Kafka Streams。Kafka Streams上手轻但问题在于它本质是一个库而不是计算引擎复杂的状态管理、窗口计算、精确一次语义都要自己拼装而且它强绑定Kafka如果数据源或结果存储要多样化写起来很别扭。这个项目要做窗口聚合、维表关联、异步IOKafka Streams不是不能用但开发效率低不少。Spark Streaming包括后来主推的Structured Streaming最大的痛点是微批。我们之前做过POC在窗口切换时会有秒级的数据延迟而且对于事件时间event time的处理Spark Streaming在当时版本下远不如Flink成熟。监控大屏上要求刷新延迟在5秒以内微批模式下很难稳定做到。Flink这边四个优势正好命中需求真流式计算延迟在毫秒到秒级配合窗口可以达到5秒内的大屏刷新要求。原生的event time watermark机制天然处理乱序和迟到数据这正是GPS数据最常见的窘境。强一致性的状态管理做去重、做路段状态维护都方便。Flink SQL的成熟度足够高我们的核心计算逻辑80%用SQL完成。最终架构定下来是Kafka数据接入 - Flink实时计算 - Doris结果存储 - Web可视化。Doris的选择后面单独说这里先卖个关子。1.3 整体架构图与数据流向设计我习惯在动手前把数据流画清楚就算不画正式架构图白板上也要理一遍。这个平台的数据流向分成六层数据源层浮动车GPS终端每5秒上报一条定位记录包含车辆ID、经纬度、时间戳、瞬时速度、方向角。卡口设备抓拍过车记录包含车牌、卡口ID、通过时间。两类数据都进入Kafka。接入层Kafka两个topic分别对应GPS原始数据和卡口过车数据分区数按8和4设置。计算层Flink集群核心作业有三个——位置明细作业清洗坐标转换维表补全、路段统计作业窗口聚合拥堵计算、事件告警作业速度骤降检测停滞检测。存储层Doris存聚合结果和明细Redis跑热数据缓存给大屏直接读MySQL存GIS路网静态数据。服务层SpringBoot服务提供HTTP接口大屏和后台拉数据。展示层Web大屏地图图表、管理后台、移动端。这里有个重要的设计决策GPS原始数据可能每秒上千条如果全部落到明细表再查Doris压力太大。所以明细只保留最近1小时用于轨迹回放过期定期清理而聚合结果永久保留。冷热分离各取所需。2. 环境搭建与Flink部署细节2.1 集群规划与资源估算我们用的是Flink 1.17部署模式选了Flink Standalone Cluster在Kubernetes还没完全普及的团队里Standalone 高可用是最容易维护的组合。三台机器16核32G系统盘加数据盘分开其中一台跑JobManager并配置HA另外两台跑TaskManager。资源估算倒不复杂按吞吐量反向推。目标支撑峰值5000条/秒单条JSON解析后不到1KB加上窗口计算和维表关联估计单并行度处理能力在800条/秒左右。预留2.5倍缓冲和故障恢复容量6个并行度足够。实际分配给每个TaskManager 8G内存其中托管内存managed memory留了2G给RocksDB状态后端。注意这里说的并行度是计算并行度不是TaskManager slot数。我习惯给TaskManager配1个slot减少同一进程内多任务互相干扰的问题虽然会有额外的JVM开销但排查问题时极度舒服。Flink安装配置这一块网上教程一抓一大把我只说过三点实战中容易翻车的flink-conf.yaml里的taskmanager.memory.process.size要留足默认值可能触发JVM Metaspace溢出。如果跑窗口聚合state.backend.type一定要配成rocksdb别用内存HashMap存大状态GC会让你哭。JobManager HA要配合ZooKeeper1.17版本甚至可以直接用内置的Kubernetes HA但传统运维团队还是习惯ZooKeeper稳定不折腾。2.2 踩坑实录Flink是不是一定要依赖HDFS热词里有一条“flink一定要hdfs”这个我也纠结过。很多教程在部署Flink时都会顺带搭一套HDFS原因是checkpoint默认要写到HDFS。但这不代表Flink强制依赖HDFS。我们当时没有现成的HDFS集群也不想为了这个项目单独搭一套Hadoop。解决方案是用文件系统做checkpoint也就是把checkpoint存到本地磁盘或者NFS共享目录。代价是什么如果机器本地盘挂了checkpoint会丢任务能恢复到什么程度取决于最后一次成功的checkpoint在不在。我们评估后觉得可以接受checkpoint每30秒一次挂了最多丢30秒数据而上游Kafka会留存数据通过Kafka的offset重置机制可以把这30秒补回来。如果你有条件上HDFS或S3最好还是用好一点的存储。但要说Flink一定要HDFS那是误解。除了checkpoint和高可用相关的存储Flink运行时本身完全不需要HDFS。2.3 开发环境与SQL Client调试技巧开发阶段大部分时间在跟Flink SQL打交道。Flink 1.17自带的SQL Client虽然能用但交互体验一般。我们组里实际开发用两种方式关键作业用Java写DataStream API调SQL快速验证逻辑用DBeaver连接Flink Gateway执行SQL。这里额外分享一个技巧在IDE里本地调试Flink SQL作业时可以启动一个本地Flink MiniCluster然后通过Flink REST API提交SQL任务。这样既能断点调试UDF又能验证完整的SQL执行计划。当时就是靠这个方式把窗口SQL的watermark节奏和join行为调明白的。3. 核心功能与Flink SQL实现3.1 车辆位置实时接入与预处理GPS原始数据的JSON结构大概是这样的{ vehicleId: 沪A12345, lng: 121.4737, lat: 31.2304, speed: 42.5, angle: 135, timestamp: 2024-06-15 08:30:12, sourceType: gps }直接把这个JSON交给Flink解析没问题但要做两件预处理第一坐标转换。原始GPS坐标是WGS84坐标系地图底图用的是GCJ02火星坐标系如果直接叠加地图上的点会偏移几十米到几百米不等。所以我在Flink里写了个UDF用标准坐标偏移算法做了转换。第二数据清洗。GPS设备偶尔会上报乱码、经纬度为0、车速为负的脏数据。清洗逻辑不复杂经纬度超出城市边界范围就丢弃车速小于0或大于200km/h丢弃时间戳偏离当前时间超过10分钟丢弃。这些判断用SQL的WHERE条件就能完成。CREATE TEMPORARY VIEW cleaned_gps AS SELECT vehicle_id, convert_coord(lng, lat) AS (gcj_lng, gcj_lat), speed, ts FROM source_kafka_gps WHERE lng BETWEEN 120.8 AND 122.0 AND lat BETWEEN 30.7 AND 31.8 AND speed 0 AND speed 200 AND ts BETWEEN NOW() - INTERVAL 10 MINUTE AND NOW() INTERVAL 1 MINUTE;迟到的数据这一层不直接扔掉而是打上迟到标记再往下游传给窗口聚合层自己决定要不要参与计算。这样保证统计口径可追溯。3.2 路段速度统计5分钟窗口与迟到数据处理路段速度统计是最核心的作业。逻辑大概是把车辆GPS点匹配到路段上每个GPS点估算出“这个车在这条路段上的瞬时速度”然后对同一路段同一时间窗口内的所有点求平均速度。这里用到了两个关键的Flink SQL特性事件时间与水位线。GPS数据在网络上传输时天然会乱序只能靠事件时间而不是处理时间来计算。我在建Kafka源表时显式声明watermarkCREATE TABLE gps_source ( vehicle_id STRING, lng DOUBLE, lat DOUBLE, speed DOUBLE, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic gps_raw, properties.bootstrap.servers kafka1:9092,kafka2:9092, properties.group.id gps-group, scan.startup.mode latest-offset, format json, json.timestamp-format.standard SQL );水位线设置为5秒意味着最多容忍5秒的乱序。实际运行下来GPS数据的乱序程度基本在2秒以内这个参数留了余量。窗口聚合语句长这样INSERT INTO doris_sink_link_speed SELECT link_id, TUMBLE_START(ts, INTERVAL 5 MINUTE) AS window_start, TUMBLE_END(ts, INTERVAL 5 MINUTE) AS window_end, AVG(speed) AS avg_speed, COUNT(*) AS sample_count, PERCENTILE_APPROX(speed, 0.9) AS v90 FROM ( SELECT match_link(lng, lat) AS link_id, speed, ts FROM cleaned_gps ) GROUP BY link_id, TUMBLE(ts, INTERVAL 5 MINUTE);这里的match_link是一个自定义函数作用是根据经纬度匹配到最近的路段。这个函数内部用了一个内存中的空间索引结构把城市路网划分成网格先定位到网格再精确匹配单次匹配耗时在微秒级别。提示5分钟窗口和1分钟窗口的实际效果差异很大。5分钟窗口可以平滑掉瞬时波动适合展示拥堵指数1分钟窗口保留更多细节但会有不少空窗和抖动。最终我们做了双链路一条跑5分钟窗口用于大屏展示一条跑1分钟窗口用于告警检测。3.3 拥堵指数计算从速度到等级的换算逻辑有了平均速度怎么转成拥堵等级这个没有统一标准我们结合需求方的经验定了一套阈值同时考虑不同道路等级快速路、主干路、次干路、支路的差异道路类型畅通缓行拥堵严重拥堵快速路50km/h35km/h20km/h20km/h主干路30km/h20km/h10km/h10km/h次干路25km/h15km/h8km/h8km/h支路20km/h12km/h6km/h6km/h这个表不是写死在代码里的而是从MySQL维表加载用Flink的维表join实时关联。这样需求方想调阈值直接在后台改表就行不用重启作业。维表关联用LOOKUP JOIN缓存策略选的是ALL全量缓存因为路网数据总共就几十万条全量加载到内存完全没问题还能避免查MySQL的性能损耗。CREATE TEMPORARY VIEW link_with_speed AS SELECT s.link_id, s.window_start, s.avg_speed, d.road_type, d.link_name FROM doris_source_link_speed s LEFT JOIN mysql_dim_link d ON s.link_id d.link_id;3.4 卡口流量与异常事件告警告警作业是另一条链路。我们监控两类异常第一类某路段平均速度在连续两个1分钟窗口内下降超过40%。这代表发生了突发拥堵或事故。实现上并不复杂用MATCH_RECOGNIZE做模式匹配识别出“快速-更快-骤降”的趋势序列。第二类某卡口在5分钟内过车数为0但历史上这个时段平均过车数大于50。这代表卡口可能故障或者路段已经完全堵死。需要先把历史同期均值数据缓存到Redis再在Flink作业里通过异步IO访问Redis比对。第二类涉及外部存储关联Flink SQL的维表join可以搞定但延迟会稍高。后来我们改成在数据进入窗口聚合前先做一次open的生效检查把整个判断逻辑封装在一个UDF里性能提升明显。告警结果写入告警表同时通过WebSocket推给大屏大屏上弹红色气泡值班人员可以秒级看到。3.5 温故知新数据血缘如何跟踪热词里提到“flink数据血缘”这个项目也做了点尝试。实时计算作业多了以后最痛苦的不是写SQL而是维护“这张Doris表的数据到底从哪来、经过了哪些计算、依赖哪些原始表”。Flink 1.17有内置的Lineage机制可以在Catalog中自动记录作业的输入输出血缘。我们在此基础上做了一层轻量的血缘管理每个Flink SQL作业在提交时附带一个作业元数据JSON包含来源topic、目标表、窗口策略、责任人。最终血缘关系通过Doris的元数据表统一查询。别小看这个动作后面排查数据问题时省了大量时间。区域A的数据对不上打开血缘图一眼定位是该区域的GPS清洗逻辑有问题还是窗口参数配错了。4. 结果存储与数据服务设计4.1 为什么选了Doris而不是MySQL或ClickHouse聚合结果和明细数据需要支撑大屏的高并发点查同时要支持一些多维度的即席分析。一开始考虑过ClickHouse但在点查场景和并发更新上不如Doris顺手。MySQL则不太适合存这种按时间分区的、高频写入的流式结果数据。Doris的优势在于支持批量导入和高并发点查正好适配Flink写入大屏查询的模型。内置的Unique模型支持主键更新对数据重算和迟到数据修正非常友好。冷热分区和动态分区管理省心一天一个分区老数据自动过期。Doris建表时我们重点考虑了两个点。第一分区和分桶键要按查询模式来。常用的查询维度是“路段ID 时间”所以建表用时间做分区路段ID做分桶键。第二Unique模型要设置sequence列处理乱序写入避免迟到数据把新数据覆盖掉。CREATE TABLE IF NOT EXISTS link_speed_5min ( link_id BIGINT, window_start DATETIME, avg_speed DOUBLE, sample_count INT, v90 DOUBLE, congest_level INT ) ENGINEOLAP UNIQUE KEY(link_id, window_start) DISTRIBUTED BY HASH(link_id) BUCKETS 8 PROPERTIES ( dynamic_partition.enable true, dynamic_partition.time_unit DAY, dynamic_partition.end 3, dynamic_partition.start -30, replication_num 1 );4.2 Flink Doris Connector的典型报错与解决办法在集成Flink和Doris时我们遇到了一个很典型的报错也是热词里提到的“Flink type is DATEV2, but arrow type is DateDay”。这个报错的原因是Flink的Doris Connector在读取Doris表时列类型映射出现问题。Doris的DATEV2类型在和Flink的Arrow格式交互时被映射成了DateDay类型而Flink内部认为自己读到的是DATEV2两边对不上就报了异常。解决办法其实简单两种升级Connector版本。Doris官方提供的flink-doris-connector新版已经修复了DATEV2的映射问题。建表时不要用DATEV2类型改用DATE或DATETIME或确保Connector和Doris版本兼容。这个报错折腾了我们半天核心教训是Doris Connector版本和Doris集群版本要严格匹配最好用官方推荐配对的版本组合。不要拿个最新版Connector去连旧版Doris也不要反过来。4.3 TiDB Flink SQL的兼容性对比热词里有“tidb flink sql”因为最开始我们考虑过用TiDB当结果存储。TiDB的Flink SQL Connector也相当成熟尤其在实时数仓场景下跟Flink配合得很好。它兼容MySQL协议如果我们团队更熟悉MySQL生态用TiDB确实上手快。但对比下来有个关键差异Doris的Unique模型天然适配“按主键覆盖更新”的写入方式Flink直接把聚合结果写入即可而TiDB要走INSERT ... ON DUPLICATE KEY UPDATE或REPLACE INTO来做更新写入性能在某些场景下会略逊一筹而且TiDB的批量导入链路没有Doris的stream load那么顺滑。所以最终选了Doris。如果你的场景偏事务型查询且需要MySQL高度兼容TiDB是完全可行的选择但纯OLAP分析大屏场景Doris更合适。4.4 JDBC连接器异常排查经验项目里还遇到过Flink的JDBC连接器异常。一次作业运行几天后突然报错错误信息大致是连接超时或连接池耗尽。排查过程分三步第一看是不是连接泄漏。JDBC Connector每个并行实例会持有自己的连接池如果连接未正常释放运行久了必然出问题。检查SQL中是否有窗口JOIN或维表JOIN频繁占用连接并确保lookup.cache配置正确。第二看MySQL的wait_timeout。MySQL默认8小时断掉空闲连接如果Connector没做连接有效性检测就会拿到一个坏连接。解决办法是调大MySQL超时时间或在连接串中加上autoReconnecttrue。第三看并发和连接池配置。Flink作业的并行度乘以每个并行度的连接数不能超过MySQL的max_connections否则就是雪崩。一个更省心的做法维表数据量不大时优先用LOOKUP JOIN的ALL缓存一次性加载全量数据避免频繁访问MySQL。5. 常见问题与故障排查实录5.1 数据延迟越来越大问题出在哪上线第一天就遇到一个教训。大屏上的数据延迟由最初的3秒慢慢涨到30秒最后直接卡住。排查思路先看Flink UI发现某个作业的背压backpressure已经变成HIGH。再看CPU和内存发现某个TaskManager的CPU已经打满垃圾回收频繁。进一步检查是维表JOIN的缓存策略配置出了问题——LOOKUP JOIN默认可能不开缓存每条数据都要查一次MySQL吞吐上不去。解决办法把维表缓存改成ALL并设置lookup.cache.ttl为1小时。修改之后压力立即降下来。注意维表缓存不是越大越好。如果维表数据会更新比如路网等级调整TTL太大会导致Flink读到的还是旧数据。建议维度数据更新不频繁时用ALL缓存频繁更新时用LRU缓存并配置合理的TTL。5.2 窗口计算结果突然为空的诡异问题有一次某条路的5分钟统计结果一直是空的其他路正常。检查Kafka日志发现那条路的数据其实一直在进。最后定位到问题出在match_link函数那条路的坐标落在我划分的空间网格的边缘网格索引的容差设置太小导致坐标匹配到相邻网格后找不到匹配路段。修复方式很简单把网格匹配的容差从50米调到100米并且增加一个“找不到最近路段时尝试相邻网格”的兜底逻辑。这种坐标匹配类的边界问题单测很难覆盖全最好的办法是线上多留日志把匹配不到的数据单独落到一个debug表定期人工检查。5.3 背压排查与性能调优实录背压问题在实时计算里几乎躲不掉。我的排查顺序是看Flink UI的背压指标确定背压发生在算子之间还是整个作业层面。检查Sink算子的写入性能Doris批量写入如果攒批参数不合理很容易成为瓶颈。sink.batch.size和sink.batch.interval要配合调整。检查窗口算子的状态大小RocksDB状态太大时序列化和反序列化开销会明显上升。用spillable和调整RocksDB的block cache可以缓解。最后才是调整并行度和资源配置。我们最终通过把Sink的攒批大小调整为每批2000条或每2秒刷一次Doris写入吞吐提升了一倍背压问题基本消失。5.4 数据重复与精确一次语义一开始用Kafka Flink Doris的标准链路时Doris里偶尔会出现重复数据。原因是Flink的checkpoint开启后作业重启恢复时会从checkpoint和Kafka的offset重新消费如果Sink不是幂等的就会重复写入。Doris的Unique模型天然支持主键幂等重复写入相同主键的数据会覆盖所以这个问题被Doris部分掩盖了。但明细表有些场景没法用Unique模型于是开了Flink的CheckpointingMode.EXACTLY_ONCE配合Doris的stream load两阶段提交确保端到端精确一次。这里要强调精确一次不是只开一个开关就完事上下游必须配合。上游Kafka要支持从offset恢复中间Flink要开checkpoint下游Sink要支持事务或幂等写入缺一个环节精确一次都是空话。5.5 一张问题排查速查表现象可能原因排查方法解决办法大屏延迟持续增大维表JOIN缓存未开或配置不当看Flink UI背压指标LOOKUP缓存改为ALL或LRU合理TTL窗口结果为空数据匹配不到维度查debug表/日志调整匹配容差增加兜底逻辑结果数据偶尔重复Sink未做幂等检查主键约束使用Unique模型或开启精确一次作业运行几天后报错JDBC连接池泄漏/超时查连接数和wait_timeout调整连接配置或改用ALL缓存数据类型映射报错Connector版本不匹配查Doris/Connector版本升级Connector或调整建表类型迟到的数据把新数据覆盖窗口未处理迟到数据看watermark设置延迟关闭窗口或设置主键sequence6. 项目上线后的几点经验沉淀最后聊点跟技术无关但跟项目成功有关的体会。这类监控平台业务方一开始说的需求永远是模糊的但他们对延迟和准确性的要求却非常具体。建议在动手前先跟业务方对齐两个核心指标数据从产生到大屏可见的最长链路时间以及可容忍的数据误差范围。链路时间是技术问题误差范围就有点微妙了。我们当时跟业务方明确过GPS数据正常情况下的覆盖率在85%左右如果某个路段样本量太少比如一分钟内只有两三辆车经过统计结果可能失真。所以系统对样本量小于5的路段做了特殊标记在界面上显示为“样本不足”而不是给出一个看起来精确但并不可靠的拥堵指数。这种主动暴露不确定性的做法反而让业务方更信任系统。维护期最大的坑是版本兼容。Flink、Doris Connector、Kafka Client这些组件的版本升级一定要谨慎升级前务必在测试环境完整回归一遍SQL作业不要轻信“兼容”两个字的官方文档。最后再分享一个小技巧上线前把Kafka的topic数据备份一份到本地用这份数据做回放测试。这样既能验证新版本作业的正确性又不会影响线上链路。我们靠着这个习惯两次大版本升级都平稳过渡了。