改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

美团增量数仓建设新进展

ckckck2025年1月10日104 浏览

Apache Flink 中文社区

以下文章来源于 Apache Flink,作者 汤楚熙@美团

文章封面

Apache Flink 是 Apache Flink 中文社区唯一官方微信公众号,由 Flink PMC 维护。

摘要: 本文整理自美团系统研发工程师汤楚熙在 Flink Forward Asia 2022 实时湖仓专场的分享。主要内容分为四个部分:

  1. 建设背景
  2. 核心能力设计与优化
  3. 业务实践
  4. 未来展望

Tips: 点击 阅读原文 免费领取 5000CU*小时 Flink 云资源。


01 美团增量数仓的建设背景

美团数仓架构的诞生基于以下技术假设:“随着业务数据越积越多,增量数据 / 存量数据的比值呈下降趋势,采用增量计算模式性价比更高。” 同时,Flink、Hudi 等具备增量计算、更新能力的技术框架为增量数仓的落地提供了必要条件。

背景图

从时间线上看,增量数仓架构的演进过程可大致划分为三个阶段:

  • 第一阶段(2019-2020年): 业务希望在离线数仓的基础上获得更新鲜的数据,即实时数仓。我们借鉴了离线数仓的模型概念,提出了实时数仓的模型抽象。
  • 第二阶段(2020-2021年): 实时数仓的生产任务大量依赖 Java API,影响了开发效率。我们加快了 Flink SQL 的落地,提升了数仓开发效率。
  • 第三阶段(2021年至今): 随着数据湖技术的成熟,我们开始尝试整合离线和实时数仓架构,提出了增量数仓的新架构。
演进图

目前,美团内部有 M、B、C、D 端等四大类业务场景,不同场景对数据一致性和时效性的要求不同,需要寻找一套尽可能适配所有场景的技术架构。

我们首先想到的是 Lambda 架构,它通过实时链路解决高时效性需求,通过离线链路解决长业务周期的指标计算需求。然而,Lambda 架构的生产链路过于复杂,导致资源成本和运维成本高昂。

Lambda 架构问题

例如,高数据新鲜度场景高度依赖 Kafka,但其架构设计未充分考虑数据一致性问题。业务通过排序、幂等处理等手段牺牲计算资源来保证数据一致性。此外,运维门槛高,典型案例如美团某 B 端业务场景,要求 Flink 作业中保留 180 天状态数据,单任务状态大小超过 50TB。

运维问题

对于时延不敏感但需要灵活组织数据的离线场景,重度依赖 Hive,但其最初设计未考虑高效更新能力。例如,离线数仓最新快照事实表的生产场景,通常需要全量加载存量快照数据,再与增量数据做归并排序,效率低下。

离线数仓问题

我们期望增量数仓架构能够兼顾数据时效性和一致性,并低成本完成数据合并计算与高效组织。

增量数仓期望

02 核心能力设计与优化

2.1 增量数仓存储架构

为了实现增量计算和更新,我们引入了一套支持事务管理、主键和 CDC 能力的新存储引擎,内部称为 Beluga。Beluga 基于 Hudi 改造而来,主要改造动机是 Hudi v0.8 不支持 CDC,因此我们引入了 HBase 来生产 Changelog。

Beluga 架构

Beluga 包含三个核心模块:

  • Beluga Client: 运行在 Flink 作业中,处理读写请求和事务协调。
  • Beluga Server: 基于 HBase 改造,承担数据更新和 Changelog 生产能力。
  • Beluga File Store: 基于 Hudi 实现,存储 CDC 数据和快照数据。

2.2 优化分桶策略

Beluga 的增量行更新能力通过数据分桶实现。我们统一了 HBase 和 Hudi 的分桶模型,将 HBase 的 HRegion 和 Hudi 的 FileGroup 一一映射,共同组成一个分桶。新记录先经过 HBase 的 Region,然后按需生产 Changelog,刷入 Hudi。

分桶策略

这样做的好处是,Hudi 可以将 HBase 作为外部索引,提升数据更新效率。

分桶优化

前期测试发现,Hudi 原生分桶策略使用不当会导致性能问题。为此,Beluga 设计了一套固定分桶策略,有效控制了文件数的增长,并通过 HBase 减少了元数据和索引的拉取频率,提升了读写性能。

2.3 CDC 数据格式优化

Hudi 原生 CDC 能力依赖 Flink 的回撤机制,但测试中发现存在数据不一致的风险。为此,我们将 UPDATE 事件的 UPDATE_BEFORE 和 UPDATE_AFTER 合并到一条记录中,类似 MySQL 的 Binlog,保证了更好的原子性。

CDC 优化

2.4 扩展有状态计算场景

Beluga 还可以低成本解决一些有状态计算场景问题。例如,长周期多流数据关联时,Flink 状态中需要保留长时间的数据快照,导致状态量过大,影响 Checkpoint 稳定性。通过 Beluga 的 HBase 更新能力,可以缓解 Checkpoint 压力,提升资源利用率。

有状态计算

2.5 批流一体数据生产与运维能力

数据回溯是常见的运维场景。流计算模式比批模式使用更多计算资源,因此我们采用 Flink 批任务完成数据回溯。为了避免批任务覆盖更新鲜的数据,我们让业务先停掉流计算任务,再用批任务进行数据覆盖更新。

批流一体

03 业务实践

3.1 案例一:加速数据入仓

采用增量数仓架构前,业务数据入仓需要经过多个步骤,导致数据就绪时间晚。通过 Flink 流处理,结合下游需求,将清洗后的数据有选择地落 Hive 或 Beluga,提升了数据入仓效率。

数据入仓优化

3.2 案例二:提升计算效率

通过增量计算模式,提升了明细快照表和累计快照事实表的计算效率。Beluga 的增量更新能力有效节省了计算资源,提升了数据时效性。

计算效率提升

3.3 案例三:提升开发运维效率

通过 Beluga 替换 Kafka 和 Doris,简化了数据链路,提升了开发运维效率,削减了冗余资源。

开发运维优化

04 未来展望

未来,我们将持续完善 Beluga 的功能和架构,支持批量更新部分列、点查和高效并发控制能力。同时,改进 Beluga 的事务提交效率,支持秒级事务提交,帮助业务低成本迁移离线数仓任务。

未来展望

end