工作这么多年大数据平台越做越复杂见过不少团队从一套 Hadoop 集群起步跑着一堆 T1 的离线任务后来业务一催又要接实时报表又要支持分析师随时跑 SQL最后把 Spark、Flink、Presto、Doris 全堆在一个集群里任务倒是都能跑但资源互相打架、口径对不齐维护成本翻着跟头涨。我们团队也走过这段弯路最后沉淀下来的方案就是围绕“混合计算模式”来做架构设计和任务编排。这篇就把我对这个模式的理解、选型思路、落地步骤和踩坑记录都写出来给正在折腾大数据平台的同学一个参考。1. 混合计算模式到底在解决什么问题1.1 三种典型计算场景为什么不能一把梭大数据计算往前推十年基本就是 MapReduce 一把梭写个 MR 跑离线统计等上几个小时很正常。后来 Spark 出来了离线计算快了一个量级但仍然解决不了“数据实时性”的需求。业务方要的是“此刻的订单量”“最近五分钟的转化率”不可能等第二天凌晨再跑批。于是流计算引擎又成了刚需。再往后从数仓取数、BI 报表拖拽、临时分析又催生了交互式查询引擎的需求。这就形成了三种典型场景批处理数据量大、容忍分钟级到小时级延迟算全量、跑历史、做月度汇总典型代表是凌晨跑数的离线任务。流处理数据实时到达需要毫秒级到秒级的处理延迟持续运行做累计、做告警、做实时大屏。交互式查询响应时间要求秒级甚至亚秒级数据规模大但往往查的是维度聚合或明细拿数分析人员直接写 SQL。难点在于这些场景往往同时存在于同一家公司、同一个数仓、甚至同一张业务表上。拿订单数据来说白天要实时看 GMV晚上要跑全量结算分析师随时要按城市、类目、时间段切片查数据这就容不得你搞三套完全隔离的环境数据既不一致成本也扛不住。混合计算模式的核心思路就是让批、流、交互查询在同一个数据底座上协同工作用不同的引擎处理各自擅长的负载同时又保证数据的一致性、可追溯性和资源可控性。1.2 一套引擎打天下的幻想与现实的差距不少团队一开始的想法是“选一个万能引擎统一批和流”。前几年 Flink 刚喊“批流一体”的时候很多人觉得曙光来了认为能一个引擎吃掉所有负载。但真实落地下来你会发现“批流一体”解决的是“一套代码两种执行模式”的问题而不是“一个引擎擅长所有负载”的问题。Flink 跑批确实可以但跟 Spark 比在 shuffle 优化、动态资源调整、成熟稳定的批处理生态方面还是有差距让 Spark 处理毫秒级延迟的流任务用 Structured Streaming 做实时指标延迟能到秒级就不错了事件时间处理和状态管理都比 Flink 笨重。至于交互式查询两个引擎都不适合直接承担高并发的点查和即席查询那是 OLAP 引擎的主场。强行统一的结果就是每个场景都用得别扭性能和稳定性两头都不占。现实的架构是让不同引擎做各自擅长的事但共享同一份存储和元数据。批处理用 Spark流处理用 Flink交互式查询用 Presto/Trino 或者 StarRocks/Doris存储层落在 HDFS/对象存储 数据湖表格式上用统一的 Catalog 管理元数据。这样既能发挥各引擎的长处又能让数据“写一次、多处读”避免烟囱式开发。1.3 Lambda 与 Kappa混合计算的两条设计路线谈到混合计算绕不开经典的 Lambda 架构和 Kappa 架构之争。Lambda 架构的做法是批处理链路维护一份全量的离线结果实时链路维护一份增量准实时结果查询时由服务层合并两份数据。这种架构上线容易但代价是两套代码、两套口径、两份结果一旦指标定义对不上批流数据会打架排查成本极高。Kappa 架构更激进只用一套流处理链路把数据当成无限流需要重算时从 Kafka 等消息队列里回放数据。逻辑上很干净真正落地时却会被消息队列的保存时长、回放速度、状态恢复复杂度卡住尤其是要批量重算半年历史数据时从 Kafka 里回放简直是个灾难。我在实际工程里的体会是别把 Lambda 和 Kappa 当非此即彼的选择而是按业务诉求做融合。整体链路可以是 Kappa 的思路——批流共用同一套数据湖表、同一套指标定义重算场景则用 Hive/Iceberg 的存量历史数据配合 Spark 批量重跑而不是非要流引擎回溯。这个落地方案本质上是“存储底座统一 计算引擎按需选型”也是我理解中混合计算模式比较务实的形态。2. 混合计算架构怎么选型计算引擎与存储底座2.1 引擎选型对照Spark、Flink、Presto/Trino、StarRocks/Doris先整理一张选型对照表方便不同角色的同学快速对齐。计算引擎擅长的场景延迟特征典型用法短板Spark离线批处理、大规模 ETL、复杂数据清洗分钟级及以上Structured Streaming 可到秒级T1 数仓分层加工、历史数据重算纯流场景延迟不够低状态管理较弱Flink实时流处理、事件驱动计算毫秒~秒级实时指标、实时大屏、实时数仓 DWD/DWS 层批处理生态相对年轻大批量 shuffle 调优成本高Presto/Trino交互式查询、联邦查询、跨源分析秒级~分钟级Ad Hoc SQL、BI 报表取数、数据探查不擅长高并发点查大查询容易压垮 CoordinatorStarRocks/Doris高并发 OLAP、明细与聚合查询秒级甚至毫秒级报表服务、多维分析、用户画像查询数据导入与更新需要配套链路不太适合复杂的流式 ETL选择的原则是第一看任务特征不要跟风。如果你的实时需求只是分钟级Spark 就够了非要上 Flink 增加维护成本如果你有大量即席 SQLPresto 是真香但要让 BI 报表支撑高并发还是把数据同步到 Doris/StarRocks 更稳。第二看团队的技术储备。Flink 对状态编程和理解 checkpoint 机制有要求团队成员不熟悉的话运维事故会止不住。第三看底层存储。引擎最好能直接读同一份数据湖表避免各种导入导出链路。我见过不少团队把 Presto 当万能查询引擎给线上 BI 用上千并发压过来Coordinator 直接打满查询超时一片。后来把核心聚合结果放开到 DorisPresto 退回到 Ad Hoc 场景问题才缓解。这就是选型时没分清“即席探索”和“高并发服务”的差异。2.2 让批和流共用一份数据数据湖格式的角色混合计算模式里最关键的底座是数据湖表格式目前主流是 Iceberg、Hudi、Delta Lake三选一都行但必须选。为什么因为裸 HDFS 上的 Hive 表支撑不了“批流共用一份数据”这件事。流任务和批任务如果各写各的表那本质上还是两条烟囱只是把“两份结果”换成了“两份底层表”。而有了 Iceberg 这类表格式就能做到ACID 语义多任务并发写同一张表不会互相读到脏数据。快照隔离与时间旅行读任务永远读到某个快照写任务提交新快照互不阻塞出问题可以回到任意历史版本排查。增量读取流任务可以像读消息队列一样读表里新增的数据文件这让“流读批写”成为现实——Spark 凌晨批量写全量数据Flink 白天读取新增数据做实时汇总。隐藏分区和分区裁剪优化不需要人工维护分区字段查询性能更可控。举个例子实时数仓里最常见的“ods 层 Kafka 数据落 HDFS”这一步以前要么 Flink 写入一个 Hive 分区目录要么直接写 Kafka 备查。用了 Iceberg 之后Flink 把 Kafka 数据以分钟级微批写入 Iceberg 表Spark 再基于 Iceberg 做每日全量刷新和汇总计算两张引擎访问的是同一张表的同一份文件不存在拷贝和同步问题。2.3 分层架构与数据流向设计落到实际架构上我习惯按四层来组织接入层业务日志、数据库 Binlog、App 埋点统一进 KafkaKafka 既是流处理的数据源也是离线数据落地的第一站。存储层以 HDFS 或对象存储为主所有结构化数据用 Iceberg/Hudi 表格式管理。这里既要存 ODS 原始数据也要存 DWD/DWS 加工结果。计算层Spark 负责离线加工Flink 负责实时链路Presto/Trino 负责即席查询StarRocks/Doris 负责服务化查询和高并发报表。调度与治理层统一用 DolphinScheduler/Airflow 做离线调度Flink 任务由 Flink 自带的作业管理配合平台化工具管控元数据用 Hive Metastore 或 Iceberg Catalog 统一管理权限和资源队列在这个层面统一控制。数据流向按“一份数据、多处计算”来设计。Kafka 里的数据Flink 直接消费做实时指标同时一个旁路任务把 Kafka 数据以列式格式写入 Iceberg ODS 表供批处理使用。批处理产出的全量汇总结果写回 Iceberg DWS 层实时链路产出的增量结果也写进同一张 DWS 表或姊妹表。查询层根据业务对实时性、并发度的要求决定从 StarRocks 读还是直接从 Iceberg 用 Presto 读。这个设计避免了“流批各搞一套”的物理隔离但也没有天真地让所有负载挤在同一个引擎上。它的巧妙之处在于存储统一、元数据统一计算按需路由写路径是同一个抽象读路径则根据场景分流。3. 实操落地一个可复现的混合计算平台案例3.1 业务场景和整体目标为了说清楚怎么落地我用一个典型的网约车数据分析平台来举例。业务方需求很明确白天要看实时订单量、完单率、各城市实时 GMV每晚需要跑全量的日结算报表分析师随时要按城市、时段、司机等级、车型等多个维度组合查询历史数据。这是一个再典型不过的混合计算场景。平台目标可以拆成四点实时指标延迟控制在分钟级核心指标 1 分钟内可见。离线任务T1 全量结算和统计早晨 8 点前必须产出。即席分析分析师直接在平台写 SQL返回时间控制在 30 秒内覆盖近一年的历史数据。高并发报表线上看板支撑几百 QPS核心指标秒级响应。存储层统一走 Iceberg 表底层是 HDFS。离线引擎用 Spark实时引擎用 Flink即席查询用 Trino服务化查询用 StarRocks。调度用 DolphinScheduler元数据统一进 Iceberg Catalog 并同步到 Hive Metastore。3.2 表设计ODS/DWD/DWS 三层怎么划表结构设计是三段式。ODS 层ods_order_info存订单原始数据分区字段为天由 Flink 从 Kafka 写入 Iceberg。这一层的价值是保留完整原始数据供重算、排查、回溯使用。为了控制成本和文件数量Flink 写入时做 10 分钟一个 checkpoint、5 分钟一次 commit。DWD 层dwd_order_clean存清洗、过滤、规范化之后的明细数据离线链路用 Spark 每天凌晨跑一次实时链路 Flink 从 ODS 增量读数据做同样的清洗逻辑后追加写入。DWS 层dws_city_order_stats_daily和dws_city_order_stats_realtime分别存天级汇总和分钟级汇总。两张表的指标口径完全一致只是统计粒度不同。这里最核心的细节是口径统一。举例指标“完单率” 完单量 / 订单量订单量定义为下单时间在当日的订单数完单时间在数据产出时刻算完单。无论实时还是离线这个定义写死在 DWD 层的 SQL/代码里不允许下游各算各的。3.3 离线批量链路的具体实现离线链路从 Kafka 落 Iceberg 开始。这一步我用 Flink 做配置大概是CREATE CATALOG iceberg_hdfs WITH ( typeiceberg, catalog-typehive, urithrift://metastore:9083, clients10, warehousehdfs://nameservice/warehouse ); CREATE TABLE IF NOT EXISTS iceberg_hdfs.dw.ods_order_info ( order_id STRING, city_id INT, driver_id STRING, passenger_id STRING, order_status STRING, amount DECIMAL(10,2), order_time TIMESTAMP(3), finish_time TIMESTAMP(3), partition_day STRING ) PARTITIONED BY (partition_day) WITH ( write.format.defaultparquet, write.target-file-size-bytes268435456 );Flink 写入时用 Kafka 的partition_day作为分区字段按事件时间分配。这里做过一次重要调优一开始write.target-file-size-bytes用默认值结果很多小文件后来统一设成 256MB文件数量立刻降下来后续 Spark 读的效率明显提升。Spark 每天凌晨读 ODS 表做清洗核心逻辑是过滤掉恶意刷单、状态异常的记录然后写 DWD 表val odsDF spark.read.table(dw.ods_order_info) .filter($partition_day bizDate) .filter($order_status ! invalid) odsDF .select( $order_id, $city_id, ... // 规范化字段 ) .writeTo(dw.dwd_order_clean) .option(partitionBy, partition_day) .append()DWD 到 DWS 的汇总就是把订单明细按城市和小时做聚合写进 DWS 天级表。这里用 Spark 的动态资源分配spark.dynamicAllocation.enabledtrue同时设了spark.sql.shuffle.partitions200避免小文件多和资源空转的问题。3.4 实时计算链路的具体实现实时链路的源头同样是 Kafka。Flink 消费订单消息实时清洗和维表关联后用 1 分钟滚动窗口做城市级聚合把结果同时写到两个地方StarRocks供线上看板查询和 Iceberg DWS 实时表供对账和历史重建。写入 StarRocks 用官方 StarRocks Connector写入 Iceberg 用 Iceberg Flink Connector。Flink 部分的关键配置是我反复踩坑后定下来的execution.checkpointing.interval: 60s execution.checkpointing.min-pause: 30s execution.checkpointing.tolerance-failed-checkpoints: 3 state.backend.type: rocksdb state.checkpoints.dir: hdfs://nameservice/flink-checkpoints/order_stats/checkpoint 间隔设 60 秒min-pause 设 30 秒避免 checkpoint 过于频繁把 IO 打满。状态后端用 RocksDB因为窗口聚合的状态量不算小放在堆内存里容易 OOM。这个配置上线后窗口聚合的实时性稳定在 1 分钟内故障恢复最多丢 1 分钟数据业务可接受。实时链路还有个容易忽略的细节维表关联。城市名称、司机等级这类维度信息变动不快我用 Flink 的lookup join关联 MySQL 维表但加了缓存缓存 30 秒刷新一次。如果不加缓存维表的高频查询会反压到业务库很危险。3.5 统一查询与指标层查询层分开两条路。分析师即席查历史明细直接用 Trino 查 Iceberg 表SELECT city_id, count(*) AS order_cnt, sum(amount) AS gmv FROM dw.dwd_order_clean WHERE partition_day 2024-01-01 GROUP BY city_id ORDER BY gmv DESC;数据量大时 Trino 的分布式查询能撑住30 秒返回基本没问题。高并发报表则查 StarRocks 里实时链路的聚合结果响应时间稳定在 200ms 上下。用 StarRocks 的 External Catalog 也能直接映射 Iceberg 表但实际项目里我更推荐“计算结果落到 StarRocks 内表”因为查询性能差距很大且日报表更新逻辑也更可控。实时、离线两份 DWS 数据Day 级数据每天凌晨由 Spark 产出分钟级数据由 Flink 持续写入 StarRocks。为了对账每天凌晨跑一个校验任务比对实时表和离线表前一天的数据差异。差异率超过阈值就告警定位是口径问题还是链路问题。这个对账任务看着不起眼实际是流批融合场景的定海神针。4. 常见问题与排查技巧实录4.1 小文件失控NameNode 内存告警混合计算场景的小文件问题比单纯离线场景更严重。因为实时任务会高频 commit 文件如果 Flink checkpoint 太频繁、提交间隔太短、或者目标表文件大小设置不合理小文件就会飞快累积。排查时先看表和目录的文件数再定位写入任务。根据我的经验80% 的小文件问题都可以通过以下几点解决Flink 写入时调大write.target-file-size-bytes我常用 256MB。适当拉长 commit 频率比如从每 2 分钟提一次改成每 5~10 分钟提一次。离线跑完后对分区做一次合并用 Iceberg 的rewrite_data_files动作合并小文件。对特别碎的分区定期跑一次 compaction 任务。Iceberg 本身提供rewrite_data_files存储过程写起来也简单关键是把它纳入调度CALL dw.system.rewrite_data_files( table dw.ods_order_info, where partition_day current_date, options map(min-input-files, 5, target-file-size-bytes, 268435456) );4.2 流批任务口径对不上报表打架这是混合计算模式里最伤脑筋的问题。实时报表明明显示今日 GMV 1.2 亿第二天离线报表却显示 1.15 亿差了 500 万业务方直接找过来。靠谱的解决思路有三步指标定义先行。任何指标在开发前必须有统一的定义文档包括统计口径、时区、业务判定规则。例如订单量按“下单时间去重”算撤销订单不计入这类规则绝不允许各链路自行解释。同源同逻辑。流和批从同一份 ODS 读取数据执行同一套清洗逻辑最好同一份代码模板生成的 SQL。实时链路如果用 Flink SQL离线链路用 Spark SQL那么尽量让 SQL 片段保持一致不要一版一个写法。强制对账。每日对账任务雷打不动。对账表就是 DWS 层实时结果和离线结果按天比对发现差异自动标记。对不上先查“晚到数据”再查业务变化最后查代码逻辑版本。我见过把对账做得很漂亮的团队对账任务以“实时链路结果”为基准差异超过 0.5% 才告警低于阈值看趋势如果连续一周偏差很小就把预警阈值收紧。这就是在成本和风险之间做平衡。4.3 资源抢占导致实时任务延迟YARN 或 K8s 集群如果共用一套资源池白天的实时任务和夜晚批处理任务之间偶尔会互相影响尤其是大查询或大批任务突然占满资源时Flink 作业背压飙升延迟层层放大。我的对策是三层隔离队列隔离给 Flink 实时任务单独分一个队列配高优先级离线批任务走自己的队列限制最大资源水位。错峰调度重要的批量重算任务放在凌晨低峰期避开实时任务的高峰窗口。动态资源可控Spark 开启动态资源分配设好最大 executor 数量避免晚上资源被一顿吃光。还有个小技巧Flink 任务的并行度和资源大小要按峰值流量评估不要卡着平均值配否则流量一波动就完蛋。实时链路最忌讳“降级式”运维宁可平时空闲一点也要保证峰值不崩。4.4 历史数据重算与补数业务规则调整之后历史一年的数据要重算这是混合计算最考验底座的场景。传统 Hive 表补数要 drop 分区、重建分区、跑任务中途出错还会污染数据。我在 Iceberg 上体验就完全不同用快照机制先确认当前表的最新快照 ID然后用 Spark 直接读取历史某个时间点的快照做重算算出的新结果生成一个新快照业务方验证无误后再切换到新快照。切换前还可以随时回滚旧快照容错空间大得多。补数的 Spark 任务注意两点第一写入模式要保证幂等overwrite方式只覆盖指定分区或指定快照别误伤其他分区第二重算任务要设置合理的资源配额因为补数往往比日常批处理数据量大几倍防止它把实时链路资源也拖垮。4.5 一张问题排查速查表现象可能原因快速排查手段解决方案实时看板数据延迟变大Flink 背压或 checkpoint 超时查看 Flink Web UI 的 Backpressure 和 Checkpoint 指标调整并行度、增加资源、检查外部系统 IO离线日常任务突然变慢源表小文件多、HDFS 读放大查看执行计划扫描文件数与大小执行 rewrite_data_files 合并检查分区裁剪是否生效流批报表数字不一致指标口径不一致或晚到数据比对对账任务明细查 ODS 数据到达时间分布统一口径、修正流批逻辑、加长实时链路 watermark 容忍度查询结果长时间不返回Trino 大查询占满资源查看 Trino query details 与资源组状态设置资源组、拆分大查询、收敛数据范围StarRocks 集群 CPU 飙升大量精确去重或没走物化视图查看慢查询和大查询队列改造 SQL、建物化视图、常用查询结果落表这些坑基本都是工作里长期反复出现的不是一篇文章能覆盖完的但把这五个点治理好整个平台的稳定性就能上一个台阶。5. 个人实操心得与后续扩展5.1 从 0 到 1 落地混合计算模式的顺序建议如果你所在的团队正打算上一个数据平台或者想把现在的离线数仓升级成批流混合架构我建议别急着把 Spark、Flink、Presto、Doris 一次性全铺上那大概率会造成运维灾难。稳健的顺序是先把存储底座和表格式定好。HDFS 或对象存储二选一然后选定 Iceberg/Hudi/Delta Lake 其中一个把 Hive Metastore 或对应 Catalog 跑起来。这一步是整个方案的地基后面所有引擎都要连它换起来最伤筋动骨。离线链路先跑通。用 Spark 接 Iceberg把 ODS/DWD/DWS 三层建好跑一段时间的 T1 任务确认数据质量和调度稳定性。期间把权限、队列、元数据管理这些基础设施补齐。引入实时链路。Kafka 到 Flink 到 StarRocks 的这条线可以单独搭核心是让 Flink 写 ODS 的实时积累和 Spark 读 ODS 做离线清洗互不冲突。这一步先做一两个核心指标不要全量并行。最后再上即席查询和服务化查询。Trino 接 Iceberg 表实现 Ad Hoc然后再把常用 DWS 结果同步到 StarRocks 供线上报表用。这时候你已经有了完整的数据链路查询这层做起来水到渠成。每步之间要有足够的观察期。代码可以快速写但稳定性需要时间验证尤其是 Flink 的状态恢复和 Iceberg 的快照管理多跑几周才能摸清脾气。5.2 更进一步的思路往后扩展可以从两个方向走。一是引入数据质量平台把 DWD 层的完整性、唯一性、及时性做成自动校验跑批和实时都要过质量门槛不合格数据不进 DWS。这套机制越早做后期省心越多。二是尝试用 Catalog 来做真正的“数据资产化”把所有表模型、血缘、指标定义都维护进统一元数据中心业务方自助找数、自助分析才不会天天来问你“这张表跟那张表有什么区别”。还有一个方向是流批一体的趋势演进。Flink 和 Spark 都在向批流一体靠近未来未必需要维护两条物理链路。但至少现阶段混合计算模式依然是工程上最稳妥的答案——让合适的引擎干合适的活底座统一、口径统一、流程受控这条路能支撑绝大多数企业级数据平台的体量。我个人实际用下来最大的感触是一次性地把“存储统一、口径统一、资源隔离、对账兜底”这四件事做好比追什么“新架构概念”有用得多。别急着上新技术先把四条链路各自跑稳再谈融合这个顺序倒了后期填坑会填到怀疑人生。