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 CDC | Canal + Kafka | Flink CDC 一条链路搞定采集+解析,支持断点续传,不用额外维护 Canal 服务 |
| GPS 轨迹、App 日志 | Kafka 直连 | Filebeat + Kafka | 设备端SDK直接写入Kafka,减少中间环节,降低延迟 |
| 外部 API 数据 | DataX 周期性拉取 | SeaTunnel | DataX 稳定成熟,配置简单,适合低频批量拉取 |
日志和数据接入这块,我提一下最容易犯的错:很多人喜欢在 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,压缩格式选了 Zstandard(zstd)。我用同一份 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 模型天然支持主键更新,而且查询并发能力更好,刚好命中我们的全部需求。
核心组件的版本和功能说明,我做了一个对照,方便按需选用:
| 组件 | 版本 | 核心用途 | 关键配置 |
|---|---|---|---|
| Kafka | 3.5 | 消息总线 | log.retention.hours=168,单分区吞吐预估 |
| Flink | 1.18 | 实时计算 | state.backend=RocksDB,checkpoint间隔120s |
| Doris | 2.1 | OLAP查询 | Unique Key模型,分区分桶设计 |
| Hive | 3.1 | ODS存储 | Parquet + zstd,分区按天 |
| Spark | 3.5 | 离线ETL | 动态资源分配,shuffle分区自适应 |
| DolphinScheduler | 3.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_id(UUID),Kafka 生产端开启enable.idempotence=true,这样同一批次内不会产生重复消息。
第二层,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表,字段就是设备原始上报的那些字段,分区是dt=2024-06-01,Parquet 格式,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核16GB | NameNode、ResourceManager、Doris FE、DolphinScheduler | 负责集群管理、任务调度 |
| Worker1 节点 | 8核32GB | DataNode、NodeManager、Kafka Broker、Doris BE | 负责存储和计算 |
| Worker2 节点 | 8核32GB | DataNode、NodeManager、Kafka Broker、Doris BE、Flink TaskManager | 负责存储和计算 |
这个配置下,Kafka 只有 2 个 Broker,分区副本数建议设成 2,保证单点故障不丢数据。Flink TaskManager 的 Slot 数设 4,跑三个实时作业没问题。运行内存规划上,我给 HDFS 的 DataNode 分配 4GB,给 Kafka 分配 6GB,Flink TaskManager 分配 12GB,Doris BE 分配 8GB,剩下的给操作系统和自用。
6.2 存储容量与计算资源配置的预估方法
资源预估的猛药是“算清楚需要多少磁盘”。我按增量扩充的方式给出一套计算逻辑,帮助大家不拍脑袋做判断。
单个组件的存储预估方法是:日增量 × 保留天数 × 文件副本数 × 压缩比系数。
以 GPS 轨迹为例:日增 20GB 原始数据,保留 30 天,HDFS 默认副本数为 3,Parquet + zstd 压缩比约为 1/4.5。所以轨迹数据的实际存储占用是:20GB × 30 天 × 3 副本 × 0.22 ≈ 396GB。加上订单、仓储日志,总存储需求约 800GB 到 1TB。
给三节点集群配磁盘时,我直接每台节点给了 4TB 的 SATA 盘,看起来冗余很大,但大数据集群最怕的就是磁盘不够。磁盘 IO 速度对实时链路的影响非常明显——Kafka 的页缓存,Doris 的 compaction,HDFS 的写入,全部依赖磁盘性能。有条件就上 SSD,至少也要 7200 转的 SATA 企业盘,不要用笔记本盘。
CPU 和内存的计算逻辑:Flink 实时作业的单并行度大约需要 1.5GB 内存 + 1 核 CPU。我的实时作业总共设置了 20 个并行度,因此需要 30GB 内存和 20 核 CPU。Spark 离线作业跑在晚上,白天不占资源,利用动态资源分配(spark.dynamicAllocation.enabled=true),夜间最大可以申请 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.incremental=true要开,增量 checkpoint 比全量快很多,特别是状态大的时候。
坑三:Doris BE 节点 OOM,查询直接把大屏卡死。原因是 Doris 的 Buffer Pool 和 Page Cache 默认配置太大,两台 BE 各配了 50% 内存做缓存,在并发大查询时会被打爆。调整buffer_pool_size到 20% 总内存,storage_page_cache_limit设成 15% 总内存,并加上查询超时限制(query_timeout=30s),之后没有出现过这个问题。
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 可能要在业务方投诉之后才会被发现。