ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink实时数仓实战:从GPS轨迹到城市交通监控平台全链路解析

Flink实时数仓实战:从GPS轨迹到城市交通监控平台全链路解析 城市交通实时监控平台这个项目做出来之后效果还不错今天抽空把整个链路设计、核心代码、以及上线前后踩过的坑都整理一遍。项目的核心就是用Flink把车载GPS轨迹、路口线圈流量、卡口识别数据实时接入经过流式清洗、坐标匹配、窗口聚合之后落到Doris和Redis里最后推到地图大屏和告警服务。适合正在学Flink想找个完整实战项目参考的人也适合需要设计实时监控类业务的同学哪怕你不是做交通的换到物流、外卖、工厂设备监控这套框架基本也能直接搬。1. 先搞清楚这类项目到底在解决什么问题1.1 “实时监控”的实时到底要多“实时”交通场景里的“实时监控”和做离线报表完全是两码事。离线报表是T1今天看昨天的数据凌晨跑批早上出报告慢一点没人在意。但交通监控不同路口堵没堵、事故有没有发生、信号灯要不要调整这些问题拖五分钟可能就已经造成了二次拥堵。所以这个项目里我定的实时目标分了三档。第一档是秒级指标比如在线车辆数、某个路口的瞬时车流量这部分主要用于大屏展示延迟控制在5秒以内。第二档是分钟级指标比如路段平均车速、拥堵指数、区域热力开5分钟窗口做一次聚合用来支持调度决策和拥堵预警。第三档是准实时的事件告警比如车速骤降、车辆异常停留、逆行这些需要比较复杂的规则判断允许30秒到1分钟的延迟但必须准确。这里要特别强调一点不要为了“实时”而实时。我们在设计的时候反复讨论过是不是所有指标都要做到秒级答案是否定的。秒级计算的成本是分钟级的几倍对状态后端、Kafka分区、并行度配置都有更高要求。合理的做法是把指标按业务用途分级不同级别匹配不同的窗口大小和计算频率。这也是这个项目架构设计的一个核心理念。1.2 为什么是Flink而不是Spark Streaming或Storm选型的时候团队内部确实吵过一轮。Storm是老牌的流计算框架延迟确实低但吞吐上不去开发又很痛苦一个简单的窗口聚合要写一大堆Java代码维护成本高团队里能写Storm的人也越来越少。Spark Streaming用微批模拟流处理生态好、上手快但它的本质还是批窗口的实时性做不到真正的秒级而且状态管理、精确一次处理都要绕不少弯。最后选了Flink核心是看中了它这几点。第一它是真正的流处理引擎事件从进入到产出是一条流水线低延迟和高吞吐能同时兼顾。第二状态管理是Flink的看家本领RocksDB状态后端可以支撑海量key的状态存储做车辆轨迹跟踪、去重、乱序处理都很顺手。第三Flink SQL让实时计算的开发门槛降了一个量级团队里只会写SQL的同学也能参与到实时指标开发中来。第四生态太完整了Kafka、Doris、Redis、JDBC都有官方连接器我们这个项目的存储层正好全用得上。Flink的典型使用场景比如实时数仓、实时风控、实时推荐还有这个项目所在的实时交通监控本质上都是“数据实时进来需要立刻算出一个结果并对外服务”的诉求。这和批处理是完全不同的思路也是Flink能做起来的原因。2. 整体架构与核心模块拆解2.1 端到端链路从路侧设备到地图大屏整个项目的链路简单来说就是“数据源 - Kafka - Flink - Doris/Redis - 应用”。数据源分成三路第一路是车载GPS终端按照每秒或每三秒上报经纬度、速度、方向第二路是路口埋设的线圈和地磁检测器输出车流量和占有率第三路是卡口摄像头的车牌识别结果带时间、位置和车牌号。这些数据通通先进Kafka。我们建了三个Topicods_trajectory存轨迹ods_flow存流量ods_event存卡口识别的事件数据。这里用Kafka做缓冲层的原因很简单上游设备经常抽风数据可能突然暴涨也可能断流Kafka像一个蓄水池能把流量峰值削平让Flink消费的时候有一个相对平稳的速率。Flink这边我采用的是“Flink SQL DataStream API”混用的方式。清洗、过滤、窗口聚合这些逻辑用SQL写代码量少、逻辑清晰而坐标匹配这种需要复杂空间计算的部分则用DataStream API写ProcessFunction实现。混用不是乱用而是因为这两类算子用SQL表达的成本和用代码表达的成本差距太大混用正好能发挥各自优势。计算完的结果分两路存储。聚合指标写入Doris用作大屏查询和历史回溯实时排行的数据写入Redis比如全城拥堵Top10路段、当前在线车辆数大屏每一次刷新直接读Redis毫秒级响应。Doris和Redis配合一个管深度的分析查询一个管高频的简单读取各司其职。2.2 架构选型里的几个关键权衡第一要不要用Lambda架构就是离线批处理和实时流处理各维护一套代码最后把结果合并。我个人的观点是除非你已经有非常成熟的离线数仓体系否则不要轻易上Lambda。两套代码意味着两套逻辑同一个指标离线算出来一个值实时算出来又是另一个值到时候对数据都是个灾难。这个项目采用的是Kappa架构只走实时链路如果某天需要重算历史数据直接用Savepoint把作业恢复到某个时间点从Kafka对应位点重放即可。第二为什么结果存储用Doris而不是ClickHouse或TiDB。ClickHouse的单表查询确实快但Join能力弱实时更新也不方便。交通监控这种场景车辆轨迹明细和路段聚合指标是要频繁更新的Doris的Unique模型正好支持主键更新Aggregate模型则天然适配聚合结果写入。而且Doris的部署维护比ClickHouse简单一套FE加BE就能跑起来很适合中小团队。TiDB当然也很好HTAP能力很强如果团队本身已经在用TiDBFlink SQL通过JDBC连接器直接写TiDB也完全可行但对我们当前这个项目来说Doris的导入生态和预聚合能力更顺手所以最终选了Doris。第三Flink的状态和Checkpoint到底要放哪里。这个细节放到后面的踩坑章节详细说但架构设计阶段就要定好不能等出了问题再想。生产环境一定要用分布式文件系统HDFS或者对象存储都可以。3. 核心模块实现从原始轨迹到路况指标3.1 数据模型设计轨迹、路段、围栏数据模型是整个实时计算的地基模型设计不好后面写什么SQL都别扭。我们定义了三个核心对象。第一个是轨迹数据。Kafka里消费的原始轨迹消息长这样{ vehicle_id: 沪A12345, lng: 121.473701, lat: 31.230416, speed: 12.5, direction: 180, source_type: taxi, event_time: 2024-05-20T08:35:12Z }direction是车辆航向角0到360度正北为0顺时针递增这个字段后面做坐标匹配方向约束时会用到。时间是绝对关键的信息我们同时保留了事件时间event_time和摄入时间Flink SQL里必须基于事件时间开窗口。第二个是路网数据存在MySQL里。一张路段表dim_section字段包括路段ID、起点经纬度、终点经纬度、道路等级、自由流速度、路段长度。这张表不会频繁变动所以我们在Flink里面把它当作维表通过Lookup Join关联。同时考虑到坐标匹配的需要我又把路网数据做成了广播流方便在ProcessFunction里做网格索引查询。第三个是电子围栏也就是区域热力计算时需要把城市划分为网格或行政区。网格ID的计算规则很简单用经纬度除以网格边长取整拼接就行比如边长500米grid_id concat(cast(floor(lat/0.005) as string), _, cast(floor(lng/0.005) as string))。行政区划则作为静态维表存在MySQL里。这里有个经验数据模型一定要区分哪些是静态数据、哪些是动态数据。静态数据路网、围栏可以做成维表或者广播流状态很小查起来没有压力动态数据轨迹、流量才是流处理的重点状态管理要谨慎设计该设TTL就设TTL否则状态会无限膨胀。3.2 实时路况计算坐标匹配与平均车速路况计算的第一步是把一个GPS点匹配到它所在的路段上。这一步在纯SQL里面很难做我们是用ProcessFunction实现的。先给路网建立网格索引。把城市按经纬度切成一个二层的网格第一层网格边长约1公里在内存里维护一个“网格ID - 路段ID列表”的映射。一个路段通常跨越多个网格所以一个路段会出现在多个网格的路段列表中。轨迹点进来后第一步就是根据点的经纬度算出它所在的网格取到这个网格关联的所有候选路段。然后是垂距计算。对每个候选路段计算GPS点到线段的最短距离。这块我会用JTS这个空间计算库它提供了一个DistanceOp可以直接计算点到线段的距离。如果最短距离小于阈值城市道路取20米高速取30米就认为这个点匹配到了该路段如果多个路段都满足阈值则结合上一个匹配点的路段ID和方向角做约束选出方向最一致的那个。匹配完成后输出一条带section_id的匹配轨迹然后回到Flink SQL做5分钟窗口的平均车速计算CREATE VIEW matched_trajectory AS SELECT vehicle_id, section_id, speed, event_time FROM trajectory_source CROSS JOIN LATERAL TABLE( match_section(lng, lat, direction, event_time) ) AS m(section_id);这里match_section是我们注册的UDTF输入经纬度和方向输出匹配到的路段ID。然后正常的窗口聚合CREATE TABLE doris_section_agg ( section_id BIGINT, window_start TIMESTAMP(3), sample_count BIGINT, avg_speed DOUBLE, max_speed DOUBLE ) WITH ( connector doris, fenodes doris-fe:8030, table.identifier rt_traffic.section_agg, username readonly, password 123456, sink.label-prefix flink_section_agg, sink.properties.format json, sink.properties.read_json_by_line true, sink.enable.batch-mode false ); INSERT INTO doris_section_agg SELECT section_id, TUMBLE_START(event_time, INTERVAL 5 MINUTE) AS window_start, COUNT(*) AS sample_count, AVG(speed) AS avg_speed, MAX(speed) AS max_speed FROM matched_trajectory GROUP BY section_id, TUMBLE(event_time, INTERVAL 5 MINUTE);聚合之前一定要过滤异常数据。比如速度大于200公里/小时的、经纬度明显超出城市边界的、还有速度小于2公里/小时的“僵尸车”——这些车辆停在路边如果不剔除会把路段的平均车速拉低拥堵指数就会误报。我们在轨迹清洗阶段会统一处理。3.3 拥堵指数与区域热力二次聚合的经典套路拥堵指数业界常用的是实际行程时间和自由流行程时间的比值。自由流速度在路网表里已经有了比如某个城市快速路自由流速度是80公里/小时5分钟窗口内的平均车速是40公里/小时那这个路段的拥堵指数就是80除40约等于2.0。再按数值映射到0到10的评级1到2是畅通2到4是缓行4到6是拥堵6以上是严重拥堵。这个计算过程是基于路段聚合表doris_section_agg做的不需要在Flink里再搞一遍SQL里一条更新语句就能解决UPDATE section_agg_view SET congestion_index ROUND(free_flow_speed / avg_speed, 2) WHERE sample_count 10;这里有个边界问题如果某个路段在5分钟内只有一两辆车经过算出来的平均速度没有统计意义所以sample_count小于10的路段我们直接不更新拥堵指数保持上一次的值。这种“样本量门槛”的做法在交通场景里特别重要否则一个偶然的极值就会导致整条路的指标异常。区域热力相对简单就是把网格ID加到轨迹数据上按网格做5分钟的去重车辆数统计。去重车辆数要用COUNT(DISTINCT vehicle_id)或者为了性能用APPROX_COUNT_DISTINCT精度略低但在大屏展示场景完全够用。然后写入Doris的Aggregate模型表等大屏查询的时候再按区域聚合展示。这里要提一个容易踩的坑COUNT(DISTINCT ...)在实时流计算里会消耗大量状态因为每一个车辆的ID都要存在状态里去重。如果车辆规模达到几十万甚至上百万这个状态会非常庞大。我们后来优化成基于Redis的HyperLogLog做近似去重Flink计算的时候只输出每个网格的原始车辆数由Redis维护去重逻辑状态压力一下子就降下来了。4. 存储选型与查询设计结果写到哪里去4.1 明细/聚合双轨Doris加Redis的配合结果存储的设计决定了实时监控平台能撑住多大的查询压力。这个项目里我采用了“Doris管深度、Redis管速度”的双轨策略。Doris这边建两类表。一类是轨迹明细表用Unique模型主键是vehicle_id加event_time组合按天分区。明细表不会给大屏直接查询用它主要服务“回放某辆车某天轨迹”这类事后分析需求以及作为问题排查的数据底座。另一类是路段5分钟聚合表用Aggregate模型按路段ID和窗口开始时间聚合这条路径的查询频次最高大屏的地图热力、路况列表都是从这张表读的。Doris建表的时候有几个细节要提前设计好。分区一定要按时间范围做查询都会带上window_start条件没有分区裁剪的话Doris全表扫描一次性能很难看。排序键也要按查询模式设计路段聚合表的排序键是section_id加window_start这样单条路段的多天查询能利用前缀索引。写入的批次大小要调Doris Stream Load在数据量不是特别大的情况下可以开sink.enable.batch-mode true减少导入次数提升吞吐。Redis这边的角色更偏向“实时快照”。我们用Zset存实时路况排行key是全城维度member是路段IDscore是拥堵指数这样大屏要展示“当前全城最堵的10条路”一条ZREVRANGE命令就够了。用String存在线车辆数每次轨迹心跳到达的时候通过Flink的Redis Sink自增。用Set或者Hash存当前活跃的电子围栏区域。Redis的数据量不大但QPS非常高大屏一秒刷一次甚至一秒刷五次都扛得住。补充说一句如果你们团队已经在生产环境用了TiDB那Flink SQL通过JDBC连接器直接写TiDB完全可行TiDB的HTAP能力也足够覆盖“实时写入加复杂查询”的需求。但就我们这个项目来说Doris的实时导入生态和预聚合模型更成熟所以没有额外引入TiDB。4.2 关于“Flink一定要HDFS吗”的正确答案这个问题在社区里被问烂了但我觉得还是要认真回答一次。Flink本身不强制要求HDFS它只要求你配置一个支持分布式语义的文件系统来存放Checkpoint和Savepoint。HDFS只是其中一个选项S3、OSS、MinIO都行。Checkpoint不能放本地磁盘。原因很简单Flink任务跑在集群的多个节点上如果Checkpoint写到某个TaskManager的本地磁盘一旦这个节点宕机状态数据就丢了作业恢复的时候找不到状态就只能从头开始消费相当于白干。而分布式文件系统天然具备跨节点访问能力任何节点挂了其他节点都能从文件系统里取回状态作业才能做到快速恢复。RocksDB状态后端倒真的是用本地磁盘它把状态存到TaskManager本地但每次Checkpoint的时候会把快照同步到分布式文件系统。所以生产环境的标准配置是RocksDB做状态存储HDFS或S3做Checkpoint存储。开发环境随便跑一下无所谓但生产环境一定要按这个标准来。我见过一个团队因为图省事没有配Checkpoint存储Flink任务跑了一个月某天Kafka集群抖动导致作业重启结果所有状态全丢了窗口数据全部重新计算线上指标乱成一锅粥。所以这个问题根本不是“要不要HDFS”而是“要不要自己给自己挖坑”。4.3 数据血缘排查问题的重要依据数据血缘在这个项目里不是做给别人看的治理文档而是我们自己排查问题的工具。实时链路长环节多一个指标异常你得知道它来自哪个Topic、经过哪个作业、写入了哪张表。我们做了两件很简单但很有效的事。第一给所有Flink作业规范命名。通过SET pipeline.name rt_traffic_section_agg这条语句每个作业在JobManager上都有一个含义清晰的名字日志、指标、Checkpoint目录里都能看到。第二启动作业的时候把作业的输入源表、输出目标表、作业名写入MySQL的元数据表里形成一张“作业血缘表”。这张表平时没人看但出问题的时候能救命。有一次线上反馈某个路口的拥堵指数连续半小时异常偏高我第一反应是算错了。通过血缘表定位到这个指标对应的是rt_traffic_section_agg作业再看这个作业消费的Topic是ods_trajectory然后查Kafka消费组的Lag发现有一个分区在持续积压原因是有两台设备上报的频率太高把Kafka某一个分区的吞吐打满了。如果从头开始看代码这个问题的排查至少要多花半天。所以做实时项目请一定把血缘信息维护好。哪怕你早期不引入Atlas这类重量级工具用一张MySQL表记录作业维度的血缘关系投入产出比也非常高。5. 实战踩坑那些文档里不会明说的细节5.1 Flink SQL写Dorisdatev2和dateday的类型映射坑这个报错我要单独拎出来说因为太经典了。我们第一次用Flink SQL把数据写入Doris的时候作业跑着跑着就报错flink type is datev2, but arrow type is dateday. at org.apache.doris.flink.什么意思呢Doris的Flink连接器在读取数据时默认使用Arrow格式而Arrow会把Doris的DATE类型映射成DateDay类型。但Doris 2.0之后新版本里DATE字段被建议使用DATEV2类型。于是在Flink侧建表的时候如果字段类型写成了DATEV2而连接器读取到的Arrow类型是DateDay两边就对不上直接抛异常。解决方案有三个按优先级排序。第一个如果Doris表是新建的建议表结构里日期字段直接用DATEV2同时在Flink SQL建表语句里把对应字段也写成DATEV2并升级Doris连接器到较新的版本新的连接器会做自动转换。第二个如果Doris表是历史表已经是DATE类型那Flink侧建表语句里把该字段的类型写成STRING或者DATE不要在SQL里显式写DATEV2避免触发类型映射冲突。第三个如果业务上允许干脆在Flink SQL里用DATE_FORMAT(event_time, yyyy-MM-dd)转成字符串再写入简单粗暴但会损失一部分在Doris侧做日期计算的能力。这个坑的根因其实就是类型映射不一致。Doris发展太快正处在新旧类型切换的过渡期连接器还没完全跟上。遇到类似报错先检查两边类型是否一致不要一上来就怀疑是网络或者并行度的问题。5.2 JDBC连接器在高并发写入时的连接耗尽热词里有个“flink的jdbc连接器异常”这个我太有体会了。我们起初用Flink JDBC连接器把指标写入MySQL用来支持一些简单的后台查询。表结构很简单字段也不多本地测试一切正常一上生产就不停报错错误信息是Too many connections和Communications link failure。一开始以为是MySQL的连接数上限不够直接调大了max_connections结果只是临时缓解过几分钟又开始报错。后来看Flink的监控页面发现TaskManager的JDBC连接数一直往上飙才意识到问题出在Flink侧。我们一度把并行度调到了24Flink每个并发都会维护自己的JDBC连接24个并发就已经很吓人了。但更关键的是我们当时没有开启写入端的flush缓冲默认每条记录都会直接走JDBC写入导致连接申请和释放非常频繁连接池根本扛不住。解决办法是两板斧。第一降低并行度到合理的范围比如8让并发数和数据库的连接池容量匹配。第二打开Flink JDBC Sink的批量写入参数sink.buffer-flush.max-rows 500, sink.buffer-flush.interval 2s, sink.max-retries 3这样一来数据先攒到500条或者2秒一次才真正写数据库连接复用率大幅提升问题就解决了。记住Flink的JDBC连接器默认行为并不适合高并发生产环境凡是写入MySQL这种传统数据库的一定要配置好flush参数。5.3 窗口乱序与迟到数据的处理交通数据是典型的乱序数据。GPS终端和卡口设备都依赖网络传输网络一抖动数据就会乱序到达。我们上线第一天就发现了这个问题某个路段明明已经拥堵但窗口计算出来的平均车速还是很高原因是早高峰的轨迹点延迟了十几秒才到被算进了后一个窗口。解决乱序问题的核心是Watermark和allowedLateness。我们在Kafka源表上设置了5秒的Watermark延迟容忍度WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND这表示允许事件时间比当前处理时间慢5秒以内的数据参与窗口计算。然后窗口侧设置ALLOWED LATENESS为30秒并开启侧输出流INSERT INTO doris_section_agg SELECT section_id, TUMBLE_START(event_time, INTERVAL 5 MINUTE) AS window_start, COUNT(*) AS sample_count, AVG(speed) AS avg_speed FROM matched_trajectory GROUP BY section_id, TUMBLE(event_time, INTERVAL 5 MINUTE);迟到超过30秒的数据会进入侧输出流我们额外写一个算子把这些数据单独处理或者直接丢弃。这里有个经验allowedLateness不能设太大否则一个窗口要等很久才能关闭窗口状态持续占用内存还会导致下游指标长时间不稳定。我们的场景里30秒够用如果你的设备网络质量比较差可以放宽到1分钟但不要超过2分钟。另外一个和乱序强相关的坑是Kafka的scan.startup.mode如果是latest-offsetFlink作业启动初期会因为等不到数据而拿不到Watermark窗口不触发。大屏上指标一直不更新。这时候不要慌给作业一点预热时间或者改成从最早位点消费一小段历史数据把Watermark顶起来。5.4 Checkpoint频繁失败RocksDB和状态膨胀的连锁反应项目稳定运行两周后突然有一天作业的Checkpoint频繁失败恢复时间越来越长背压一路飙升到90%。查了日志发现是RocksDB状态后端的状态太大了。我们有两个地方大量使用了状态。一个是轨迹数据里的车辆状态用来做异常停留检测另一个是坐标匹配过程中的上一匹配点缓存。因为当初没有给状态设置TTL每辆车的状态会一直保留两周下来状态文件已经占了几百GB每次Checkpoint都要把这些状态快照到HDFS时间长、失败率高。解决办法是给状态设置合理的TTL。Flink的状态API支持配置过期时间StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.hours(2)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptorTrajectoryState descriptor new ValueStateDescriptor(vehicle-state, TrajectoryState.class); descriptor.enableTimeToLive(ttlConfig);然后调整了TaskManager内存配置把RocksDB的taskmanager.memory.managed.fraction从默认的0.4提高到0.6给RocksDB更多的内存做缓存。这样调整之后Checkpoint恢复正常背压也降下来了。这里有一个容易被忽略的细节开启TTL之后状态数据并不会立刻变小要等过期数据被清理掉之后才会明显下降。所以上线之前就应该把TTL设计好不要等出了问题再补状态膨胀的恢复过程是很痛苦的。6. 常见问题与排查技巧速查表项目做完了我整理了一张问题速查表平时排查问题基本上就是对着这张表找方向。不一定覆盖所有情况但能解决大多数日常问题。现象可能原因解决办法实时指标突然变成0Kafka消费组发生Rebalance、Topic断流、窗口无数据触发检查消费组Lag如果Lag持续增长说明计算跟不上了如果Lag为0但指标为0那就检查上游设备窗口结果跳变严重Watermark配置不合理乱序数据处理不当确认源表事件时间字段是否解析正确适当调大Watermark延迟容忍度写入Doris报datev2/dateday类型不匹配Doris和Flink连接器类型映射不一致统一字段类型升级Doris连接器或将日期字段先用字符串表示JDBC写入超时、连接耗尽并行度过高、未开启批量flush降低并行度配置sink.buffer-flush.max-rows和intervalCheckpoint持续失败状态过大、RocksDB内存配置不合理设置状态TTL调大managed.fraction检查是否把Checkpoint放在了本地磁盘大屏查询越来越慢Doris没有按分区裁剪查询全表扫描按事件时间分区查询语句带window_start条件合理设置排序键Redis写入延迟高Sink未开启Pipeline、批量写规模太小改用Redis Sink的批量写入模式或者在下游用管道命令合并请求作业重启后状态丢失未配置分布式Checkpoint存储配置HDFS、S3或MinIO作为Checkpoint路径禁止使用本地磁盘某个路段指标长期不更新坐标匹配失败率过高GPS点没匹配到路段检查路网网格索引是否覆盖了GPS点的范围打印匹配失败样本分析阈值是否过小数据延迟持续累积并行度与Kafka分区数不匹配将Flink作业并行度与Kafka分区数对齐优先保证单分区被单并发消费排查这类问题我个人的习惯是先从结果反推链路。先看指标是不是为零、是不是跳变然后顺着血缘表找到对应的Flink作业再看作业的背压、Checkpoint、Kafka Lag这三个核心监控项。背压高说明计算能力不够Checkpoint失败说明状态或者存储有问题Lag持续增长说明消费速率跟不上生产速率。这三个监控项任何一个异常都能快速锁定问题的大方向不至于在代码层面大海捞针。这个项目上线之后我最深的体会是实时监控平台真正的难点通常不在算法而在链路。任何一环抖动——网卡问题、Kafka分区不均、状态膨胀、类型不匹配——都会传导到最终指标上而且排查起来比离线任务费劲得多。所以做这类项目监控告警和历史回放能力一定要提前做进去不要等上线之后再去补。再分享一个小建议新团队做实时项目不要一上来就上最复杂的技术栈。先用Flink SQL把主链路跑通让数据从Kafka顺利流到Doris大屏能看到基本指标再把坐标匹配、维表关联、精确一次这类深水区的功能一个个加上去。先把链路打通稳定比什么都重要。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进