美团增量数仓建设新进展
Apache Flink 中文社区
以下文章来源于 Apache Flink,作者 汤楚熙@美团。

Apache Flink 是 Apache Flink 中文社区唯一官方微信公众号,由 Flink PMC 维护。
摘要: 本文整理自美团系统研发工程师汤楚熙在 Flink Forward Asia 2022 实时湖仓专场的分享。主要内容分为四个部分:
- 建设背景
- 核心能力设计与优化
- 业务实践
- 未来展望
Tips: 点击 阅读原文 免费领取 5000CU*小时 Flink 云资源。
01 美团增量数仓的建设背景
美团数仓架构的诞生基于以下技术假设:“随着业务数据越积越多,增量数据 / 存量数据的比值呈下降趋势,采用增量计算模式性价比更高。” 同时,Flink、Hudi 等具备增量计算、更新能力的技术框架为增量数仓的落地提供了必要条件。

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

目前,美团内部有 M、B、C、D 端等四大类业务场景,不同场景对数据一致性和时效性的要求不同,需要寻找一套尽可能适配所有场景的技术架构。
我们首先想到的是 Lambda 架构,它通过实时链路解决高时效性需求,通过离线链路解决长业务周期的指标计算需求。然而,Lambda 架构的生产链路过于复杂,导致资源成本和运维成本高昂。

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

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

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

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

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,保证了更好的原子性。

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
