ARTICLE DETAIL

资讯详情

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

物流大数据双轨架构实战:从Kafka到Doris的实时离线一体化

物流大数据双轨架构实战:从Kafka到Doris的实时离线一体化 1. 为什么物流场景是“实时离线”双轨架构的最佳试验场先交代一下背景我去年年中开始接手公司智慧物流大数据平台的搭建从需求梳理到集群规划再到链路开发全程踩了一遍。当时项目组就三个人一个负责后端接口一个负责前端大屏我负责整个数据架构和数仓建设。三个月时间把一套覆盖订单、轨迹、车辆、仓储四类核心数据的实时离线双轨平台跑通了这里把我的架构设计思路和实操细节整理出来。这套平台要解决的业务痛点很明确物流调度中心需要实时掌握全国网点的车辆位置、订单状态、异常滞留件大屏上要能看到秒级刷新的单量曲线而财务结算、运营复盘、线路优化又需要基于全量历史数据做离线分析。一个业务场景同时逼着你既要实时又要全量这正是当前超流行框架组合——Kafka、Flink、Hive、Spark、Doris 等最典型的应用战场。先说结论我最终采用的是 Lambda 架构的变体离线链路走 Hive 数仓分层 Spark SQL 批处理实时链路走 Kafka Flink Doris两条链路共用 ODS 层数据源在 Doris 里做数据汇合大屏和 BI 都从 Doris 取数。这套组合扛住了日均 1.5 亿条轨迹数据的压力实时指标延迟控制在 10 秒以内离线日批作业在凌晨 2 点前全部跑完。1.1 物流数据的四类天然矛盾不做物流行业的人可能觉得数据量不大但实际一梳理就知道坑在哪。我按数据特征把物流数据拆成四类每一类都有各自的处理难点第一类是订单状态数据。订单从创建、揽收、中转、派送到签收中间要经过十几个状态节点每个节点还带时间戳、操作人、网点ID。这类数据存储在业务库 MySQL 里QPS 不算高但状态变更频繁而且下游的实时大屏、实时预警都要依赖它。难点在于业务库不能直接扛分析查询必须做增量同步。第二类是GPS轨迹数据。全国几千台车、几万个快递员每台设备 5 到 10 秒上报一次经纬度算下来一天就是上亿条记录。这类数据是典型的时序数据量大、写入频繁、价值密度低但实时监控车辆位置、计算里程、判断偏航全靠它。难点在于写入吞吐和存储成本。第三类是仓储作业数据。出库、入库、盘点、拣货这些操作散落在 WMS 系统里字段多、口径杂同一个“出库成功”在不同系统里可能定义不一样。这类数据适合离线清洗和标准化实时性要求不高。第四类是外部接口数据。比如天气预报、交通管控、第三方运力报价这些数据格式不固定通过 API 拉取更新频率低但对线路规划和时效预测有辅助价值。四类数据的写入频率、数据量、实时性要求完全不同单一架构根本没法同时满足。这就是为什么物流平台一定要双轨——实时链路保证“看得见”离线链路保证“算得清”。1.2 为什么我选了 Lambda 而不是 Kappa业界处理实时离线的方案无非三种Lambda 架构、Kappa 架构、混合架构。我先排除了纯 Kappa——只保留实时链路所有历史数据重放都用 Kafka 解决。听起来很优雅但在物流场景不现实Kafka 消息默认只保留 7 天你要重放三个月的 GPS 轨迹做线路优化Kafka 存不下重新灌数据代价太高而且离线分析的复杂 SQL多表 Join、窗口函数、几十个维度的 olap 切片用 Flink SQL 写起来远不如 Hive/Spark 灵活稳定。我最后采用的是 Lambda 变体但做了两个关键改良一是ODS 层共用。不管是实时链路还是离线链路数据源都从 Kafka 消费同一份数据既落 HDFS 进 Hive也直接进 Flink 做实时计算从源头保证两条链路的数据口径一致。二是用 Doris 做数据汇合层。实时链路 Flink 计算完结果写入 Doris离线链路 Hive 跑完日批任务也写入 Doris业务方只对接 Doris 一个查询入口不需要关心数据是从实时来的还是离线来的。这个设计大幅降低了业务方的使用成本。2. 技术选型超流行框架的搭配逻辑与理由物流大数据平台的分工很清晰数据要先进得来、存得下、算得动、出得快。围绕这四件事我把选型做成了对比表格直接展示我在每个环节的思考过程。2.1 数据接入层业务增量用 Flink CDC日志采集用 Kafka 直连数据源类型推荐方案备选方案选型理由MySQL 业务库订单、仓储Flink CDCCanal KafkaFlink CDC 一条链路搞定采集解析支持断点续传不用额外维护 Canal 服务GPS 轨迹、App 日志Kafka 直连Filebeat Kafka设备端SDK直接写入Kafka减少中间环节降低延迟外部 API 数据DataX 周期性拉取SeaTunnelDataX 稳定成熟配置简单适合低频批量拉取日志和数据接入这块我提一下最容易犯的错很多人喜欢在 Kafka 前面再加一层 Flume 或者 Filebeat觉得这样可以缓冲。但 GPS 设备端上报走的是长连接 批量发送本身就有缓冲能力Kafka 直接扛写入完全没问题。多加一层只是多一个故障点没有实际收益。只要客户端 SDK 做重试和批量Kafka 写入端不需要额外代理。2.2 离线链路Hive 数仓 Spark SQL 批处理离线计算引擎我之前对比过 Hive on Tez、Spark SQL、Flink Batch最后选了 Spark SQL。原因有三一是物流数仓的 ETL 以 Hive SQL 为主Spark SQL 兼容 Hive SQL 语法迁移成本低二是 Spark 的资源复用做得好白天实时任务占用资源不多晚上批处理可以申请全部资源跑大任务调度上更灵活三是 Spark 对复杂 Join 的优化成熟处理亿级表关联不容易 OOM。Hive 我只让它承担最原始的 ODS 建表存储计算全部上推给 Spark。数仓存储这块文件格式选了 Parquet压缩格式选了 Zstandardzstd。我用同一份 7 天的 GPS 数据做过测试Parquet zstd 比 Parquet snappy 节省约 30% 存储压缩和解压速度几乎没有差别对于日增 20GB 的轨迹数据来说一个月能省下近 200GB 空间。2.3 实时链路Kafka Flink Doris 三件套这一套组合在 2024 年后基本是行业标配了几乎每个招聘 JD 里都会出现。我解释一下为什么是这三件套而不是其他替代品Kafka 负责消息缓冲和数据分发。选 Kafka 没什么争议吞吐量高、生态最完善、和 Flink 集成度最好。版本用的 3.5三个 Broker 节点扛了每秒 2 万条写入没问题。注意一下 Kafka 的分区数是关键调优参数我后文会专门讲。Flink 负责实时计算。选 Flink 而不用 Spark Streaming核心原因是物流场景大量依赖事件时间比如 GPS 上报时间晚于业务发生时间和状态管理比如计算车辆连续行驶时长Flink 的 Watermark 和 Checkpoint 机制处理这类问题最成熟。我们用 Flink SQL 写实时 ETL用 DataStream API 写复杂的状态计算两种方式在同一个作业里可以混用。Doris 负责实时OLAP查询。这个选型很多人问为什么不选 ClickHouse。我的真实对比结论是物流大屏和 BI 报表既需要大宽表的聚合查询也需要订单明细的点查比如查某个运单现在到哪了还需要实时更新订单状态从“运输中”改成“已签收”。ClickHouse 在聚合查询上性能极强但点查和实时更新是短板Doris 的 Unique Key 模型天然支持主键更新而且查询并发能力更好刚好命中我们的全部需求。核心组件的版本和功能说明我做了一个对照方便按需选用组件版本核心用途关键配置Kafka3.5消息总线log.retention.hours168单分区吞吐预估Flink1.18实时计算state.backendRocksDBcheckpoint间隔120sDoris2.1OLAP查询Unique Key模型分区分桶设计Hive3.1ODS存储Parquet zstd分区按天Spark3.5离线ETL动态资源分配shuffle分区自适应DolphinScheduler3.2任务调度工作流依赖失败重试3. 数据接入层订单、轨迹、车辆多源数据如何统一入湖入仓链路设计得再漂亮数据进不来都是空谈。这一节我详细讲三类主要数据源的接入方案每一步都是实操过的。3.1 业务库增量变更Flink CDC 同步订单和仓储数据订单数据在 MySQL 里如果每次全量同步单表 5000 万行的订单表要跑 20 分钟而且会对业务库造成压力。增量同步用 Flink CDC 是最省事的方案直接监听 MySQL 的 binlog把 insert、update、delete 变更实时捕获。具体这样搭在 Flink SQL 里定义一张 CDC 源表语法大概是CREATE TABLE order_cdc ( order_id BIGINT, status STRING, update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 192.168.1.10, port 3306, username flink, password ***, database-name logistics, table-name t_order, scan.startup.mode initial, debezium.snapshot.fetch.size 4096 );这里有一个关键点scan.startup.mode第一次跑用initial会自动先做一次全量快照再无缝切换成增量监听不需要手动处理全量与增量的衔接。后续重启任务用latest-offset只从当前 binlog 位置开始监听。但实际运行中我踩了一个大坑当订单表数据量超过 2000 万且业务库有大量历史归档数据时initial 模式做快照时会把整个表的数据扫描一遍导致业务库 CPU 飙升接近 100%。后来换成两台业务库一台专门承担 CDC 读取压力或者在低峰期凌晨 2 点初始化任务问题就解决了。3.2 高并发 GPS 轨迹日志的采集管道GPS 轨迹数据走的是 Kafka 直连生产端是车辆 T-Box 设备每台车 5 秒上报一条上报内容包括车辆ID、经纬度、速度、方向角、上报时间。高峰期 3000 台车同时在线每秒产生 600 条数据每条大约 150 字节对 Kafka 来说毫无压力。Kafka Topic 设计上我按业务域拆分而不是按数据类型拆分topic-gps-raw存原始 GPS 数据保留 7 天下游消费后清理topic-order-event存订单状态变更事件Flink CDC 写入后转发到这里topic-vehicle-status存车辆的在线/离线状态、实时位置聚合结果这个设计的好处是下游 Flink 任务只订阅自己关心的 Topic互不干扰而且 Kafka 的消费组机制天然支持多消费者并行处理。比如大屏服务需要实时位置直接消费topic-vehicle-status不用去读原始 GPS 流。3.3 接入层的幂等与去重设计物流数据接入层一个容易忽略的问题就是“数据重复”。设备断线重连会把缓存的上报记录重新发一遍Flink CDC 重启也会重复读取 binlog 最后一段这些都会造成数据重复。我的处理方案分三层第一层Kafka 生产端做幂等。设备端 SDK 上传时带一个唯一的message_idUUIDKafka 生产端开启enable.idempotencetrue这样同一批次内不会产生重复消息。第二层Flink 消费端做去重。Flink 作业里用状态存储过去 5 分钟内见过的message_id出现重复直接丢弃。RocksDB 状态后端存几百万个 ID 完全没压力。第三层存储层做去重兜底。Doris 的 Unique Key 模型用order_id或message_id做主键即使前面两层失效重复写入也会自动按主键覆盖查询结果不会有重复记录。三层都做了之后我验证过数据准确率能到 99.99%剩下那 0.01% 的差异主要来自跨天边界和时区处理不影响业务决策。4. 离线链路数仓分层模型与调度体系的搭建细节离线链路是物流平台数据分析的底座财务结算、线路优化、网点绩效考核全部依赖离线数仓。这一节给大家讲清楚数仓怎么分层、任务怎么调度、存储怎么治理。4.1 ODS/DWD/DWS/ADS 四层数仓模型标准的数仓分层设计我直接应用到了物流场景每一层的作用和表结构设计如下。ODS 层原始数据层Kafka 里的原始数据原封不动落到 HDFS按天分区。表名统一加ods_前缀字段和上游保持一致不做任何加工。这一层只做一件事保证数据不丢。比如ods_gps_trace表字段就是设备原始上报的那些字段分区是dt2024-06-01Parquet 格式zstd 压缩。DWD 层明细数据层对 ODS 层做清洗、去重、标准化形成业务口径一致的明细数据。这里要重点处理三类问题枚举值标准化。比如订单状态字段MySQL 里可能叫1、2、3、4需要映射成已创建、已揽收、运输中、已签收脏数据过滤。GPS 经纬度超出中国范围经度不在 73°E~135°E纬度不在 3°N~53°N的记录直接丢弃维表补充。把订单明细关联上网点名称、区域、城市等维度字段方便后续多维分析DWD 层典型表如dwd_trade_order_flow一个订单一行包含订单 ID、状态流转时间线、所属网点、区域、时效节点时间。DWS 层汇总数据层按业务维度做轻度汇总。比如按“城市日期小时”汇总订单量、GMV、妥投量、平均时效这些指标。这一层是查询性能的关键因为物流分析 80% 的报表查询都是按城市、网点的维度聚合提前聚合能减少大量计算。ADS 层应用数据层面向具体业务应用比如大屏指标、财务结算表、绩效考核表。这一层的表可以直接被业务方查询字段命名高度业务化比如ads_transport_daily_summary是“运输日报汇总”包含当日订单量、妥投率、异常占比等。4.2 DolphinScheduler 工作流设计与调度周期调度我用 DolphinScheduler 3.2没用 Airflow原因是 DolphinScheduler 对大数据任务的原生支持更好——直接拖拽编排 Spark、Flink、Hive 任务不需要额外写 Python 包装器。离线任务的调度周期分三档小时级任务每小时跑一次 DWD 层增量清洗。从 ODS 分区读最近 2 小时数据清洗后 append 到 DWD 对应表。为什么要跑最近 2 小时而不是 1 小时因为延迟上报的数据比如断网车辆恢复后补传 GPS经常晚到 1-2 小时读最近 2 小时可以尽可能把这些数据捞进来。日级任务每天凌晨 1 点开始跑 DWS 汇总和 ADS 报表。工作流依赖关系是ODS 层父任务检查 Kafka 落 HDFS 的数据完整性对比 Kafka 消息数和 HDFS 文件行数不一致则告警触发上游重推DWD 层任务依赖 ODS 任务完成后执行做全量清洗DWS 层任务依赖 DWD 完成后执行按维度汇总ADS 层任务依赖 DWS 完成后执行产出应用报表Doris 数据同步ADS 层完成后通过 Stream Load 方式把结果写入 Doris这 5 个任务串成一个工作流任何一个失败都会阻断下游DolphinScheduler 配置失败自动重试 2 次间隔 5 分钟。我们跑了半个月凌晨批作业的成功率稳定在 99% 以上。4.3 离线链路的存储优化与小文件治理离线链路最大的存储消费来自 ODS 层的原始数据。GPS 轨迹一天新增约 20GB加上订单、仓储日志一天总增量约 35GB一个月就是 1TB。如果不做治理半年后存储成本就会失控。我做了三件事控制存储一是压缩格式优化。前面提到的 Parquet zstd实测压缩比能达到 4.5:1一张 500GB 的 Hive 表压缩后只有 110GB。二是分区裁剪。查询时必须带分区条件禁止全表扫描。比如分析“6 月份华东区订单量”SQL 里必须写WHERE dt 2024-06-01 AND dt 2024-06-30 AND region 华东Spark 才能做到只扫描对应分区文件。如果业务经常要查长期趋势我就在 DWS 层额外做一张按月分区的汇总表查询走汇总表而不是扫 ODS 明细。三是小文件治理。Flink 写 HDFS 时默认并发度高会产生大量小文件。我设置了 StreamingFileSink 的sink.rolling-policy.rollover-interval 60min和sink.rolling-policy.inactivity-interval 30min强制文件按时间和大小滚动每个文件控制在 256MB 左右。另外每周跑一次小文件合并任务用 Spark 的repartition控制输出文件数量。5. 实时链路从 Kafka 到 OLAP 引擎的秒级数据管道实时链路是这套架构里技术含量最高的部分。实时指标的计算链路是Kafka 消费 → Flink 计算 → 结果写入 Doris 和 Redis → 大屏和接口读取。下面把作业设计细节展开。5.1 Flink 作业的拓扑设计与状态管理实时计算按业务场景拆成三个 Flink 作业每个作业独立部署、独立 checkpoint避免一个作业故障拖垮全部作业一实时订单状态流处理。消费topic-order-event清洗后关联维表网点表、区域表按订单 ID 分组输出订单实时状态、状态变更时间。写入 Doris Unique 模型表doris_order_realtime同时把“最近一小时内各城市订单量”的聚合结果写入 Redis给大屏接口实时访问。作业二车辆实时监控与偏航预警。消费topic-gps-raw按车辆 ID 分组用一个 Flink 状态保存每辆车当前的位置和上一位置计算速度和累计里程并判断是否偏离预设路线路线数据从 Redis 读取。发现偏航或车速超过 80km/h 时写入预警 Kafka Topic由预警服务推送告警。这个作业是状态管理最复杂的——每辆车都要保存一个移动窗口的状态我用 RocksDB 做状态后端将状态存储在本地磁盘配合每 10 分钟一次的增量 checkpoint既不占堆内存也能快速恢复。作业三大屏核心指标的实时聚合计算。消费订单状态流和 GPS 流通过 Flink SQL 按分钟粒度做累计聚合比如当前总订单量、全国平均妥投时效、在途车辆数、异常滞留件数。结果写入 Doris 预聚合表doris_realtime_indicator大屏前端每隔 5 秒轮询一次 Doris 接口。5.2 迟到数据和乱序数据的处理策略实时链路最常见的坑是数据乱序——物流设备传输延迟经常导致事件时间戳比处理时间早很多。比如一辆车 14:00 经过某地但 GPS 数据 14:03 才上传到服务器而同一车辆 14:02 的数据反而先到了。如果按处理时间计算就会得到错误的轨迹顺序偏航判断也会出错。我的解决方案是 Flink 的Event Time Watermark机制。在 Flink SQL 中定义 Watermark 策略CREATE TABLE gps_source ( vehicle_id STRING, lng DOUBLE, lat DOUBLE, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 30 SECOND ) WITH (...);这里 Watermark 设置为事件时间减 30 秒意味着最多容忍 30 秒的延迟数据进入窗口计算超过 30 秒的迟到数据会被丢弃。这个 30 秒不是拍脑袋定的是在分析了设备上报延迟分布后确定的95% 的数据延迟在 10 秒以内99% 在 30 秒以内所以设 30 秒可以在性能和准确性之间取得平衡。如果设太长窗口计算要等很久实时性变差设太短会有 1% 左右的数据被丢弃影响轨迹还原准确率。被丢弃的迟到数据我另外做了一个旁路Flink 侧输出流side output捕获被丢弃的数据写入另一个 Kafka Topic由离线链路兜底处理。这样即使实时指标漏算了一两条离线日批时也会修正回来。5.3 实时结果如何写入 Doris 和 Redis实时计算结果写存储时有讲究写得好不好直接决定了大屏的延迟和接口的响应速度。Doris 写入用 Stream Load。Flink 官方提供了 Doris Connector底层调用 Doris 的 Stream Load 接口支持微批写入。我配置的是每 15 秒或每 1 万条触发一次 flush这样既不会频繁建 HTTP 连接也不会因为缓冲太久导致大屏数据延迟。Doris 表模型的选择很关键。实时订单状态表我用Unique Key 模型以order_id作为唯一键订单状态每次变更都执行 upsert查询时永远拿到最新状态。大屏聚合指标表我用Aggregate Key 模型以indicator_code指标编码stat_time统计时间作为维度列指标值为 SUM 累加这样 Flink 每次写入一条增量数据Doris 自动累加到对应维度上大屏读的时候直接就是累计值不需要再聚合。Redis 设计成两级缓存。第一级缓存存大屏最近一次查询的完整 JSON 结果比如 5 秒内相同请求直接返回缓存不查 Doris第二级缓存存热门 Key 的明细比如某城市某网点的近期单量。Flink 作业每算出一批结果就主动更新 Redis保证接口取数永远是新鲜的。这样做的收益很明显Doris 的查询压力大幅下降大屏接口 P99 响应时间稳定在 200ms 以内。6. 集群部署与资源预估从单机 Demo 到生产环境很多人在做毕业设计或者小规模 Demo 时直接把生产部署方案往上套结果资源浪费严重、维护成本高。我做这套平台的集群规划时分了三套方案按阶段选择。6.1 三节点起步的最小集群角色分配如果是学习或者毕业设计一台机器也可以跑但体验很差——Kafka、Flink、Hive 挤在一起经常因为内存不够导致任务失败。我建议至少三台机器推荐配置如下节点配置部署组件说明Master 节点8核16GBNameNode、ResourceManager、Doris FE、DolphinScheduler负责集群管理、任务调度Worker1 节点8核32GBDataNode、NodeManager、Kafka Broker、Doris BE负责存储和计算Worker2 节点8核32GBDataNode、NodeManager、Kafka Broker、Doris BE、Flink TaskManager负责存储和计算这个配置下Kafka 只有 2 个 Broker分区副本数建议设成 2保证单点故障不丢数据。Flink TaskManager 的 Slot 数设 4跑三个实时作业没问题。运行内存规划上我给 HDFS 的 DataNode 分配 4GB给 Kafka 分配 6GBFlink TaskManager 分配 12GBDoris BE 分配 8GB剩下的给操作系统和自用。6.2 存储容量与计算资源配置的预估方法资源预估的猛药是“算清楚需要多少磁盘”。我按增量扩充的方式给出一套计算逻辑帮助大家不拍脑袋做判断。单个组件的存储预估方法是日增量 × 保留天数 × 文件副本数 × 压缩比系数。以 GPS 轨迹为例日增 20GB 原始数据保留 30 天HDFS 默认副本数为 3Parquet zstd 压缩比约为 1/4.5。所以轨迹数据的实际存储占用是20GB × 30 天 × 3 副本 × 0.22 ≈ 396GB。加上订单、仓储日志总存储需求约 800GB 到 1TB。给三节点集群配磁盘时我直接每台节点给了 4TB 的 SATA 盘看起来冗余很大但大数据集群最怕的就是磁盘不够。磁盘 IO 速度对实时链路的影响非常明显——Kafka 的页缓存Doris 的 compactionHDFS 的写入全部依赖磁盘性能。有条件就上 SSD至少也要 7200 转的 SATA 企业盘不要用笔记本盘。CPU 和内存的计算逻辑Flink 实时作业的单并行度大约需要 1.5GB 内存 1 核 CPU。我的实时作业总共设置了 20 个并行度因此需要 30GB 内存和 20 核 CPU。Spark 离线作业跑在晚上白天不占资源利用动态资源分配spark.dynamicAllocation.enabledtrue夜间最大可以申请 30 核 60GB 内存。两个计算框架叠加起来三台 Worker 节点的 32GB 内存刚好卡线够用。6.3 部署过程中的高频踩坑点集群搭建的坑太多了我挑三个印象最深的说每一个都是真金白银踩出来的坑一Kafka 分区数设少了Flink 并行度提不上来。Flink 消费 Kafka 的并行度受限于 Kafka 分区数分区只有 3 个Flink 算子的并行度设成 10 也没用只有 3 个并发在消费。后来我把订单事件 Topic 的分区扩到 12 个然后用 Kafka 自带的kafka-reassign-partitions.sh做数据重分布才解决消费瓶颈。Kafka 分区数建议按目标峰值吞吐量来定假设单分区吞吐 5MB/s业务峰值需要 50MB/s那至少 10 个分区留出两倍余量设 20 个。坑二Flink Checkpoint 频繁超时任务一直在重启恢复循环。检查后发现是 checkpoint 存储路径被放在了 HDFS 上而 HDFS 的 NameNode 内存不够频繁 GC 导致 checkpoint 提交超时。后来把 checkpoint 存储改成本地 RocksDB 定期备份到 HDFS问题解决。Flink 默认的state.backend.incrementaltrue要开增量 checkpoint 比全量快很多特别是状态大的时候。坑三Doris BE 节点 OOM查询直接把大屏卡死。原因是 Doris 的 Buffer Pool 和 Page Cache 默认配置太大两台 BE 各配了 50% 内存做缓存在并发大查询时会被打爆。调整buffer_pool_size到 20% 总内存storage_page_cache_limit设成 15% 总内存并加上查询超时限制query_timeout30s之后没有出现过这个问题。7. 大屏可视化与业务应用的落地效果架构的最终价值要体现在业务使用上。这一节讲一讲物流驾驶舱大屏的指标设计和数据刷新方案算是给整套架构做一个应用层落地的收尾。7.1 物流驾驶舱的核心指标设计大屏指标的选取不是随便放的要能回答调度中心最关心的三个问题单量够不够、时效稳不稳、车辆在不在。指标分类具体指标实时/离线数据来源单量监控今日实时订单量、昨日同期单量、环比实时Doris 实时聚合表时效监控平均妥投时长、超时件数、24小时妥投率实时Doris 实时聚合表车辆监控在线车辆数、空闲车辆数、行驶总里程实时Redis 缓存车辆状态网络质量各网点妥投率排名、异常网点 Top10离线Doris 日汇总表库存周转各仓库存量、出库订单数、库存周转天数离线Doris 日汇总表这里有一个细节大屏展示“异常滞留件”指标时我做了实时和离线两条链路的对比发现实时链路算出的滞留件数总比离线少 1% 左右。原因很隐蔽——实时链路判断滞留用的是“当前时间 - 最近状态变更时间”而离线链路用的是“业务日期截止 23:59:59 的状态”跨天数据在实时和离线两条链路里的归属日期不同。后来统一成“超过 48 小时未更新状态则判定为滞留”两条链路的数据差异就缩小到千分之一以内。7.2 ECharts 大屏的实时数据刷新方案大屏前端用的是 Vue ECharts数据刷新方案是“轮询 WebSocket 推送”的组合设计。轮询方案对于核心指标今日订单量、在途车辆数、妥投率前端每 5 秒轮询一次后端接口后端先从 Redis 取数Redis 没有再到 Doris 查询并回填缓存。这种方式实现简单延迟取决于轮询间隔5 秒刷新对调度大屏来说完全够用。但轮询有天然缺陷如果某个时刻后端 Doris 重启或者慢查询接口响应超过 5 秒前端就会积压一堆 pending 请求把服务拖垮。所以我给接口加上了 2 秒超时控制超时直接返回上一次的缓存值前端显示“数据更新于 XX 秒前”来提示业务方。WebSocket 推送方案对于实时预警类数据偏航预警、温度异常预警、车辆故障用 WebSocket 做服务端主动推送。Flink 检测到异常后写入 Kafka 预警 Topic预警服务消费后通过 WebSocket 实时推送给大屏前端弹出预警卡片并配合地图定位展示异常车辆。这个方案比轮询的准实时性更高而且服务端可以精确控制推送频率不会像轮询那样出现大量无效请求。前端不推荐的方案很多博客喜欢用“定时器 Ajax 轮询”来做大屏数据刷新我给个结论数据量小的时候可以一旦指标多、图表多还是老老实实上 WebSocket——不是因为性能而是因为可维护性和服务端压力好控制得多。最后分享一个运维上的小心得实时和离线两套链路一定要每天都对账。第二天早上离线日批跑完后写一个对账脚本对比“离线数仓统计的昨日订单量”和“实时 Doris 里昨天累计的订单量”差异超过 0.5% 就告警。这个动作能帮你尽早发现 Flink 丢数据、Kafka 重复消费、Doris 覆盖写入失败等隐藏问题。我上线的第三周就是靠这个对账脚本发现了一个 Flink Checkpoint 恢复后重复消费的 bug——实时订单量虚高 2%离线对出来数字不对才定位到。没有对账这个 bug 可能要在业务方投诉之后才会被发现。
返回列表