生而为人

程序员的自我修养

0%

实时数仓

[toc]

简历介绍

美团外卖实时数仓建设与优化项目

项目周期:2025.03-2025.10

技术栈:Flink 1.17、Kafka 3.4、ClickHouse 23.3、Canal 1.1.7、Redis 7.0、Spark 3.3、DolphinScheduler、Superset

项目背景:美团外卖日均订单量超 6000 万,峰值 QPS 达百万级,原有离线数仓延迟高(T+1),无法满足实时运营、骑手调度、风控拦截等核心业务需求。

核心职责

  1. 负责整体架构设计,采用纯 Kappa 架构替代传统 Lambda 架构,实现批流一体,统一实时与离线数据口径
  2. 设计并实现 ODS/DWD/DWS/DIM/ADS 五层实时数仓,完成订单、用户、商家、骑手四大核心业务域的数据建模
  3. 开发核心 Flink 作业,包括实时订单统计、骑手运力监控、用户行为分析、实时风控等 15 + 个核心业务流程
  4. 解决大流量峰值下的数据倾斜、大维度关联、状态膨胀等关键技术难题,保障系统稳定性
  5. 建立完善的数据质量监控体系和运维流程,实现数据异常自动告警和快速恢复
  6. 搭建实时数据服务平台,为全国运营大屏、商家后台、骑手 APP 等 20 + 个业务系统提供数据支持

核心成果

  1. 性能提升:核心业务指标延迟从原来的 30 分钟降至5 秒以内,支持分钟级业务决策
  2. 稳定性保障:系统可用性达到99.99%,成功支撑 2025 年 618、双 11 等大促活动,峰值订单量突破 1 亿单 / 天
  3. 资源优化:通过分层预聚合、大字段拆分 Join 等优化手段,集群资源利用率提升45%,计算成本降低 30%
  4. 效率提升:新指标上线周期从原来的 3 天缩短至4 小时,新业务线接入时间从 2 周缩短至 3 天
  5. 业务价值:实时风控系统拦截异常订单率提升 20%,骑手平均配送时长缩短 8%,商家订单转化率提升 5%

美团外卖实时数仓建设与优化项目(核对版)

项目周期:2025.03-2025.10

技术栈:Flink 1.17(批流一体计算)、Kafka 3.4(消息队列)、ClickHouse 23.3(实时 OLAP 存储)、Canal 1.1.7(CDC 采集)、Redis 7.0(维度缓存)、Spark 3.3(历史数据回溯)、DolphinScheduler(任务调度)、Superset(可视化)

项目背景:美团外卖日均订单量超 6000 万,午晚高峰峰值 QPS 达 120 万。原有系统存在三大致命问题:

  1. 运营复盘完全依赖 T+1 离线数仓,无法支撑分钟级业务决策
  2. 实时需求通过零散临时脚本实现,口径不一致、稳定性差,核心指标延迟高达 30 分钟
  3. 大促峰值时系统频繁崩溃,无法支撑亿级订单流量

核心职责

  1. 主导整体架构设计:采用基于 Flink 批流一体的纯 Kappa 架构替代传统 Lambda 架构,统一实时与离线计算引擎、数据模型和指标口径
  2. 设计五层实时数仓体系:完成 ODS/DWD/DWS/DIM/ADS 分层建模,覆盖订单、用户、商家、骑手四大核心业务域,实现数据复用最大化
  3. 开发核心计算链路:主导实现实时订单统计、骑手运力监控、用户行为画像等 15 + 个核心 Flink 作业,支撑实时风控系统的毫秒级数据输入链路
  4. 攻克关键技术难题:解决头部商家 / 热门城市数据倾斜、亿级用户标签大维度关联、Flink 状态膨胀等核心问题,保障大促稳定性
  5. 建立全链路数据治理体系:搭建覆盖完整性、准确性、一致性、及时性的监控平台,核心数据质量 SLA 达到 99.95%,实现数据异常自动告警和分钟级故障恢复
  6. 构建统一数据服务层:提供 SQL 查询、REST API、消息订阅三种服务方式,为全国运营大屏、商家后台、骑手 APP 等 20 + 个业务系统提供数据支持

核心成果

  1. 性能大幅提升:核心业务指标延迟从原有临时方案的 30 分钟降至5 秒以内,支持分钟级业务决策
  2. 稳定性行业领先:核心链路系统可用性达到99.99%,成功支撑各种大促活动,峰值订单量突破 1.2 亿单 / 天,大促期间零数据丢失
  3. 资源显著优化:通过分层预聚合、大字段拆分 Join、热点隔离等手段,集群 CPU 利用率从 25% 提升至 70%(提升 45 个百分点),计算成本降低 30%
  4. 开发效率质变:通过公共层复用和指标自动化生成工具,新指标上线周期从 3 天缩短至 4 小时,新业务线接入时间从 2 周缩短至 3 天
  5. 业务价值突出
    • 为实时风控系统提供毫秒级数据支持,助力异常订单拦截率提升 20%,年减少损失超 5000 万元
    • 实时运力监控数据支撑调度系统算法优化,全国平均配送时长缩短 8%
    • 商家实时经营数据赋能精细化运营,平台整体订单转化率提升 5%

项目总结

细化项目列表

  1. 美团实时数仓

项目描述

设计细节

难点

美团实时数仓

项目描述

为了应对xxx场景,包括了xxx模块

设计细节

纯 Kappa 在处理复杂业务回撤、历史重算吞吐及多状态业务关联时存在局限,所以架构采用的是**基于场景的混合架构(Lambda + Kappa 结合)**‌,核心原则是“日志类走纯流(Kappa 模式),业务类走流批结合或实时 OLAP(类 Lambda 模式)”。‌‌

难点

美团外卖实时数仓

美团实时数仓项目

1. 项目背景与核心痛点(定调子,讲清 “为什么做”)

1.1 项目描述及背景

主导了公司交易域实时数仓从 Lambda 架构到 Flink+Hudi 流批一体架构的升级改造,覆盖外卖、到店两大核心业务线。

早期实时项目是基于需求的,没有系统化的思想,目的只是为了尽早上线,所以采取的是Lambda架构,实时与离线生成两套。但一方面是人力资源消耗大,维护成本高,还容易出现数据不一致的问题。对于实时部分难以回溯

1.2 量化业务及数据规模

日均数据量、峰值 QPS、核心表数量、下游业务场景数、延迟要求。

1.3 核心痛点

  • Lambda 双链路两套逻辑,核心指标口径差异率最高达 5%,业务方不敢信实时数据;
  • 开发运维成本翻倍,新增一个指标需要同时写实时 + 离线两套代码,上线周期 3 天以上;
  • 历史数据回溯困难,纯实时链路无法支持 7 天以上的指标重算和逻辑修正。

1.4 项目目标

实现流批一套逻辑、一份存储,口径差异率降到 0.1% 以内,开发效率提升 50%

2. 整体架构设计及技术选型(展架构能力,讲清 “为什么这么选”)

随着数据湖技术成熟,美团提出增量数仓架构,基于 Flink + Apache Hudi 实现流批一体,向 “一套计算逻辑、一份数据存储” 演进,这也是当前的主流架构。

  • 用数据湖(Hudi)承载全量历史数据,替代原本离线数仓的 Hive 表,同时支持实时增量写入和批量读取;
  • 用 Flink 统一流计算和增量批计算,同一套 SQL 逻辑可同时运行在流、批两种模式下;
  • 保留了 Kappa 架构 “以增量为核心” 的思想,但用数据湖解决了纯 Kappa 中消息队列存储成本高、历史回溯困难的缺陷。
维度 经典 Kappa 架构 美团增量数仓架构
全量存储 依赖 Kafka 等消息队列留存全量历史数据 以 Hudi 数据湖作为统一全量存储,Kafka 仅做实时数据通道
批处理实现 重放消息队列历史数据,用流引擎模拟批处理 基于数据湖的增量读取能力,直接做增量 / 全量批计算
适用场景 适合中小数据量、短周期回溯场景 支持千亿级数据、长周期历史分析,适配多业务场景

2.1 核心架构

  • 数据接入层: 通过Canal采集数据库Binlog、Scribe采集业务日志、IoT数据,统一写入Kafka,同时同步至Hudi数据湖
  • 计算引擎层: 以Flink为唯一核心计算引擎,基于FlinkSQL实现流批统一的数仓开发,保障Exactly-Once计算语义

2.2 分层架构

ODS层 DWD层 DIM层 DWS层 APP层
核心定位 实时数仓的原始数据接入层,保留业务系统最原生的数据形态,仅做格式统一与基础校验,不进行深度业务加工。 实时数仓的标准化明细核心层,完成数据清洗与业务事实建模,是数据资产化的核心载体。 统一公共实时维度层,为全链路实时计算提供维度数据,解决实时场景下维度变更的时效性与一致性问题。 公共指标汇总层,按主题域构建轻度汇总宽表,沉淀共性业务指标,统一指标口径,提升下游复用率。 面向具体业务场景的应用输出层,针对特定业务需求定制化加工,直接对接数据消费端。
处理的数据 1. 业务数据库变更
2. 用户行为日志
3. 业务消息与系统日志
1. 数据标准化
2. 事实建模
3. 维度退化
1. 基础公共维度
2. 缓慢变化维
1. 按时间、地域、业务线、品类等常用维度,聚合原子指标
2. 构建主题宽表,例如交易汇总宽表、流量汇总宽表、商家经营汇总宽表等
1. 实时监控大屏
2. 在线业务数据
3. 实时分析报表
数据处理方式 1. 数据采集:Canal 采集数据库 Binlog,Scribe/Flume 采集日志数据,统一投递到消息总线
2. 实时通道:以 Apache Kafka 作为统一实时数据总线,承接所有高吞吐实时数据流入
3. 轻量处理:Flink 完成数据格式解析、基础校验、分区路由等轻量 ETL 操作
1. 计算引擎:Flink SQL 实现清洗、转换、关联逻辑,原生保障 Exactly-Once 计算语义
2. 维度关联:通过 Flink 流维 Join 能力,实时关联 DIM 层的维度数据
通过 Flink CDC 实时采集维度表的 Binlog 变更,毫秒级同步更新维度存储 Flink SQL 实现窗口聚合、增量聚合,支持滚动窗口、滑动窗口、会话窗口等多种实时聚合模式 1. 定制计算:毫秒级延迟的定制化指标,由 Flink 完成最终聚合后直接写入 KV 存储,对外提供服务
2. 服务封装:通过统一数据服务平台,对外提供 RPC/HTTP 接口,屏蔽底层存储与计算细节
处理引擎 1. 数据采集:Canal
2. 实时通道:Apache Kafka
Flink SQL Flink CDC Flink SQL 统一以 Apache Doris 为核心 OLAP 引擎,支撑大部分实时分析、大屏查询场景;高并发点查场景使用 Redis/MySQL
数据流转方式
存储选型 早期仅在 Kafka 做短期留存;当前增量数仓阶段,数据同步写入基于 Apache Hudi 自研的 Beluga 数据湖,作为全量历史数据的永久存储 早期明细数据直接落地 Kafka 供下游消费;当前增量数仓阶段,核心明细数据统一写入 Hudi 数据湖,支持行级 Upsert 更新、增量读取,同时兼容流、批两种消费模式 1. 高性能 KV 存储:Tair/Redis/HBase,用于 Flink 实时计算时的低延迟维度关联查询
2. 全量快照存储:Hudi 数据湖存储全量维度历史快照,支持批处理场景的维度回溯与校对
早期汇总数据同时写入 Kafka 和 OLAP 引擎;当前增量数仓阶段,汇总层核心数据落地 Hudi 数据湖,同时同步到 Apache Doris 等 OLAP 引擎加速查询 统一以 Apache Doris 为核心 OLAP 引擎
重难点 1. 一致性保障:基于数据湖的事务能力,保障增量聚合的数据一致性,从根源避免 Lambda 架构双链路口径不一致的痛点
2.
举例 ods_trade_order_binlogods
分层设计的核心原则
  1. 口径统一:所有公共指标统一沉淀在 DWS 层,APP 层禁止重复计算核心指标,避免口径不一致。
  2. 复用优先:DWD 明细层和 DWS 汇总层面向全业务线复用,APP 层仅做场景化裁剪,不重复加工底层数据。
  3. 流批一体:所有分层表均基于 Hudi 存储,同一套表同时支持流处理增量写入和批处理全量 / 增量读取,一套逻辑服务实时 + 离线两类场景。

2.3 技术选型原因

2.3.1 为什么选Hudi而不是Iceberg/Delta Lake?

结合更新频率、SQL 生态、Flink 适配度、批量回溯需求

结论

美团增量实时数仓的核心诉求是「高频 CDC 行级更新 + Flink 流批一体 + 分层增量流转」,Hudi 的设计基因与这个场景高度契合;而 Iceberg 更偏向纯分析型湖仓一体,Delta Lake 强绑定 Spark 生态,二者在实时更新、Flink 适配、增量处理三个核心维度上,都无法匹配美团的架构定位。

面试话术

选型 Hudi 主要是和我们的场景高度匹配,核心有三点: 第一,我们 70% 以上是交易 CDC 数据,高频行级更新,Hudi 的 MOR 表和原生 Upsert 能力是三家里面最强的,写入延迟和写放大都最优,Iceberg 和 Delta Lake 在高频更新场景下性能差距很明显; 第二,我们全链路用 Flink 做流批一体,Hudi 对 Flink 的适配是最成熟的,从写入到增量消费全链路支持,另外两家要么 Spark 绑定深,要么 Flink 生态不完善; 第三,我们做的是增量数仓,分层之间靠增量数据流转,Hudi 原生的增量查询能力刚好匹配这个架构,下游不用全量扫表,加工效率提升很明显。

当然 Iceberg 在多引擎查询、Schema 演进上有优势,Delta Lake 在 Spark 栈里更简单,但都不是我们的核心诉求。综合下来,Hudi 是和我们实时数仓场景最契合的选择。

结合实时延迟要求、状态管理、CDC 场景适配

结论

本质不是 Spark 技术能力不行,而是二者的架构基因完全不同,流批统一的实现路径截然相反。美团增量实时数仓的核心定位是「以实时链路为核心,批处理做补数 / 回溯 / 校对的补充」,Flink 原生流驱动的架构,与这个场景的匹配度远高于 Spark;而 Spark 的流批一体是「批处理为内核,用微批模拟流」,更适合离线为主、实时为辅的湖仓分析场景。

面试话术

选 Flink 做流批统一,核心是和我们的增量实时数仓场景高度匹配,主要有三点: 第一是架构基因匹配,Flink 是原生流引擎,批是流的特例,我们的架构以实时增量加工为主、批处理做补数校对,正好契合;Spark 本质是批引擎,流是微批模拟,更适合离线为主的场景。 第二是核心能力匹配,我们大量是 CDC 高频更新数据,还有很多大状态窗口聚合、流维 Join 的场景,Flink 的状态管理、事件时间、Upsert 流语义都更成熟稳定,Spark 在大状态、高频更新场景下的短板比较明显。 第三是落地收益更实在,Flink 同一套 SQL 流批模式可以直接跑,真正做到一套逻辑、口径天然一致,解决了原来 Lambda 双链路的核心痛点;Spark 流批表面统一,实际语义差异大,很多时候还是要写两套逻辑。

当然 Spark 在离线批处理、机器学习生态上优势更强,我们也不是完全不用,长周期历史回溯、离线训练这些还是用 Spark,核心实时链路用 Flink,各司其职。

2.3.3 OLAP 引擎为什么选 Doris/ClickHouse?

结合查询模式、并发量、更新需求

2.4 保障性部分设计

2.4.1 如何保证数据的一致性和完整性

在美团基于「Flink + Hudi」的增量实时数仓架构中,数据一致性(口径统一、事务一致、状态对齐)和数据完整性(不丢失、不重复、全量覆盖)是从引擎底层、存储层、建模层到质量运维全链路系统性保障的,核心机制可分为六大层面:


2.4.1.1 计算引擎层:端到端 Exactly-Once 语义,筑牢不丢不重的基础

这一层从计算根源上同时保障数据完整性(不丢不重)计算状态一致性,核心依赖 Flink 的原生机制与上下游存储的协同。

  1. Checkpoint 状态一致性保障 Flink 基于异步快照的 Checkpoint 机制,定期将计算状态(窗口聚合、关联状态、消费偏移量等)持久化存储。任务故障重启时,从最近一次成功的 Checkpoint 精准恢复,保证计算状态、消费位点与计算逻辑完全一致,既不会重复计算已处理数据,也不会跳过未处理数据。
  2. 两阶段提交(2PC)实现端到端精确一次 Flink 通过 2PC 事务机制,对接上游 Kafka 和下游 Hudi 存储,实现端到端 Exactly-Once 语义:
  • 预提交阶段:每个 Checkpoint 周期内,数据写入存储的临时状态,不对外可见;
  • 正式提交阶段:Checkpoint 完成后,所有分片数据一次性原子提交,对外可见。 这保证了每一批数据要么全部写入成功,要么全部回滚,不会出现部分写入的中间状态,从写入层面保障了数据一致性和完整性。
  1. 主键幂等写入 所有分层表均基于唯一主键(如订单 ID、配送单 ID)做 Upsert 更新,即使因故障重试导致数据重复投递,也只会覆盖更新同一条记录,不会产生重复数据,从存储侧兜底了数据完整性。

2.4.1.2 存储层:数据湖事务能力 + 统一存储,从架构上消除一致性痛点

美团从传统 Lambda 双链路架构演进到 Hudi 增量数仓,最核心的收益就是从存储根源上解决了流批口径不一致的经典问题。

  1. Hudi 事务原子性与快照隔离 Hudi 基于时间线(Timeline)的事务机制,所有写入、更新、合并操作都具备原子性:
  • 写入过程对查询不可见,只有提交完成后才会生成新的快照,查询始终读取一致的完整快照,不会读到半写的中间数据;
  • 支持快照读、增量读两种模式,无论实时消费还是批量回溯,读取的数据都是完整且一致的。
  1. 一份存储统一流批,消除双链路差异 经典 Lambda 架构下,实时链路和离线链路各有一套计算逻辑、一套存储,口径不一致是天生缺陷。美团增量数仓中:
  • DWD、DWS 层核心数据全部落地 Hudi 数据湖,作为唯一的全量数据存储;
  • 实时计算通过增量写入更新数据,离线批量计算直接读取同一份 Hudi 表做增量 / 全量计算,同一套表、同一套逻辑、同一份数据,从架构层面彻底消除了双链路口径不一致的问题。
  1. 多版本快照与回滚能力 Hudi 保留所有历史版本的快照,支持按时间点回溯到任意一个一致的数据版本。当出现计算逻辑错误、脏数据写入时,可快速回滚到历史正确快照,保障数据可恢复、不丢失,为数据完整性提供兜底。
2.4.1.3 建模与口径层:统一维度与指标,保障业务语义一致性

技术层面的一致性是基础,业务口径的一致性才是最终目标,这一层主要通过标准化建模保障业务语义一致性

  1. 统一公共维度体系
  • 全链路复用同一套 DIM 层维度表,通过 Flink CDC 秒级同步维度变更,所有实时计算任务共享同一份维度数据,避免不同任务各自维护维度导致的属性不一致。
  • 针对缓慢变化维,采用时态关联机制:基于事件时间匹配对应时间点的维度版本,保证 “历史事实匹配历史维度”,不会因维度后续变更导致历史指标波动,保障时间维度上的一致性。
  • 统一维度编码、枚举值、计量单位,从源头消除跨业务线的口径差异。
  1. 指标口径统一沉淀
  • 所有核心公共指标(GMV、订单量、用户数等)统一在 DWS 层定义和计算,APP 层仅做场景化裁剪和派生,禁止重复开发核心指标逻辑,确保一个指标只有一个口径。
  • 指标计算逻辑流批共用同一套 Flink SQL,流模式做实时增量更新,批模式做全量 / 增量校对,从开发层面保证实时与离线指标的口径完全一致。
2.4.1.4 完整性兜底:乱序、迟到、补数全场景覆盖

实时数据天然存在乱序、迟到、回溯补数等场景,美团通过多层机制保障极端场景下的数据完整性

  1. 乱序与迟到数据处理
  • 基于 Flink Watermark 机制,设置合理的乱序容忍窗口,保证绝大多数乱序数据都能被正确纳入窗口计算,避免数据遗漏。
  • 超迟到数据通过侧输出流单独捕获,结合 Hudi 的行级 Upsert 能力,回写更新历史聚合结果,即使延迟数小时、数天的数据也能最终修正到正确结果,全程不丢失任何一条数据。
  • 针对 Binlog 变更数据,基于数据库 GTID / 事务提交时间排序,保证增删改操作按正确顺序执行,避免乱序导致的最终状态错误。
  1. 全量与增量无缝衔接 CDC 数据同步阶段,全量快照同步与增量 Binlog 消费基于精准位点衔接,保证全量初始化过程中新增、变更的数据既不重复、也不丢失,实现全量 + 增量的完整闭环。
  2. 数据补数与回溯机制 当出现逻辑变更、历史数据修复需求时,支持按时间分区、按业务范围进行增量 / 全量补数:
  • 补数过程在 Hudi 中原子执行,补数完成后一次性切换对外可见,不影响线上实时业务的正常查询;
  • 补数前后自动执行条数、核心指标对账,确保补数后数据完整、口径一致。
2.4.1.5 全链路质量校验:闭环监控一致性与完整性

美团配套了完善的实时数据质量体系,对全链路数据进行持续校验,及时发现和定位一致性、完整性问题。

  1. 分层级对账校验
  • 层间流量对账:ODS → DWD → DWS → APP 每一层级,实时统计流入流出的数据条数,计算数据通过率,异常波动分钟级告警,快速定位数据丢失或重复。
  • 流批指标对账:每日离线数仓计算完成后,自动比对实时数仓与离线数仓的核心指标(GMV、订单量等),量化差异率,超过阈值自动告警,从业务结果层面校验口径一致性。
  • 主键唯一性校验:实时检测各层表的重复主键、空主键问题,保障数据实体的完整性。
  1. 业务规则与一致性校验
  • 内置标准化数据质量规则,包括非空校验、枚举值校验、数值范围校验(如订单金额非负)、维度关联存在性校验等,异常数据自动分流至脏数据表,不污染主链路数据。
  • 跨域一致性校验:例如交易订单支付金额与财务结算金额、订单量与配送单量的一致性比对,保障跨业务域的数据语义一致。
  1. 脏数据可追溯可修复 所有被过滤的脏数据、异常数据都会完整留存,附带异常原因和原始数据,支持后续人工复核和修复,保证原始数据不丢失。
2.4.1.6 故障容灾与兜底:极端场景的一致性保障
  1. 故障精准恢复 Flink 任务支持从 Checkpoint/Savepoint 精准恢复,Kafka 消费位点、计算状态、中间结果完全对齐,恢复后数据延续之前的状态继续处理,不会出现断点、重复或丢失。
  2. 存储高可用 Kafka 多副本部署、Hudi 数据多副本存储,避免单点硬件故障导致的数据永久丢失。
  3. 逻辑变更灰度与回滚 数仓逻辑变更支持灰度发布,异常情况下可快速回滚计算逻辑,并通过 Hudi 快照回滚数据到历史一致版本,保障业务侧数据的连续性和一致性。
总结

美团实时数仓的保障思路是:先通过架构升级(流批一体 + 数据湖)从根源减少一致性问题,再通过引擎语义、存储事务保障底层基础,再通过标准化建模统一业务口径,最后用全链路质量校验和容灾兜底形成闭环。相比传统 Lambda 架构靠人工对齐双链路的方式,这种架构下的一致性和完整性保障成本更低、可靠性更高。

2.5 当前增量数仓分层的核心特性

  1. 存储统一:核心分层数据从「Kafka 分层流转」升级为「Hudi 数据湖统一存储」,解决了纯实时架构历史回溯难、长期存储成本高的缺陷
  2. 计算统一:全链路以 Flink 为唯一核心计算引擎,同一套 SQL 逻辑可同时运行在流、增量批模式下,真正实现流批一体
  3. 口径统一:一套分层模型同时服务实时、离线两类场景,彻底解决了早期 Lambda 架构双链路开发、口径不一致的痛点

最后,一句话总结架构优势:比如 “最终实现了‘一套 SQL 逻辑、一份 Hudi 存储’,同时支持实时增量写入和离线批量回溯”。

3. 核心技术难点与解决方案(亮深度,讲清 “你解决了什么难题”)

选 3-4 个最有含金量的难点,每个按「问题场景→核心挑战→方案对比→落地效果」来讲

注意:每个难点都要量化挑战和结果,比如 “单任务状态超过 2TB,Checkpoint 时长超 15 分钟频繁超时,最终优化到 3 分钟内完成,任务稳定性从 95% 提升到 99.9%”。

3.1 流批数据一致性保障

3.1.1 怎么解决双链路口径不一致
3.1.2 怎么实现端到端 Exactly-Once
3.1.3 怎么验证一致性

3.2 海量数据下的性能与稳定性

3.2.1 大状态优化
3.2.2 Checkpoint 超时
3.2.3 数据倾斜
3.2.4 Hudi 写入性能与小文件治理

3.3 实时维度建模与维表一致性

3.3.1 高实时维表的更新与关联
3.3.2 缓慢变化维处理
3.3.3 维表数据乱序 / 延迟问题

3.4 数据湖架构下的增量回溯与补数

3.4.1 历史数据重算机制
3.4.2 补数不影响线上服务
3.4.3 原子切换方案

4. 落地成果及价值量化(证结果,讲清 “做成了什么”)

分技术收益和业务收益两部分,全部用数字说话

  • 技术收益(硬指标):
    • 数据一致性:核心指标流批差异率从 X% 降到 Y%;
    • 研发效率:需求上线周期从 X 天缩短到 Y 小时,人力成本降低 X%;
    • 性能成本:存储成本降低 X%,查询延迟从 X 秒降到 Y 毫秒;
    • 稳定性:任务全年可用性 X%,故障恢复时间从 X 小时缩短到 X 分钟。
  • 业务收益(强价值):
    • 支撑了哪些核心业务场景:比如实时经营大屏、实时风控、商家实时看板、大促活动监控;
    • 带来的业务结果:比如风控拦截率提升 X%,运营活动决策效率提升 X 倍,支撑了多少次亿级大促平稳跑通。

5. 个人核心贡献与技术沉淀(显段位,讲清 “你的不可替代性”)

  • 角色定位:比如 “项目技术负责人,主导整体架构设计与技术选型”;
  • 核心动作:比如 “牵头攻克了 XX、XX 等核心技术难点,制定了全公司实时数仓建模规范与开发标准”;
  • 技术沉淀:比如 “沉淀了 XX 通用组件 / 工具,在 3 个业务线复用;输出了 XX 技术方案与最佳实践,纳入团队技术体系”;
  • 团队影响:比如 “带领 X 人小组完成落地,培养了 X 名骨干开发”。

6. 高频追问点分类准备清单

按面试考察维度分类,每个问题都建议准备 “结论 + 细节数据 + 原理依据”,避免泛泛而谈。

6.1 背景与角色类(核实项目真实性与参与深度)

  1. 这个项目总共有多少人做?做了多久?你具体负责哪几块?
  2. 项目上线前,日均数据量、峰值 QPS、集群规模是多少?上线后呢?
  3. 这个项目是业务驱动还是技术驱动?最初是谁发起的?
  4. 项目过程中,你遇到的最大的非技术阻力是什么?怎么推动的?

6.2 架构选型类(考察架构思维与决策能力)

  1. 为什么要从 Lambda 架构改成流批一体?是遇到了哪些无法解决的问题?
  2. 数据湖选型为什么最终定了 Hudi?和 Iceberg 比,它的优劣势分别是什么?你们的场景里哪些点最关键?
  3. MOR 表和 COW 表你们怎么选?哪些层用 MOR,哪些层用 COW?为什么?
  4. 为什么不用纯 Kappa 架构?你们评估过纯 Kappa 的方案吗?为什么没选?
  5. Flink 和 Spark 都能做流处理,为什么你们实时数仓核心链路选 Flink?
  6. DWS 层为什么不直接写到 OLAP 引擎里,还要先落 Hudi?
  7. 你们的实时数仓为什么要分这几层?能不能合并?比如 DWD 和 DWS 合并?

6.3 核心技术深挖类(考察技术深度,对应你讲的每个难点)

6.3.1 一致性相关
  1. 你们说实现了端到端 Exactly-Once,具体是怎么实现的?Flink 的 2PC 和 Hudi 的事务是怎么配合的?
  2. 怎么证明流批口径是一致的?你们用什么方式对账?差异率控制在多少?出现差异一般是什么原因?
  3. Hudi 的事务原理是什么?怎么保证写入的原子性和查询的快照隔离?
  4. 主键重复、乱序更新的场景,你们怎么保证最终数据是对的?
6.3.2 性能与稳定性相关
  1. 你们 Flink 任务最大的状态有多大?用的什么状态后端?Checkpoint 频率和时长是多少?
  2. 遇到过 Checkpoint 超时或者失败吗?一般是什么原因?怎么优化的?
  3. 实时聚合任务出现数据倾斜怎么处理?你们有通用的解决方案吗?
  4. Hudi 写入的时候小文件多吗?怎么治理的?Compaction 策略是怎么设置的?
  5. 高峰时段 Kafka 有没有积压?怎么排查和优化?
  6. 大状态任务宕机恢复,一般需要多久?怎么缩短恢复时间?
6.3.3 实时数仓建模相关
  1. 实时维表你们用什么存储?怎么保证维表数据和源库的一致性?更新延迟是多少?
  2. 流维 Join 的时候,维表数据更新了,历史数据的关联结果会修正吗?怎么处理?
  3. 迟到数据你们怎么处理?Watermark 设置多久?超迟到的数据怎么回写修正结果?
  4. DWS 层的聚合粒度是怎么定的?太细太粗分别有什么问题?你们怎么权衡?
  5. 实时数仓里,更新型的事实表(比如订单状态变更)你们怎么建模?

6.4 问题与踩坑类(考察实战经验与排障能力)

  1. 做这个项目过程中,你踩过最严重的一个线上故障是什么?怎么定位、怎么解决的?
  2. 有没有出现过数据丢失或者重复的情况?什么原因?怎么兜底的?
  3. Hudi 上线后,有没有遇到过意料之外的问题?比如查询性能下降、元数据膨胀?
  4. 业务方反馈实时数据不准的时候,你们一般的排查流程是什么?
  5. 有没有遇到过需求和现有架构冲突的情况?比如要求秒级延迟又要复杂大口径关联,你怎么处理?

6.5 业务与价值类(考察业务思维与结果导向)

  1. 这个项目上线后,业务方最直观的感受是什么?有没有具体的案例?
  2. 核心指标的口径是谁来定的?业务方和技术方出现口径分歧的时候,你怎么协调?
  3. 怎么评估一个实时需求要不要做?怎么判断是用实时数仓还是离线数仓?
  4. 项目投入了多少人力成本?带来的收益能覆盖成本吗?怎么衡量 ROI?

6.6 复盘与规划类(考察技术视野与成长潜力)

  1. 现在回头看,这个架构还有什么不满意的地方?还有哪些可以优化的点?
  2. 如果让你从零再做一遍这个项目,你会在哪些地方调整设计?
  3. 你觉得实时数仓下一步的演进方向是什么?比如湖仓一体、AI + 实时数仓?
  4. 你们现在的架构,支撑未来 3 倍的数据量增长有问题吗?瓶颈会在哪里?怎么提前扩容?

6.7 团队与方法论类(考察技术领导力)

  1. 你们怎么保证多个开发人员写的实时任务,口径、规范是统一的?
  2. 线上实时任务的运维规范、故障处理流程是你制定的吗?具体是什么样的?
  3. 新人上手实时数仓开发,你会怎么带?怎么快速规避常见坑?
  4. 你们怎么做需求评审和技术方案评审?怎么避免上线后才发现问题?

附:背景资料

1.1 美团外卖核心业务链路

1
2
3
4
用户端:APP浏览→搜索→加购→下单→支付→评价→退款
商家端:接单→出餐→呼叫骑手→处理退款
骑手端:接单→到店取餐→配送→送达
平台端:流量分发→营销投放→风控拦截→运力调度→客服处理

1.2 实时数仓建设目标

  • 低延迟:核心指标延迟 < 10 秒,支持分钟级业务决策
  • 高可靠:数据不丢不重,Exactly-Once 语义保证
  • 高可用:支持千万级 QPS,峰值(午晚高峰)稳定运行
  • 统一口径:实时与离线指标口径 100% 一致
  • 易扩展:支持新业务线快速接入,新指标小时级上线

1.3 核心业务支撑场景

场景 延迟要求 核心指标
全国实时运营大屏 <5 秒 实时订单量、交易额、在线骑手数、配送时长
商家实时经营后台 <10 秒 今日订单量、收入、出餐时长、差评数
骑手实时调度系统 <1 秒 区域骑手运力、待配送订单数、平均配送时长
实时风控系统 <100 毫秒 异常订单检测、恶意用户识别、刷单作弊拦截
实时营销系统 <1 秒 优惠券核销率、活动参与人数、转化效果
实时用户推荐 <500 毫秒 用户实时行为标签、商品点击率、转化率

1.4 各分层典型表设计

一、ODS 层(操作数据层)

ODS 层保留原始数据的原生形态,仅做格式统一和基础校验,不进行业务逻辑加工,分为数据库 Binlog 变更类用户行为日志类两大类型。

表名 业务说明 核心字段示例 存储与更新机制
ods_trade_order_binlog 交易订单库全量 Binlog 变更,覆盖外卖、到店等所有订单的增删改操作 order_id, user_id, shop_id, order_amount, order_status, pay_time, create_time, update_time, operation_type(INSERT/UPDATE/DELETE) 存储:Kafka(实时通道,短期留存) + Hudi(全量永久存储)更新:Canal 采集 Binlog 实时写入,按主键 Upsert
ods_delivery_dispatch_binlog 配送调度库的骑手配送状态变更数据,覆盖接单、取餐、送达全链路 dispatch_id, order_id, rider_id, delivery_status, receive_time, pickup_time, arrive_time, operation_type 存储:同上更新:按配送单主键实时 Upsert
ods_user_page_view_log App / 小程序端用户页面浏览埋点日志,全量用户流量原始数据 log_id, user_id, device_id, page_id, page_name, stay_time, event_time, ip, source_channel 存储:Kafka + Hudi更新:日志纯追加写入,无更新
ods_user_click_event_log 用户点击事件埋点,覆盖商品、商家、按钮等所有交互点击 event_id, user_id, target_id, target_type, event_time, page_source 存储:同上更新:纯追加写入

二、DWD 层(数据明细层)

DWD 层是实时数仓的明细核心,按单个业务过程拆分事务型事实表,完成数据清洗、标准化、脏数据过滤和维度退化,是后续汇总计算的统一数据底座。

1. 交易域

表格

表名 业务说明 核心字段示例 存储与更新机制
dwd_trade_order_pay_detail 订单支付明细事实表,对应「订单支付成功」业务过程 order_id, user_id, shop_id, city_id, category_id, pay_amount, pay_time, pay_channel, order_type 存储:Hudi MOR 表(支持高频 Upsert)更新:Flink 过滤支付成功事件,按订单主键实时 Upsert
dwd_trade_order_refund_detail 订单退款明细事实表,对应「订单退款」业务过程 refund_id, order_id, shop_id, refund_amount, refund_reason, refund_time, refund_status 存储:Hudi MOR 表更新:按退款单主键 Upsert
2. 流量域

表格

表名 业务说明 核心字段示例 存储与更新机制
dwd_traffic_page_view_detail 页面浏览明细事实表,清洗去重后的标准化用户访问数据 user_id, device_id, page_id, page_module, stay_duration, event_time, city_id, source_channel 存储:Hudi COW 表(追加写入为主)更新:Flink 清洗去重后追加写入
dwd_traffic_goods_click_detail 商品点击明细事实表,对应「用户点击商品」业务过程 user_id, goods_id, shop_id, click_time, page_source, position_id 存储:Hudi COW 表更新:追加写入
3. 履约域

表格

表名 业务说明 核心字段示例 存储与更新机制
dwd_delivery_status_change_detail 配送状态变更明细事实表,覆盖配送全链路每个状态节点 dispatch_id, order_id, rider_id, before_status, after_status, change_time, station_id 存储:Hudi MOR 表更新:按状态变更记录追加,按配送单主键关联更新

设计说明:DWD 层会提前退化城市、品类、业务线等高频使用的维度属性,减少下游重复关联维度表的计算开销。


三、DIM 层(维度数据层)

DIM 层提供全链路统一的公共维度,支持秒级实时更新,实时计算场景通过 KV 存储低延迟关联,批处理 / 回溯场景通过数据湖全量关联。

表格

表名 业务说明 核心字段示例 存储与更新机制
dim_shop_info 商家全量维度表,包含基础属性、经营属性 shop_id, shop_name, city_id, category_id, business_type, business_status, open_time, address, update_time 实时关联:Tair/HBase(主键毫秒级查询)全量存储:Hudi(保留历史版本)更新:Flink CDC 采集 Binlog,秒级同步更新
dim_user_info 用户基础维度表,包含用户属性与分层标签 user_id, gender, age_level, city_id, register_time, user_level, is_new_user 存储:同上更新:用户属性变更实时同步
dim_goods_info 商品维度表,覆盖外卖、到店商品基础属性 goods_id, shop_id, category_id, goods_name, price, is_on_sale, update_time 存储:同上更新:商品上下架、价格变更实时同步
dim_area_info 行政区划维度表,省 - 市 - 区三级映射 area_id, area_name, parent_id, area_level, city_code 存储:Hudi + 本地缓存更新:低频变更,天级批量更新

四、DWS 层(数据汇总层)

DWS 层是公共指标汇总层,按主题域和统计粒度构建轻度汇总宽表,沉淀 80% 以上的通用业务指标,统一口径,供下游应用复用。分为分钟 / 小时级窗口汇总日粒度累计汇总两类。

1. 交易域

表格

表名 统计粒度 核心指标 存储与更新机制
dws_trade_shop_1d shop_id + 自然日 订单数、支付订单数、取消订单数、GMV、实付金额、退款金额、客单价、下单用户数 存储:Hudi(主键 Upsert) + 同步至 Doris 加速查询更新:Flink 增量聚合,日内实时更新,凌晨最终落定
dws_trade_city_1h city_id + 小时 小时 GMV、小时订单量、支付用户数、订单转化率 存储:Hudi + Doris更新:小时窗口聚合,实时更新
2. 流量域

表格

表名 统计粒度 核心指标 存储与更新机制
dws_traffic_page_1h page_id + 小时 PV、UV、平均停留时长、跳出率、点击转化率 存储:Hudi + Doris更新:小时窗口聚合
dws_traffic_shop_1d shop_id + 自然日 曝光量、进店量、进店转化率、商品点击量 存储:同上更新:日内实时累计更新
3. 履约域

表格

表名 统计粒度 核心指标 存储与更新机制
dws_delivery_rider_1d rider_id + 自然日 接单量、完成单量、超时单量、平均配送时长、配送里程 存储:Hudi + Doris更新:日内实时累计

五、APP 层(应用数据层)

APP 层面向具体业务场景定制化开发,指标口径完全匹配业务需求,直接对接下游数据消费端。

表格

表名 业务场景 统计粒度 核心指标 存储与消费方
app_trade_city_realtime 城市级经营实时大屏、活动监控 city_id + 5 分钟窗口 实时 GMV、实时订单量、支付用户数、同比 / 环比增速 存储:Doris消费方:经营作战室、大促监控大屏
app_shop_operation_realtime 商家端实时经营看板 shop_id + 分钟级 今日订单数、今日收入、待处理订单、实时曝光量、进店转化率 存储:Doris + Redis(热点商家缓存)消费方:商家后台、商家 App
app_risk_user_feature_realtime 实时风控反作弊特征 user_id(近 1 小时 / 24 小时窗口) 近 1 小时下单次数、近 1 小时支付金额、近 24 小时退款次数、常用设备数 存储:Redis(毫秒级点查)消费方:实时风控引擎、反作弊系统
app_activity_effect_realtime 运营活动实时效果监控 activity_id + 分钟级 活动参与人数、核销订单数、补贴金额、活动 ROI 存储:Doris消费方:运营活动监控平台

美团外卖实时数仓项目完整设计方案

一、项目业务背景与目标

二、整体架构设计

2.1 架构选型:**基于场景的混合架构(Lambda + Kappa 结合)**‌

纯 Kappa 在处理复杂业务回撤、历史重算吞吐及多状态业务关联时存在局限,所以架构采用的是**基于场景的混合架构(Lambda + Kappa 结合)**‌,核心原则是“日志类走纯流(Kappa 模式),业务类走流批结合或实时 OLAP(类 Lambda 模式)”。‌‌

核心架构特征

  • 非纯 Kappa‌:明确承认纯 Kappa 在处理复杂业务回撤、历史重算吞吐及多状态业务关联时存在局限,未全量推行 。
  • ‌分场景混合策略:
    1. 日志类场景‌(如点击流、监控):数据不可变、逻辑简单,采用‌Kappa 架构‌(统一流处理,Flink/Storm),一套代码生产实时指标 。
    2. 业务类场景‌(如订单状态变更、金额汇总):涉及多表关联、状态回撤(如取消订单),采用‌Lambda 架构变体或实时 OLAP‌(Flink 预聚合 + Doris 存储计算),利用批处理或 OLAP 引擎解决复杂回溯与关联问题 。
  • 流批一体演进‌:近期实践强调“流批一体”开发体验(一套逻辑编译为流/批任务),但底层执行仍根据需求区分流计算与批重跑机制,非理论上的纯 Kappa 。‌‌

技术选型与落地逻辑

  • 计算引擎‌:Flink(主流)、Storm(存量);‌存储与服务层‌:Doris(实时 OLAP,解决业务回撤与即席查询)、Redis(点查)、HBase(状态存储)。
  • 决策依据‌:日志数据量大连同态少,适合 Kappa;业务数据关联强、状态多变,纯流处理成本高且难维护,需引入批层或 OLAP 能力兜底准确性与回溯效率 。
  • 架构本质‌:是‌务实的混合架构‌,在统一数据接入(Kafka)和基础明细层构建后,上层计算链路根据业务特性分流,而非教条地套用单一架构范式 。‌‌

简言之,美团是"‌混合架构‌",仅在特定日志场景使用 Kappa 模式,关键业务场景保留了 Lambda 架构的批层优势或借助实时 OLAP 弥补纯流不足 。‌‌

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
┌─────────────────────────────────────────────────────────────────┐
│ 数据采集层 │
│ 业务数据库(MySQL) → Canal/Debezium → Kafka │
│ 用户行为日志 → Filebeat/Flume → Kafka │
│ 骑手GPS日志 → Flume → Kafka │
│ 第三方系统 → API网关 → Kafka │
└─────────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────────┐
│ 消息队列层 │
│ Kafka集群(多租户隔离,按业务线分区) │
└─────────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────────┐
│ 实时计算层 │
│ Flink集群(YARN部署,支持动态资源扩缩容) │
│ ├─ ODS层:原始数据清洗、格式转换 │
│ ├─ DWD层:明细数据标准化、维度关联、数据脱敏 │
│ ├─ DWS层:多维度预聚合、指标计算 │
│ └─ DIM层:实时维度表管理、缓慢变化维度处理 │
└─────────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────────┐
│ 数据存储层 │
│ ├─ 实时明细存储:ClickHouse │
│ ├─ 实时聚合存储:ClickHouse/Doris │
│ ├─ 维度数据存储:Redis/HBase │
│ ├─ 原始数据归档:HDFS/S3 │
│ └─ 离线数据修正:Spark │
└─────────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────────┐
│ 数据服务层 │
│ ├─ 实时查询引擎:ClickHouse JDBC/HTTP │
│ ├─ 数据API网关:统一接口封装、权限控制、限流熔断 │
│ └─ 数据订阅服务:Kafka消息订阅 │
└─────────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────────┐
│ 数据应用层 │
│ 实时大屏、商家后台、骑手APP、风控系统、营销系统、推荐系统 │
└─────────────────────────────────────────────────────────────────┘

2.2 核心组件与工具选型

层级 组件 选型理由
数据采集 Canal 阿里开源,MySQL CDC 采集成熟稳定,支持增量同步
数据采集 Filebeat 轻量级日志采集工具,资源占用低,与 Elastic 生态兼容
消息队列 Kafka 高吞吐量、高可靠性,支持百万级 QPS,Flink 原生支持
计算引擎 Flink 1.17+ 实时计算事实标准,支持 Exactly-Once、状态管理、CEP
实时存储 ClickHouse 23.3+ 列式存储,查询性能优异,适合实时 OLAP 分析
维度存储 Redis 7.0+ 高性能 KV 存储,支持毫秒级维度查询
离线计算 Spark 3.3+ 用于历史数据回溯、数据修正、离线指标验证
调度系统 Apache DolphinScheduler 分布式任务调度,支持 DAG、定时任务、依赖管理
元数据管理 Apache Atlas 数据血缘、数据字典、数据质量监控
监控告警 Prometheus + Grafana 全面监控集群状态、任务运行情况、数据质量
可视化 Apache Superset 开源 BI 工具,支持丰富的图表类型和交互式查询

三、数仓分层详细设计

3.1 ODS 层(原始数据层)

设计原则:保持数据原样,不做任何修改,便于数据回溯和问题排查。

数据来源与 Topic 设计

Topic 名称 数据来源 分区数 保留时间 数据量
ods_mysql_order_binlog 订单库 MySQL binlog 24 7 天 5000 万条 / 天
ods_mysql_user_binlog 用户库 MySQL binlog 12 7 天 1000 万条 / 天
ods_mysql_merchant_binlog 商家库 MySQL binlog 12 7 天 500 万条 / 天
ods_mysql_rider_binlog 骑手库 MySQL binlog 12 7 天 500 万条 / 天
ods_app_user_behavior APP 用户行为日志 48 3 天 10 亿条 / 天
ods_rider_gps_log 骑手 GPS 日志 24 1 天 5 亿条 / 天
ods_payment_log 支付系统日志 12 7 天 5000 万条 / 天

数据格式:统一使用 JSON 格式,包含ts(时间戳)、data(数据内容)、type(操作类型)、table(表名)等字段。

3.2 DWD 层(明细数据层)

设计原则:数据清洗、标准化、脱敏、维度关联,生成干净的明细数据。

核心处理逻辑

  1. 数据清洗:过滤脏数据、空值、异常值,去重
  2. 数据标准化:统一时间格式、统一编码格式、统一单位
  3. 数据脱敏:对手机号、身份证号、地址等敏感字段进行脱敏
  4. 维度关联:关联静态维度表(地区、品类、渠道等)
  5. 数据分流:按业务线和数据类型分流到不同的 DWD 表

核心 DWD 表设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
-- 订单明细事实表
CREATE TABLE dwd_order_info_di (
order_id STRING COMMENT '订单ID',
user_id STRING COMMENT '用户ID',
merchant_id STRING COMMENT '商家ID',
rider_id STRING COMMENT '骑手ID',
order_time TIMESTAMP COMMENT '下单时间',
pay_time TIMESTAMP COMMENT '支付时间',
accept_time TIMESTAMP COMMENT '商家接单时间',
fetch_time TIMESTAMP COMMENT '骑手取餐时间',
finish_time TIMESTAMP COMMENT '订单完成时间',
order_amount DECIMAL(10,2) COMMENT '订单金额',
pay_amount DECIMAL(10,2) COMMENT '支付金额',
order_status INT COMMENT '订单状态:1-待支付 2-待接单 3-待取餐 4-配送中 5-已完成 6-已取消',
province_code STRING COMMENT '省份编码',
city_code STRING COMMENT '城市编码',
area_code STRING COMMENT '区域编码',
category_id STRING COMMENT '品类ID',
channel_id STRING COMMENT '渠道ID',
dt STRING COMMENT '日期分区',
hour STRING COMMENT '小时分区'
) COMMENT '订单明细事实表'
PARTITIONED BY (dt, hour)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

-- 用户行为明细事实表
CREATE TABLE dwd_user_behavior_di (
user_id STRING COMMENT '用户ID',
device_id STRING COMMENT '设备ID',
session_id STRING COMMENT '会话ID',
event_type STRING COMMENT '事件类型:view-浏览 click-点击 add_cart-加购 purchase-下单',
event_time TIMESTAMP COMMENT '事件时间',
page_id STRING COMMENT '页面ID',
item_id STRING COMMENT '商品ID',
merchant_id STRING COMMENT '商家ID',
stay_time INT COMMENT '停留时长(毫秒)',
province_code STRING COMMENT '省份编码',
city_code STRING COMMENT '城市编码',
dt STRING COMMENT '日期分区',
hour STRING COMMENT '小时分区'
) COMMENT '用户行为明细事实表'
PARTITIONED BY (dt, hour)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

3.3 DWS 层(聚合数据层)

设计原则:按业务主题和维度预聚合,减少上层查询压力,提高查询速度。

聚合粒度

  • 时间粒度:1 分钟、5 分钟、1 小时、1 天
  • 维度粒度:全国、省份、城市、区域、商家、骑手、品类、渠道

核心 DWS 表设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
-- 商家维度1分钟聚合表
CREATE TABLE dws_merchant_order_1min (
merchant_id STRING COMMENT '商家ID',
city_code STRING COMMENT '城市编码',
category_id STRING COMMENT '品类ID',
window_start TIMESTAMP COMMENT '窗口开始时间',
window_end TIMESTAMP COMMENT '窗口结束时间',
order_cnt BIGINT COMMENT '订单量',
total_amount DECIMAL(10,2) COMMENT '总交易额',
pay_cnt BIGINT COMMENT '支付订单量',
cancel_cnt BIGINT COMMENT '取消订单量',
avg_order_amount DECIMAL(10,2) COMMENT '平均客单价',
dt STRING COMMENT '日期分区'
) COMMENT '商家维度1分钟订单聚合表'
PARTITIONED BY (dt)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

-- 城市维度1小时聚合表
CREATE TABLE dws_city_order_1h (
city_code STRING COMMENT '城市编码',
province_code STRING COMMENT '省份编码',
window_start TIMESTAMP COMMENT '窗口开始时间',
window_end TIMESTAMP COMMENT '窗口结束时间',
order_cnt BIGINT COMMENT '订单量',
total_amount DECIMAL(10,2) COMMENT '总交易额',
rider_cnt BIGINT COMMENT '在线骑手数',
avg_delivery_time INT COMMENT '平均配送时长(分钟)',
dt STRING COMMENT '日期分区'
) COMMENT '城市维度1小时订单聚合表'
PARTITIONED BY (dt)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

3.4 DIM 层(维度数据层)

设计原则:统一管理所有维度数据,支持实时更新和缓慢变化维度处理。

维度类型与处理方式

维度类型 示例 更新频率 处理方式
静态维度 地区、品类、渠道 天级 / 周级 全量加载 + 广播 Join
缓慢变化维度 商家信息、用户信息 小时级 / 天级 Flink 状态 + Redis 缓存
快速变化维度 骑手位置、订单状态 秒级 实时流处理 + 状态管理

核心维度表设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
-- 商家维度表
CREATE TABLE dim_merchant_info (
merchant_id STRING COMMENT '商家ID',
merchant_name STRING COMMENT '商家名称',
city_code STRING COMMENT '城市编码',
category_id STRING COMMENT '品类ID',
address STRING COMMENT '商家地址',
phone STRING COMMENT '商家电话',
status INT COMMENT '商家状态:1-营业中 2-休息中 3-关闭',
create_time TIMESTAMP COMMENT '创建时间',
update_time TIMESTAMP COMMENT '更新时间',
start_date STRING COMMENT '生效日期',
end_date STRING COMMENT '失效日期',
is_latest BOOLEAN COMMENT '是否最新版本'
) COMMENT '商家维度表'
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

3.5 ADS 层(应用数据层)

设计原则:直接面向业务应用,提供最终的指标数据。

核心 ADS 表设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
-- 全国实时运营大屏表
CREATE TABLE ads_national_realtime_dashboard (
ts TIMESTAMP COMMENT '统计时间',
total_order_cnt BIGINT COMMENT '今日总订单量',
total_amount DECIMAL(10,2) COMMENT '今日总交易额',
online_rider_cnt BIGINT COMMENT '在线骑手数',
avg_delivery_time INT COMMENT '平均配送时长(分钟)',
order_completion_rate DECIMAL(5,2) COMMENT '订单完成率(%)'
) COMMENT '全国实时运营大屏表'
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

-- 商家实时经营报表
CREATE TABLE ads_merchant_realtime_report (
merchant_id STRING COMMENT '商家ID',
dt STRING COMMENT '日期',
order_cnt BIGINT COMMENT '今日订单量',
total_amount DECIMAL(10,2) COMMENT '今日收入',
avg_order_amount DECIMAL(10,2) COMMENT '平均客单价',
cancel_rate DECIMAL(5,2) COMMENT '取消率(%)',
avg_meal_time INT COMMENT '平均出餐时长(分钟)'
) COMMENT '商家实时经营报表'
PARTITIONED BY (dt)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='ZSTD');

四、核心业务场景实现

4.1 实时订单统计

业务需求:实时统计全国、省份、城市、商家的订单量、交易额、订单状态分布。

数据流程

  1. ODS 层消费 Kafka 的订单 binlog 数据
  2. DWD 层清洗过滤,关联地区、品类维度
  3. DWS 层按不同维度和时间粒度预聚合
  4. ADS 层生成最终的实时报表数据
  5. 写入 ClickHouse 供实时大屏和商家后台查询

Flink 核心代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 1分钟滚动窗口聚合商家订单指标
val orderAggStream = dwdOrderStream
.keyBy("merchant_id", "city_code", "category_id")
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(
new OrderAggregateFunction(),
new OrderWindowFunction()
)

// 写入ClickHouse
orderAggStream
.addSink(ClickHouseSink.builder()
.setUrl("jdbc:clickhouse://clickhouse-host:8123/meituan_waimai")
.setUsername("default")
.setPassword("password")
.setTableName("dws_merchant_order_1min")
.build())

4.2 实时骑手运力监控

业务需求:实时监控各区域的骑手数量、待配送订单数、平均配送时长,为骑手调度提供数据支持。

数据流程

  1. ODS 层消费骑手 GPS 日志和订单状态变更日志
  2. DWD 层清洗过滤,关联地区维度
  3. DWS 层按区域和时间粒度聚合骑手运力指标
  4. ADS 层生成实时运力报表
  5. 通过 Kafka 消息推送给骑手调度系统

4.3 实时用户行为分析

业务需求:实时分析用户的浏览、点击、加购、下单行为,为个性化推荐和营销活动提供数据支持。

数据流程

  1. ODS 层消费 APP 用户行为日志
  2. DWD 层清洗过滤,关联用户、商品、商家维度
  3. DWS 层按用户、商品、商家维度聚合行为指标
  4. 实时更新用户行为标签到 Redis
  5. 推荐系统从 Redis 读取用户标签进行个性化推荐

4.4 实时风控系统

业务需求:实时检测异常订单、恶意用户、刷单作弊行为,保障平台安全。

数据流程

  1. ODS 层消费订单、支付、用户行为日志
  2. DWD 层清洗过滤,提取风控特征
  3. 使用 Flink CEP 进行复杂事件模式匹配
  4. 调用风控模型进行实时评分
  5. 对高风险订单进行拦截或预警

五、关键技术难点与解决方案

5.1 数据倾斜问题

问题场景

  • 头部商家(如麦当劳、蜜雪冰城)订单量占比过高
  • 热门城市(如北京、上海)订单量占比过高
  • 午晚高峰时段数据量暴增

解决方案

  1. 两阶段聚合(加盐法):给大 key 加上随机前缀,打散数据后再聚合
  2. 动态分区裁剪:只处理需要的分区数据,避免全表扫描
  3. Flink 自适应负载均衡:开启 Flink 1.17 + 的自适应调度功能
  4. 热点 Key 拆分:将热点 Key 拆分为多个子 Key,分别处理后再合并

5.2 实时维度关联问题

问题场景

  • 维度表数据量大,无法全量广播
  • 维度表更新频繁,需要实时同步
  • 大字段关联导致内存溢出

解决方案

  1. 广播小维度表:对于 < 100MB 的小维度表,使用广播 Join
  2. Redis Lookup Join:对于大维度表,将维度数据存储在 Redis 中,实时查询
  3. 预加载 + 定时刷新:将维度表预加载到 Flink 状态中,定时刷新
  4. 拆分 Join + 回表:先轻量 Join 找关联关系,再按需回表捞取大字段

5.3 数据一致性问题

问题场景

  • 数据丢失或重复
  • 实时与离线指标不一致
  • 订单状态变更导致数据不准确

解决方案

  1. Exactly-Once 语义保证:使用 Flink 的 Checkpoint 和 Kafka 的事务性生产
  2. 幂等性写入:使用唯一主键写入 ClickHouse,避免重复数据
  3. 实时离线口径统一:抽取公共指标计算逻辑为 UDF,实时和离线共用
  4. 数据修正机制:每天凌晨用离线数据修正前一天的实时数据

5.4 状态管理问题

问题场景

  • Flink 状态过大,导致 Checkpoint 时间过长
  • 状态恢复慢,影响任务可用性
  • 状态内存溢出

解决方案

  1. RocksDB 状态后端:使用 RocksDB 作为状态后端,将状态存储在磁盘上
  2. 增量 Checkpoint:开启增量 Checkpoint,只保存变化的状态
  3. 状态 TTL:设置合理的状态过期时间,清理无用状态
  4. 状态分区:将大状态拆分为多个小状态,分散存储

六、开发流程与规范

6.1 项目开发流程

  1. 需求分析:明确业务需求、指标口径、数据来源、延迟要求
  2. 数据调研:梳理数据结构、数据质量、更新频率
  3. 架构设计:设计数仓分层、数据流程、组件选型
  4. 开发测试:编写 Flink 作业、SQL 脚本,在测试环境验证
  5. 性能测试:模拟峰值流量,测试系统性能和稳定性
  6. 部署上线:在生产环境部署任务,逐步切换流量
  7. 监控运维:监控任务运行状态、数据质量,及时处理问题
  8. 文档编写:编写需求文档、设计文档、操作手册

6.2 开发规范

  1. 命名规范:表名、字段名使用小写字母加下划线,见名知意
  2. SQL 规范:使用标准 SQL 语法,添加必要的注释,避免复杂嵌套
  3. 代码规范:使用 Scala 或 Java 编写 Flink 作业,遵循代码规范,添加注释
  4. 版本管理:所有代码和配置文件纳入 Git 管理,使用分支开发模式
  5. 测试规范:编写单元测试、集成测试,确保代码质量
  6. 部署规范:使用容器化部署,统一配置管理,自动化部署

七、监控与运维

7.1 监控体系

  1. 集群监控:监控 Kafka、Flink、ClickHouse 集群的 CPU、内存、磁盘、网络使用率
  2. 任务监控:监控 Flink 作业的运行状态、吞吐量、延迟、Checkpoint 时间
  3. 数据监控:监控数据量、数据质量、指标准确性
  4. 业务监控:监控核心业务指标的变化趋势,异常告警

7.2 告警机制

  1. 告警级别:分为紧急、重要、一般三个级别
  2. 告警方式:短信、邮件、企业微信、电话
  3. 告警阈值:根据业务需求设置合理的告警阈值
  4. 告警升级:对于未及时处理的告警,自动升级到更高级别

7.3 运维流程

  1. 日常巡检:每天检查集群和任务运行状态
  2. 故障处理:及时处理故障,记录故障原因和解决方案
  3. 容量规划:根据业务增长情况,提前规划集群容量
  4. 版本升级:定期升级组件版本,修复已知问题,提升性能

八、面试高频问题及答案

8.1 架构与设计

Q1:为什么选择 Kappa 架构而不是 Lambda 架构?

A1:Lambda 架构需要维护实时和离线两套链路,开发和维护成本高,且容易出现口径不一致的问题。Kappa 架构采用统一的流处理引擎处理所有数据,开发和维护成本低,口径统一。现在 Flink 已经非常成熟,支持批流一体,可以很好地处理历史数据回溯和修正的需求。

Q2:实时数仓和离线数仓的区别是什么?

A2:

  • 延迟要求:实时数仓延迟在秒级甚至毫秒级,离线数仓延迟在小时级或天级
  • 数据处理方式:实时数仓处理流数据,离线数仓处理批数据
  • 计算引擎:实时数仓主要使用 Flink,离线数仓主要使用 Spark
  • 存储系统:实时数仓主要使用 ClickHouse、Doris 等 OLAP 数据库,离线数仓主要使用 Hive
  • 应用场景:实时数仓用于实时监控、实时决策、实时风控等场景,离线数仓用于历史数据分析、报表生成、数据挖掘等场景

Q3:为什么选择 ClickHouse 作为实时数仓的存储引擎?

A3:ClickHouse 是一款列式存储的 OLAP 数据库,具有以下优势:

  • 查询性能优异:支持秒级甚至毫秒级的查询响应
  • 高吞吐量:支持每秒数百万条数据的写入和查询
  • 压缩比高:列式存储 + 高效压缩算法,节省存储空间
  • 支持 SQL:支持标准 SQL 语法,学习成本低
  • 开源免费:社区活跃,生态完善

Q4:Flink 的 Exactly-Once 语义是怎么实现的?

A4:Flink 的 Exactly-Once 语义主要通过以下三个方面实现:

  1. Checkpoint 机制:定期将作业的状态保存到持久化存储中,故障时可以从最近的 Checkpoint 恢复
  2. 两阶段提交(2PC):对于支持事务的外部系统(如 Kafka、MySQL),Flink 使用两阶段提交协议保证数据的原子性写入
  3. 幂等性写入:对于不支持事务的外部系统(如 HDFS),Flink 通过幂等性写入保证数据不重复

Q5:如何处理 Flink 的数据倾斜问题?

A5:处理 Flink 数据倾斜的方法主要有:

  1. 预聚合:在 shuffle 之前先进行局部聚合,减少 shuffle 的数据量
  2. 加盐法:给大 key 加上随机前缀,打散数据后再聚合
  3. 动态负载均衡:开启 Flink 的自适应调度功能,自动将数据均匀分配到各个 Task
  4. 拆分大 key:将热点 key 拆分为多个子 key,分别处理后再合并
  5. 调整并行度:合理设置作业的并行度,避免单个 Task 处理过多数据

Q6:Flink 的状态后端有哪些,怎么选择?

A6:Flink 的状态后端主要有三种:

  1. MemoryStateBackend:将状态存储在 TaskManager 的内存中,速度快,但容量有限,适合小状态场景
  2. FsStateBackend:将状态存储在文件系统中,支持大状态,但 Checkpoint 时间较长
  3. RocksDBStateBackend:将状态存储在 RocksDB 中,支持超大状态,支持增量 Checkpoint,适合大状态场景

选择建议:

  • 小状态(<1GB):使用 MemoryStateBackend
  • 中等状态(1GB-10GB):使用 FsStateBackend
  • 大状态(>10GB):使用 RocksDBStateBackend

8.3 业务与实践

Q7:美团外卖实时数仓如何处理午晚高峰的流量峰值?

A7:主要通过以下方式处理流量峰值:

  1. 集群弹性扩缩容:基于 YARN 的动态资源调度,高峰时自动增加资源,低峰时释放资源
  2. 流量削峰:使用 Kafka 作为缓冲层,削峰填谷,避免计算引擎被峰值流量冲垮
  3. 预聚合:在 DWD 层和 DWS 层进行多层预聚合,减少上层查询的计算量
  4. 数据分流:将不同业务线的数据分流到不同的 Flink 作业,避免相互影响
  5. 降级策略:当系统负载过高时,自动降级非核心指标,保证核心指标的正常运行

Q8:如何保证实时数仓的数据质量?

A8:主要通过以下方式保证数据质量:

  1. 数据校验:在 DWD 层进行数据校验,过滤脏数据、空值、异常值
  2. 数据监控:监控数据量、数据延迟、指标准确性,异常告警
  3. 数据对比:每天将实时数据与离线数据进行对比,验证指标一致性
  4. 数据回溯:保留原始数据,支持数据回溯和问题排查
  5. 数据审计:记录数据的处理流程和变更历史,便于审计和问题定位

需要我把这个项目整理成一份可直接用于简历的项目描述,并补充3 道 P7 级别的深度面试题及答案吗?

实时风控系统

P7 级别深度面试题及参考答案

面试题 1:美团外卖午晚高峰(11:00-13:00、17:00-19:00)流量是平时的 5-10 倍,你如何设计系统架构保障大流量峰值下的稳定性?

参考答案

这是外卖实时数仓最核心的挑战之一,我会从事前、事中、事后三个维度进行全链路保障:

1. 事前:容量规划与弹性准备
  • 精准容量预估:基于历史数据建立流量预测模型,提前 7 天预测大促和日常高峰的流量峰值,预留 30% 的冗余资源
  • 集群弹性扩缩容:基于 YARN 的动态资源调度和 Flink 的自适应调度功能,实现资源的自动扩缩容。高峰前 1 小时自动扩容 50% 的资源,高峰结束后 1 小时自动释放
  • 预压测:每次大促前进行全链路压测,模拟 1.5 倍峰值流量,找出系统瓶颈并提前优化
  • 数据预热:提前将热点维度数据(如热门商家、热门城市)加载到 Redis 和 Flink 状态中,避免高峰时大量冷查询
2. 事中:流量控制与降级策略
  • 多层流量削峰

    • Kafka 层:设置合理的分区数和副本数,开启消息压缩,使用分区限流功能避免单个分区过载
    • Flink 层:使用反压机制,当下游处理能力不足时,自动减慢上游数据消费速度
  • 分级降级策略

    :制定明确的降级规则,按优先级从低到高依次降级:

    • 一级降级(非核心):关闭用户行为分析、推荐数据等非核心指标的计算
    • 二级降级(次核心):降低非核心维度的聚合粒度(如从 1 分钟改为 5 分钟)
    • 三级降级(核心):只保留全国、省份级别的核心指标,关闭城市、商家级别的细粒度指标
  • 热点隔离:将头部商家、热门城市的流量单独拆分到独立的 Flink 作业和 Kafka 主题中,避免热点影响整体系统

  • 大字段过滤:在 ODS 层就过滤掉不需要的大字段(如用户行为日志中的原始请求体),减少网络传输和内存占用

3. 事后:故障复盘与持续优化
  • 全链路监控:建立从数据采集、消息队列、计算引擎到存储系统的全链路监控体系,实时监控吞吐量、延迟、错误率等关键指标
  • 快速故障恢复:使用 Flink 的 Savepoint 机制,定期保存作业状态,故障时可以在 5 分钟内恢复
  • 故障复盘:每次高峰后进行故障复盘,记录问题原因和解决方案,持续优化系统架构

面试题 2:外卖业务中订单状态会频繁变更(下单→支付→接单→取餐→送达→取消),如何保证实时数仓中订单数据的一致性和准确性?

参考答案

订单状态的频繁变更是外卖实时数仓最复杂的问题之一,我会通过以下多层保障机制来解决:

1. 数据采集层:保证变更数据的完整性
  • 使用 Canal 采集 MySQL 的 binlog 数据,开启ROW 模式,记录每一行数据的完整变更前和变更后的值
  • 配置 Canal 的重试机制和故障转移,确保 binlog 数据不丢失
  • 在 Kafka 中设置足够的保留时间(至少 7 天),便于数据回溯和重放
2. 计算层:实现 Exactly-Once 语义和幂等性处理
  • Exactly-Once 语义

    • 开启 Flink 的 Checkpoint 机制,设置合理的 Checkpoint 间隔(1 分钟)和超时时间(10 分钟)
    • 使用 Flink 的两阶段提交(2PC)Sink,保证数据写入外部系统的原子性
    • 对于不支持事务的存储系统(如 ClickHouse),使用幂等性写入,通过唯一主键(order_id + update_time)避免重复数据
  • 订单状态合并

    • 使用 Flink 的状态管理,维护每个订单的最新状态
    • 当收到订单状态变更事件时,更新状态中的订单信息,并输出最新的完整订单数据
    • 设置状态 TTL(24 小时),清理已完成的订单状态,避免状态膨胀
3. 口径统一层:实时与离线口径对齐
  • 抽取公共的订单状态转换逻辑和指标计算逻辑为 UDF,实时和离线作业共用同一套 UDF
  • 建立实时离线对比机制:每天凌晨用离线数据修正前一天的实时数据,对比两者的差异,找出不一致的原因并优化
  • 对于跨天的订单(如 23:59 下单,00:01 支付),统一按下单时间归属到前一天,避免数据拆分
4. 数据质量层:全链路监控与校验
  • 数据完整性校验:监控每个环节的数据量,确保数据没有丢失
  • 数据准确性校验:对比实时和离线的核心指标(如订单量、交易额),当差异超过 0.5% 时自动告警
  • 数据一致性校验:校验订单状态的流转是否符合业务规则(如未支付的订单不能直接变为已完成)
  • 数据回溯机制:保留原始的 binlog 数据,支持任意时间点的数据回溯和重算

面试题 3:随着美团外卖业务的快速发展,新业务线(如闪购、买药、跑腿)不断接入,新指标需求层出不穷,如何设计实时数仓的架构来保证其扩展性和可维护性?

参考答案

扩展性和可维护性是实时数仓长期发展的关键,我会从架构设计、开发规范、工具链建设三个方面来解决:

1. 分层架构设计:高内聚低耦合
  • 公共层复用

    • 建设统一的 ODS 层和 DWD 层,所有业务线共用同一套原始数据和明细数据
    • 在 DWS 层建设公共聚合层,按业务主题(订单、用户、商家、骑手)和通用维度(时间、地区、品类)进行预聚合,供上层多个业务线复用
    • 避免每个业务线都从 ODS 层开始处理,减少重复计算和数据不一致
  • 业务层隔离

    • 在 ADS 层按业务线进行隔离,每个业务线有自己独立的应用层表
    • 业务线之间通过公共层进行数据交互,避免直接依赖
  • 插件化设计

    • 将指标计算逻辑封装成独立的插件,新指标只需开发对应的插件即可,无需修改核心代码
    • 支持动态加载和卸载插件,实现指标的热更新
2. 标准化开发规范:统一开发流程
  • 命名规范:制定统一的表名、字段名、主题名命名规范,见名知意
  • 数据建模规范:采用维度建模方法,统一事实表和维度表的设计规范
  • 代码规范:制定 Flink 作业和 SQL 的开发规范,使用模板化开发,减少代码冗余
  • 版本管理规范:所有代码和配置文件纳入 Git 管理,使用语义化版本号,支持版本回滚
3. 自动化工具链建设:提升开发效率
  • 元数据管理平台

    • 建设统一的元数据管理平台,管理所有表的元数据、数据血缘、数据字典
    • 支持自动解析 Flink SQL 和作业代码,生成数据血缘关系
    • 提供元数据搜索和查询功能,方便开发人员快速找到需要的数据
  • 指标管理平台

    • 建设统一的指标管理平台,对所有指标进行统一管理,包括指标定义、口径、计算公式、负责人
    • 支持指标的自动生成和上线,开发人员只需在平台上配置指标的维度和度量,即可自动生成对应的 Flink 作业
  • 自动化部署平台

    • 建设 CI/CD 流水线,实现代码的自动构建、测试、部署
    • 支持一键部署和回滚,减少人工操作失误
  • 监控运维平台

    • 建设统一的监控运维平台,集中管理所有 Flink 作业的运行状态、性能指标、告警信息
    • 支持作业的自动重启和故障转移,减少运维工作量

通过以上设计,我们的实时数仓可以支持新业务线在 3 天内完成接入,新指标在 4 小时内上线,同时保证系统的可维护性和稳定性。

HR 面高频问题及回答模板(结合本项目)

所有回答均采用STAR 法则(情境 - 任务 - 行动 - 结果),突出个人贡献、解决问题的能力和业务价值。


问题 1:你在这个美团外卖实时数仓项目中遇到的最大挑战是什么?你是如何解决的?

回答模板

情境 (S):美团外卖午晚高峰流量是平时的 8-10 倍,2025 年 618 大促期间峰值订单量突破 1 亿单 / 天。原有系统在高峰时出现严重的数据倾斜和任务堆积,核心指标延迟从 5 秒飙升至 30 分钟以上,甚至出现任务崩溃的情况,严重影响了全国运营大屏和骑手调度系统的正常运行。

任务 (T):我作为项目技术负责人,需要在 1 个月内解决大流量峰值下的系统稳定性问题,保证大促期间核心指标延迟 < 10 秒,系统可用性达到 99.99%。

行动 (A)

  1. 全链路压测定位瓶颈:搭建全链路压测平台,模拟 1.5 倍峰值流量,发现头部商家数据倾斜和大维度关联是主要瓶颈
  2. 热点隔离与打散:将订单量前 100 的头部商家流量单独拆分到独立的 Kafka 主题和 Flink 作业中,使用加盐法打散大 key 数据
  3. 分层预聚合优化:在 DWD 层增加 10 秒粒度的预聚合,减少 DWS 层的计算量
  4. 弹性资源调度:基于 YARN 实现资源的自动扩缩容,高峰前 1 小时自动扩容 50% 的资源
  5. 分级降级策略:制定了三级降级规则,当系统负载过高时自动关闭非核心指标的计算

结果 ®

  • 成功支撑了 2025 年 618 和双 11 大促,核心指标延迟稳定在 5 秒以内
  • 系统可用性从原来的 99.5% 提升至 99.99%
  • 集群资源利用率提升了 45%,计算成本降低了 30%
  • 该方案被推广到公司其他业务线,成为大流量实时数仓的标准解决方案

问题 2:这个项目中你最有成就感的部分是什么?为什么?

回答模板

情境 (S):美团外卖原有实时和离线数仓是两套独立的链路,使用不同的计算引擎和代码逻辑,导致实时和离线指标经常出现不一致的情况,差异最大时达到 10% 以上。业务部门经常质疑数据的准确性,数据团队需要花费大量时间排查和解释差异。

任务 (T):我负责设计并实现批流一体的实时数仓架构,统一实时和离线数据口径,将指标差异控制在 0.5% 以内。

行动 (A)

  1. 架构选型:采用纯 Kappa 架构替代传统 Lambda 架构,使用 Flink 作为统一的计算引擎
  2. 逻辑统一:抽取所有公共的指标计算逻辑为 UDF,实时和离线作业共用同一套 UDF
  3. 数据对齐:统一时间时区、数据过滤条件和指标定义,建立了严格的数据口径规范
  4. 对比验证机制:开发了实时离线数据对比工具,每天自动对比核心指标的差异,异常时自动告警

结果 ®

  • 实时和离线核心指标差异从原来的 10% 以上降低到 0.3% 以内
  • 彻底解决了数据口径不一致的问题,业务部门对数据的信任度大幅提升
  • 数据团队的运维工作量减少了 60%,不再需要花费大量时间排查数据差异
  • 该成果获得了公司年度技术创新奖

问题 3:你在项目中是如何与团队成员协作的?遇到过什么冲突吗?怎么解决的?

回答模板

情境 (S):这个项目涉及数据开发、运维、业务分析、产品等多个团队,共 15 名成员。在项目初期,由于各团队对需求的理解不一致,导致开发进度缓慢,出现了多次返工的情况。

任务 (T):我作为项目技术负责人,需要协调各团队的工作,解决团队之间的冲突,保证项目按时交付。

行动 (A)

  1. 建立统一的沟通机制:每周召开一次项目例会,每天进行 15 分钟的站会,及时同步进度和问题
  2. 明确职责分工:制定了详细的项目计划和职责分工表,明确每个团队和个人的任务和交付时间
  3. 需求评审机制:所有需求都需要经过产品、技术、业务三方评审,确保需求的准确性和可行性
  4. 冲突解决:当出现冲突时,我会组织相关人员进行面对面沟通,从业务价值出发,找到各方都能接受的解决方案。例如,在指标口径的问题上,业务部门希望指标越详细越好,而技术部门担心性能问题。我提出了分层聚合的方案,既满足了业务对细粒度指标的需求,又保证了系统的性能。

结果 ®

  • 项目提前 2 周完成交付,所有功能都达到了预期的效果
  • 团队之间的沟通效率大幅提升,没有再出现因为需求理解不一致导致的返工
  • 建立了良好的团队协作氛围,项目结束后团队成员的满意度达到了 95% 以上

问题 4:通过这个项目,你最大的收获和成长是什么?

回答模板

通过这个项目,我在技术、业务和管理三个方面都获得了很大的成长:

  1. 技术能力
    • 深入掌握了 Flink、Kafka、ClickHouse 等大数据技术的底层原理和最佳实践
    • 学会了如何设计和实现高可用、高并发、低延迟的实时数仓架构
    • 积累了处理大流量峰值、数据倾斜、状态膨胀等复杂技术问题的经验
  2. 业务理解
    • 深入了解了美团外卖的核心业务流程和业务痛点
    • 学会了从业务角度思考问题,将技术方案与业务需求紧密结合
    • 能够准确地将业务需求转化为技术方案,并评估技术方案的业务价值
  3. 管理能力
    • 提升了项目管理和团队协作能力,能够带领 10 人以上的团队完成复杂的技术项目
    • 学会了如何进行有效的沟通和协调,解决团队之间的冲突
    • 培养了风险意识和问题解决能力,能够提前识别项目风险并制定应对措施

这个项目让我从一个单纯的技术开发人员成长为一个能够独当一面的技术负责人,为我未来的职业发展打下了坚实的基础。


问题 5:如果让你重新做这个项目,你会在哪些方面进行改进?

回答模板

如果让我重新做这个项目,我会在以下几个方面进行改进:

  1. 更早地引入数据治理
    • 在项目初期就建立完善的数据治理体系,包括元数据管理、数据质量监控、数据安全等
    • 避免在项目后期因为数据质量问题花费大量时间进行整改
  2. 更完善的自动化工具链
    • 提前建设自动化的指标生成平台和部署平台
    • 进一步提升开发效率,减少人工操作失误
  3. 更充分的预研和测试
    • 在项目初期对关键技术进行更充分的预研和测试
    • 避免在项目中期因为技术选型问题导致的架构调整
  4. 更重视用户体验
    • 在项目初期就与业务用户进行充分的沟通,了解他们的真实需求
    • 提供更友好的数据查询和可视化界面,提升用户体验
  5. 更长远的架构规划
    • 在架构设计时考虑未来 3-5 年的业务发展
    • 预留足够的扩展空间,避免因为业务快速发展导致的架构重构

技术答辩 PPT 大纲(P6-P7 级别,15-20 页)

封面(第 1 页)

  • 标题:美团外卖实时数仓建设与优化
  • 副标题:支撑亿级订单的高可用低延迟实时数据平台
  • 汇报人:XXX
  • 日期:XXXX 年 XX 月 XX 日

目录(第 2 页)

  1. 项目背景与目标
  2. 整体架构设计
  3. 数仓分层详细设计
  4. 核心业务场景实现
  5. 关键技术难点与解决方案
  6. 项目成果与业务价值
  7. 未来规划与展望
  8. Q&A

一、项目背景与目标(第 3-4 页)

3.1 业务背景

  • 美团外卖业务规模:日均订单 6000 万 +,峰值 1 亿 + 单 / 天
  • 原有系统痛点:
    • 离线数仓延迟高(T+1),无法满足实时决策需求
    • Lambda 架构双链路维护成本高,口径不一致
    • 大流量峰值下系统不稳定,经常出现任务崩溃
    • 新指标上线周期长(3 天 +),无法快速响应业务需求

3.2 项目目标

  • 性能目标:核心指标延迟 < 5 秒,支持百万级 QPS
  • 可靠性目标:系统可用性 99.99%,数据不丢不重
  • 一致性目标:实时离线指标差异 < 0.5%
  • 效率目标:新指标上线 < 4 小时,新业务线接入 < 3 天

二、整体架构设计(第 5-6 页)

4.1 架构选型:纯 Kappa 架构

  • 对比 Lambda 架构的优势:统一计算引擎、统一数据口径、降低维护成本
  • 批流一体实现:Flink 同时处理实时流和历史批数据

4.2 整体架构图

  • 数据采集层:Canal、Filebeat、Flume
  • 消息队列层:Kafka(多租户隔离)
  • 实时计算层:Flink 1.17(YARN 部署)
  • 数据存储层:ClickHouse、Redis、HDFS
  • 数据服务层:API 网关、数据订阅
  • 数据应用层:实时大屏、商家后台、骑手 APP、风控系统

4.3 核心组件选型理由

  • Flink:Exactly-Once 语义、强大的状态管理、CEP 支持
  • ClickHouse:列式存储、高查询性能、高压缩比
  • Kafka:高吞吐量、高可靠性、Flink 原生支持

三、数仓分层详细设计(第 7-8 页)

5.1 分层设计原则

  • 高内聚低耦合
  • 数据复用最大化
  • 口径统一
  • 易于扩展

5.2 各层详细设计

  • ODS 层:原始数据原样落地,保留 7 天,支持数据回溯
  • DWD 层:数据清洗、标准化、脱敏、维度关联,生成干净的明细数据
  • DWS 层:按业务主题和维度预聚合,减少上层查询压力
  • DIM 层:统一管理维度数据,支持实时更新和缓慢变化维度处理
  • ADS 层:直接面向业务应用,提供最终的指标数据

5.3 核心表示例

  • 订单明细事实表(dwd_order_info_di)
  • 商家维度 1 分钟聚合表(dws_merchant_order_1min)
  • 全国实时运营大屏表(ads_national_realtime_dashboard)

四、核心业务场景实现(第 9-10 页)

6.1 实时订单统计

  • 数据流程:ODS 订单 binlog → DWD 清洗 → DWS 多维度聚合 → ADS 实时报表
  • 核心指标:订单量、交易额、订单状态分布、平均客单价
  • 延迟要求:<5 秒

6.2 实时骑手运力监控

  • 数据流程:ODS 骑手 GPS 日志 + 订单状态日志 → DWD 清洗 → DWS 区域聚合 → 调度系统
  • 核心指标:区域骑手数、待配送订单数、平均配送时长
  • 延迟要求:<1 秒

6.3 实时风控系统

  • 数据流程:ODS 多源日志 → DWD 特征提取 → CEP 模式匹配 → 模型评分 → 风险拦截
  • 核心能力:异常订单检测、恶意用户识别、刷单作弊拦截
  • 延迟要求:<100 毫秒

五、关键技术难点与解决方案(第 11-14 页)

7.1 大流量峰值处理

  • 问题:午晚高峰流量是平时的 8-10 倍,系统容易崩溃
  • 解决方案:
    • 多层流量削峰(Kafka 缓冲 + Flink 反压)
    • 弹性资源扩缩容
    • 分级降级策略
    • 热点隔离
  • 效果:成功支撑 1 亿 + 单 / 天的峰值流量,核心指标延迟稳定在 5 秒以内

7.2 数据倾斜问题

  • 问题:头部商家、热门城市数据量占比过高,导致任务倾斜
  • 解决方案:
    • 两阶段聚合(加盐法)
    • 热点 Key 拆分
    • 动态负载均衡
  • 效果:任务运行时间缩短 70%,不再出现单个 Task 卡死的情况

7.3 订单状态一致性问题

  • 问题:订单状态频繁变更,容易出现数据不一致
  • 解决方案:
    • Exactly-Once 语义保证(Checkpoint + 2PC)
    • 幂等性写入
    • 实时离线口径统一
    • 数据修正机制
  • 效果:实时离线指标差异 < 0.3%,数据准确性达到 99.9%

7.4 大维度关联问题

  • 问题:用户标签等大维度表无法全量广播,关联效率低
  • 解决方案:
    • Redis Lookup Join + 本地缓存
    • 拆分 Join + 回表拼接
    • 预加载 + 定时刷新
  • 效果:关联性能提升 5 倍,内存占用减少 80%

六、项目成果与业务价值(第 15-16 页)

8.1 技术成果

  • 核心指标延迟从 30 分钟降至 5 秒以内
  • 系统可用性从 99.5% 提升至 99.99%
  • 集群资源利用率提升 45%,计算成本降低 30%
  • 新指标上线周期从 3 天缩短至 4 小时
  • 新业务线接入时间从 2 周缩短至 3 天

8.2 业务价值

  • 实时风控系统拦截异常订单率提升 20%,每年减少损失数千万元
  • 骑手平均配送时长缩短 8%,用户满意度提升 5%
  • 商家订单转化率提升 5%,平台交易额增长 3%
  • 运营决策效率提升 10 倍,从 T+1 决策变为实时决策

七、未来规划与展望(第 17 页)

  1. 批流一体深化:完全统一实时和离线计算链路,实现一套代码跑遍所有场景
  2. AI 赋能:引入机器学习算法,实现智能异常检测、智能指标预测、智能资源调度
  3. 数据治理升级:建设全链路数据治理平台,实现数据全生命周期管理
  4. 云原生改造:将系统迁移到 Kubernetes 上,实现更灵活的资源调度和更高的资源利用率
  5. 开放平台:建设实时数据开放平台,赋能更多业务线和合作伙伴

八、Q&A(第 18 页)

  • 感谢聆听
  • 欢迎提问

PPT 制作注意事项

  1. 简洁明了:每页 PPT 只讲一个核心点,文字不要太多,多用图表和流程图
  2. 突出重点:用加粗、颜色等方式突出关键数据和结论
  3. 数据支撑:所有成果都要有具体的数据支撑,避免空泛的描述
  4. 逻辑清晰:按照 “为什么做 - 怎么做 - 做得怎么样 - 未来怎么做” 的逻辑展开
  5. 准备充分:提前演练,熟悉每个部分的内容,准备好可能被问到的问题

技术选型

实时数仓业务库同步技术选型全指南(P7 面试版)

一、核心选型维度(P7 必须掌握的决策依据)

在进行技术选型前,必须先明确以下 7 个核心维度,这也是面试官一定会追问的点:

  1. 数据一致性:是否支持 Exactly-Once 语义,能否保证全量 + 增量同步的数据一致性
  2. 延迟要求:业务能接受的端到端延迟(毫秒级 / 秒级 / 分钟级)
  3. 吞吐量:单表每秒变更量(TPS)和总数据量
  4. 数据源支持:是否需要支持 MySQL、PostgreSQL、Oracle 等多种数据库
  5. 数据处理能力:是否需要在同步过程中进行过滤、转换、关联等复杂操作
  6. 运维成本:团队的技术栈和运维能力,是否能支撑复杂的分布式系统
  7. 生态兼容性:与下游实时计算引擎(Flink)、存储系统(Doris/ClickHouse)的集成度

二、主流 CDC 工具深度对比

四大核心工具对比表
特性 Canal Debezium Flink CDC SeaTunnel CDC
开源组织 阿里巴巴 Red Hat/CNCF Apache Apache
架构 自研 Server Kafka Connect Flink Source 统一数据集成框架
支持数据库 仅 MySQL MySQL/PG/Oracle/MongoDB 等 10 + 种 MySQL/PG/Oracle/SQL Server 等 30 + 种数据源
全量同步 不支持(需配合 DataX) 支持 支持(无缝全量转增量) 支持(批流一体)
Exactly-Once 需自行实现 支持(Kafka 事务) 原生支持(Flink Checkpoint) 支持
数据处理能力 弱(仅简单过滤) 弱(需配合 Kafka Streams) 强(Flink SQL/Table API) 中(内置转换算子)
与 Flink 集成 需通过 Kafka 中转 需通过 Kafka 中转 原生集成(直连) 原生集成
运维复杂度 低(单节点即可) 中(依赖 Kafka 集群) 低(复用 Flink 集群) 低(独立集群)
社区活跃度 高(国内) 极高(全球) 极高(全球) 快速增长
适用场景 中小规模 MySQL 同步、阿里技术栈 大规模多源同步、微服务事件总线 实时数仓、复杂 ETL、数据湖入湖 多源异构数据集成、批流一体
各工具核心优缺点详解
1. Canal(阿里开源)

核心优势

  • 轻量级,部署运维简单,单节点即可支撑每秒数万 TPS
  • 对 MySQL 版本兼容性极好,支持 5.5 到 8.0 所有版本
  • 国内社区活跃,中文资料丰富,问题容易解决
  • 支持直接输出到 Kafka、RocketMQ、Redis 等多种下游

核心缺点

  • 仅支持 MySQL,不支持其他数据库
  • 不支持全量同步,首次同步需要配合 DataX 等工具
  • 数据处理能力弱,复杂转换需要在下游实现
  • Exactly-Once 语义需要自行实现,容易出现数据丢失或重复
2. Debezium(最主流开源 CDC)

核心优势

  • 支持最丰富的数据源,几乎覆盖所有主流数据库
  • 与 Kafka 生态深度集成,是 Kafka Connect 的标准 CDC 连接器
  • 支持全量 + 增量无缝同步,自动处理断点续传
  • 高可用和可扩展性好,支持分布式部署
  • 全球社区活跃,是企业级 CDC 的事实标准

核心缺点

  • 强依赖 Kafka 集群,架构复杂度高,运维成本高
  • 数据处理能力弱,复杂 ETL 需要配合 Flink 或 Kafka Streams
  • 对 MySQL 的某些特殊特性支持不如 Canal 完善
  • 默认输出的 JSON 格式比较复杂,需要额外解析

核心优势

  • 彻底解决了传统 CDC 架构 “全量 + 增量割裂” 的痛点,支持无缝切换
  • 原生集成 Flink 生态,可以直接使用 Flink SQL 进行复杂的数据处理
  • 无需中间 Kafka,支持直连数据库和下游存储,架构更简单
  • 原生支持 Exactly-Once 语义,数据一致性有保障
  • 支持分布式部署,可线性扩展吞吐量

核心缺点

  • 对数据库的支持不如 Debezium 丰富,某些小众数据库支持不完善
  • 早期版本存在一些稳定性问题,建议使用 1.15 以上版本
  • 全量同步阶段对源库压力较大,需要合理配置并行度
  • 没有独立的运维界面,需要依赖 Flink 的监控体系
4. SeaTunnel CDC(新兴批流一体数据集成框架)

核心优势

  • 支持最多的数据源和目标,超过 300 种连接器
  • 批流一体,同一个任务既可以做离线全量同步,也可以做实时增量同步
  • 架构极简,无需依赖 Kafka,支持 Source 到 Sink 直连
  • 性能优异,单节点吞吐量可达每秒数十万条
  • 运维简单,有可视化的管理界面

核心缺点

  • 社区相对较新,生态不如 Flink CDC 成熟
  • 复杂数据处理能力不如 Flink CDC 强大
  • 某些高级特性还在开发中

三、不同场景下的推荐架构方案

架构图

1
业务数据库 → Canal/Debezium → Kafka → Flink → 实时数仓(Doris/ClickHouse)

适用场景

  • 大规模生产环境,要求极高的稳定性和可靠性
  • 数据需要被多个下游系统消费(实时数仓、缓存、搜索等)
  • 团队已经有成熟的 Kafka 和 Flink 运维经验
  • 数据源类型比较单一(主要是 MySQL)

推荐工具组合

  • MySQL 数据源:Canal(国内首选)或 Debezium(国际首选)
  • 多数据源:Debezium
  • 消息队列:Kafka
  • 实时计算:Flink

架构图

1
业务数据库 → Flink CDC → Flink → 实时数仓(Doris/ClickHouse)

适用场景

  • 中小规模实时数仓,追求架构简洁和低运维成本
  • 数据只需要进入实时数仓,不需要被多个系统消费
  • 需要在同步过程中进行复杂的数据处理和转换
  • 团队主要使用 Flink 技术栈

推荐工具组合

  • 所有支持的数据源:Flink CDC
  • 实时计算:Flink
  • 存储:Doris/ClickHouse
方案三:SeaTunnel 批流一体架构(最适合多源异构)

架构图

1
多源业务数据库 → SeaTunnel CDC → 实时数仓/数据湖

适用场景

  • 需要同步多种不同类型的数据源(MySQL、PG、Oracle、MongoDB 等)
  • 同时有离线和实时同步需求,希望统一技术栈
  • 团队没有足够的运维能力支撑复杂的 Kafka 和 Flink 集群
  • 追求极致的性能和低延迟

四、生产环境最佳实践

1. 全量同步优化
  • 分表分库同步:对于大表,使用分表分库并行同步,提高全量同步速度
  • 限流控制:全量同步阶段对源库进行限流,避免影响业务
  • 增量先行:先开启增量同步,再进行全量同步,最后合并数据,减少业务影响
  • 断点续传:确保所有工具都支持断点续传,避免同步失败后重新开始
2. 数据一致性保障
  • 开启 Exactly-Once 语义:Flink CDC 开启 Checkpoint,Debezium 开启事务
  • 幂等写入:下游存储使用幂等写入,避免重复数据
  • 数据校验:定期进行全量数据校验,确保源端和目标端数据一致
  • 事务支持:对于需要强一致性的场景,使用支持事务的下游存储(如 Doris)
3. 性能优化
  • 并行度配置:根据表的大小和变更量合理配置并行度
  • 分区策略:Kafka 分区数与 Flink 并行度保持一致,避免数据倾斜
  • 批量写入:下游存储开启批量写入,提高写入性能
  • 数据压缩:在 Kafka 和网络传输中使用压缩算法,减少带宽占用
4. 监控与运维
  • 全链路监控:监控 CDC 工具的延迟、吞吐量、错误率
  • 源库监控:监控源库的 CPU、内存、IO 和复制延迟
  • 告警机制:建立多级告警机制,及时发现和处理问题
  • 灾备方案:制定完善的灾备方案,确保数据安全

五、面试回答技巧(结合你的经历)

当面试官问你实时数仓同步业务库的技术选型时,你可以按照以下结构回答,突出你的 P7 级别能力:

  1. 先讲选型原则:“我在做技术选型时,会首先明确业务的核心需求,包括数据一致性要求、延迟要求、吞吐量要求,然后结合团队的技术栈和运维能力,综合评估各个方案的优缺点。”
  2. 对比主流方案:“目前主流的 CDC 方案有 Canal、Debezium 和 Flink CDC。Canal 轻量简单,适合 MySQL 同步;Debezium 支持多数据源,生态完善;Flink CDC 架构简洁,与 Flink 集成最好。”
  3. 结合实际经历:“在微软工作时,我们需要同步全球多个地区的 MySQL 和 PostgreSQL 数据库,并且需要进行复杂的数据处理和转换。我们最终选择了 Flink CDC 直连架构,因为它不需要中间 Kafka,架构更简单,而且可以直接使用 Flink SQL 进行数据处理。通过这套架构,我们将端到端延迟从原来的秒级降低到了毫秒级,同时运维成本降低了 50%。”
  4. 讲遇到的问题和解决方案:“在实施过程中,我们遇到了全量同步对源库压力大的问题。我们通过分表并行同步、限流控制和增量先行的策略,成功将源库的 CPU 使用率控制在 20% 以下,没有对业务造成任何影响。”
  5. 总结选型结论:“总的来说,如果是大规模多源同步场景,我会推荐 Debezium+Kafka+Flink 架构;如果是中小规模实时数仓,我会推荐 Flink CDC 直连架构,它更简洁高效。”

面试口述版

面向 10 年经验大数据开发 / 技术专家岗,口述时长约 2.5-3 分钟,适配资深专家级面试,全程预埋可深挖的技术亮点)

我主导过公司核心交易域实时数仓从 Lambda 架构到 Flink+Hudi 流批一体架构的全面升级,覆盖外卖、到店两大核心业务线,核心是解决原来双链路架构下口径不统一、研发效率低的刚性痛点。

一、项目背景与痛点

项目启动前,我们的实时数仓是典型的 Lambda 双链路模式:离线链路用 Hive+Spark 跑 T+1 指标,实时链路用 Flink+Kafka 做秒级计算,两套链路各写一套逻辑、各存一份数据。 当时日均处理万亿级埋点日志和业务 Binlog 数据,核心交易域有近 200 层数仓表,下游支撑经营作战大屏、实时风控、商家经营看板等 30 多个核心业务场景。但旧架构有三个绕不开的问题: 第一是口径不一致,GMV、有效订单量这类核心指标,实时和离线的差异率最高冲到 5%,业务方做决策不敢完全信实时数据,每次都要等离线对账; 第二是研发效率低,新增一个公共指标,要同时开发实时、离线两套代码,做两轮测试,上线周期至少 3 天,人力成本直接翻倍; 第三是历史回溯能力弱,纯实时链路依赖 Kafka 留存数据,要修正 7 天以上的历史指标,只能全量重放消息,时间和资源成本极高,根本满足不了业务频繁的补数、逻辑修正需求。 所以这个项目的核心目标很明确:实现一套计算逻辑、一份统一存储,流批口径完全对齐,同时把研发效率提上来。

二、整体架构设计与选型

架构上我们沿用了 ODS-DWD-DIM-DWS-APP 的标准数仓分层,核心做了两个关键升级: 一是计算引擎统一,全链路用 Flink 作为唯一核心引擎,同一套 Flink SQL 既可以跑流模式做实时增量计算,也可以跑批模式做全量 / 增量回溯; 二是存储层统一,核心分层数据全部落地 Hudi 数据湖,替代原来的 Kafka 分层存储 + Hive 离线表的双存储模式。 选型阶段我们对比了 Hudi、Iceberg、Delta Lake 三款产品,最终敲定 Hudi,核心有两个原因: 第一,我们交易域 70% 以上是 CDC 更新场景,Hudi 的 MOR 表对高频行级 Upsert 的支持更成熟,和 Flink 的生态适配也更完善; 第二,它原生支持增量读取,非常贴合我们做 “增量数仓” 的定位,下游可以按增量消费,不用全量扫表。 分层存储我们也做了差异化设计:纯追加的流量日志 DWD、静态维度表用 COW 表保证查询性能;高频更新的交易 / 履约 DWD、DWS 汇总层用 MOR 表保证写入性能,两边做权衡。

三、核心难点与落地

落地过程中,我牵头攻克了三个最核心的技术难点: 第一个是端到端的数据一致性保障。我们基于 Flink 的两阶段提交机制,配合 Hudi 的事务原子性,实现了全链路 Exactly-Once 语义;同时搭建了分层级的对账体系,从层间数据条数、到核心业务指标,自动做流批对账,从机制上兜底口径一致。 第二个是海量数据下的性能与稳定性。我们最大的日粒度聚合任务,状态超过 2TB,早期 Checkpoint 经常超时失败,高峰期还会出现 Hudi 写入积压。我们通过 RocksDB 增量 Checkpoint、状态 TTL 优化,配合 Hudi 侧的异步 Compaction 策略调优、小文件合并治理,最终把 Checkpoint 时长从 15 分钟压缩到 3 分钟以内,核心任务全年稳定性达到 99.9%。 第三个是迟到数据的完整性兜底。我们用 Watermark 处理常规乱序数据,设置了合理的容忍窗口;超迟到的数据走侧输出流单独捕获,结合 Hudi 的行级 Upsert 能力回写修正历史聚合结果,保证即使延迟数天的数据也能最终落对,全程不丢数据。

四、落地成果与个人贡献

项目上线后,不管技术还是业务层面收益都很明确: 技术上,核心指标的流批差异率从 5% 降到了 0.1% 以内,真正实现了口径统一;需求上线周期从 3 天缩短到半天,研发效率提升 60%;统一存储后,整体存储成本下降了 30%。 业务上,这套架构成功支撑了多次亿级大促的实时监控零故障,同时给实时风控、商家运营等场景提供了可靠的实时数据底座,风控策略的迭代效率直接提升了一倍。 我在项目里担任技术负责人,全程牵头架构设计、技术选型和核心难点攻关,同时制定了公司统一的实时数仓建模规范、开发标准和运维流程,沉淀的通用维表关联、增量补数等组件,后来在公司 3 个业务线都完成了复用。


面试配套提示

  1. 语速控制:正常语速讲完约 2 分 40 秒,符合面试单项目描述的最佳时长,不用刻意加快;
  2. 引导追问:讲完后可以补一句 “其中一致性保障和性能优化这块我们踩了不少坑,也做了很多定制化优化”,引导面试官往你准备充分的方向提问;
  3. 细节兜底:逐字稿里提到的每个技术点(2PC 机制、MOR/COW 选型、Compaction 策略、增量 Checkpoint、小文件治理等),都要提前准备好原理细节和踩坑经历,应对深挖。