改版通知

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

基于ApachePaimon的StreamingLakehouse的探索与实践

ckckck2025年1月10日33 浏览

以下文章来源于Apache Paimon,作者石在虎&史亚光

Apache Paimon
Apache Paimon 中文社区官微,PMC 成员维护

摘要
本文主要介绍巴别时代基于 Apache Paimon(Incubating) 构建 Streaming Lakehouse 的生产实践经验。我们基于 Apache Paimon 构建 Streaming Lakehouse 的落地实践主要分为三期:

  • 第一期:在调研验证的基础上进行数仓分层,并上线一些简单的业务验证效果;
  • 第二期:实现流式数仓的基础设施建设,优先替换当前基于 Apache Kafka 构建的实时数仓;
  • 第三期:完善 Paimon 的生态建设,包括数据资产、数据服务等平台服务建设,目标是提供完整的基于 Apache Paimon(Incubating) 端到端的平台服务能力。

目前已完成第一期的数仓分层,并进行数据质量验证,基本满足业务需求。本文将详细解析基于 Paimon 构建 Streaming Lakehouse 的生产实践:

  1. 介绍基于 Apache Kafka 的实时数仓业务痛点及引入 Paimon 的背景;
  2. 剖析基于 Paimon 的 Streaming Lakehouse 架构设计;
  3. 概述 Paimon 的使用姿势和业务生产实践;
  4. 总结基于 Paimon 构建 Streaming Lakehouse 的生产实践遇到的问题及解决方案;
  5. 整理基于 Paimon 提供平台建设能力的规划展望。

01 业务背景

在基于 Apache Kafka 构建的实时数仓过程中,我们遇到了一些痛点,例如中间层数据不可分析、数据保留时间短等问题。我们的实时数仓基于 Flink + Kafka + Redis + ClickHouse 构建,难以查询和分析 Kafka 的中间层数据和 Redis 的维表数据。

目前只有 ADS 层数据最终写入到 ClickHouse 里才能分析,但 ClickHouse 对数据更新的支持不佳,需要通过写入重复数据的方式以达到更新的效果。ClickHouse 去重表执行操作是异步的,业务端需要进行数据去重,增加了业务 SQL 的复杂度和性能损耗。此外,ClickHouse 不支持事务,难以做到 Flink 到 ClickHouse 端到端的数据一致性保障。

基于以上痛点,我们希望借助当下流行的数据湖存储方案简化数仓架构,提高数据分析效率,降低数据存储和开发成本,最终选择 Apache Paimon 作为湖仓底座,主要基于以下几个方面的考量:

  1. 强大的数据更新能力:Apache Paimon (Incubating) 基于 LSM 的数据更新能力,支持基于 PK 的数据更新、Partial Update 的部分更新和 Aggregate 表的预聚合能力,能够大大简化业务开发的复杂度。
  2. 与 Flink 的高度集成:Apache Paimon (Incubating) 作为 Apache Flink 的子项目,对 Flink 集成的成熟度较高,支持所有 Flink SQL 语法,优先支持 Flink 集成。
  3. 流式数仓的演进:Flink Forward Asia 2021 的主题演讲中,Apache Flink 中文社区发起人王峰提出流式数仓的概念,Paimon 作为流批一体的存储,是 Flink 推动流批一体演进中存储领域的重要一环。
  4. 社区支持:在调研测试过程中,社区对问题的快速响应和 Bug 修复,促成了我们最终选择 Apache Paimon。

在此特别感谢之信老师、晓峰老师以及 Paimon 社区的各位开发者的支持。


02 数仓架构

目前我们已完成基于 Paimon 的数仓分层设计,包括 ODS、DWD、DIM 层的搭建以及 DWS 层一些业务模型的建设,整体架构如下:

2.1 数据来源

我们的数据源主要包括前端和后端打点日志以及业务数据库的 Binlog。打点日志按项目通过 Filebeat 采集到 Kafka 对应的 Topic,然后通过 Flink SQL 同步到湖仓的 ODS 层。业务库的数据通过 FlinkCDC 整库同步到湖仓的 ODS 层的 Paimon 表。

2.2 湖仓建设

湖仓主要基于 Apache Paimon 构建,各层通过 Flink SQL 进行数据的准实时同步:

  • ODS 层:采用 Paimon 的 Append Only 表,保留数据原貌不做更新。
  • DIM 层:采用 Paimon 的 PK 表,部分维表使用 Partial Update 能力保留最新的维表数据。
  • DWD 层:采用 Paimon 的 PK 表,ODS 层表数据经由 Flink SQL 做 ETL 清洗,并通过 Retry Lookup Join 关联维表拉宽后写入到 DWD 层对应的 Paimon 表。
  • DWS 层:主要采用 Paimon 的 Agg 表进行预聚合模型及大宽表的建设。
  • ADS 层:将 DWS 层的结果数据和 DWD 层的一些明细表数据流读到 ClickHouse 在线系统,提供在线服务使用。

2.3 在线系统

通过 Flink SQL 将 DWS 层的结果数据和 DWD 层的一些明细表数据近实时地流读到 ClickHouse 在线系统进行 OLAP 分析,提供 BI 实时报表、大屏展示以及用户行为分析系统等使用。同时扩展 Paimon 的 Presto 连接器,数据分析师可以使用 Presto 引擎进行 Adhoc 查询和数据捞取工作。

2.4 平台服务

我们的湖仓目前使用 Paimon 的 Hive Catalog,基于 HMS 做元数据的统一管理。数据开发基于 Dinky 进行二次开发,使用 Dinky 在 Flink SQL 开发方面的能力。数据指标基于不同类型的游戏进行梳理,以便构建统一的指标体系。后续考虑基于 Paimon 构建数据资产以及数据服务。


03 生产实践

在介绍业务生产实践之前,首先介绍一些 Paimon 的正确使用姿势,以便更好理解以下的业务建表实践。

Merge Engine

指定 Merge Engine 的作用是把写到 Paimon 表的多条相同 PK 的数据合并为一条,用户可以通过 merge-engine 配置项选择以何种方式合并同 PK 的数据。Paimon 支持的 Merge Engine 包括:

  • deduplicate:默认的 Merge Engine,只保留最新的记录,其他同 PK 数据则被丢弃。
  • partial-update:支持多个 Flink 流任务更新同一张表的部分列,最终实现一行完整数据的更新。
  • aggregation:支持通过聚合函数进行预聚合,每个除主键以外的列都可以指定一个聚合函数。

Changelog Producer

Changelog 主要应用在流读场景,Paimon 支持的 Changelog Producer 包括:

  • none:默认不生成 Changelog,成本较高,不建议使用。
  • input:依赖输入端的完整 Changelog,适用于业务库的 Binlog 等场景。
  • lookup:通过 Lookup 方式在数据写入时生成 Changelog,目前处于实验状态。
  • full-compaction:在 Compaction 后生成完整的 Changelog,适用于流读场景。

Append Only Table

建表时配置 write-mode = 'append-only',用户可以创建 Append Only 表。Append Only 表采用追加写的方式,只能插入完整记录,不能更新和删除,适用于无需更新的场景,如 ODS 层数据。


04 实践总结

  1. 资源压力:使用 full-compaction Changelog Producer 时,changelog-producer.compaction-intervalcheckpoint interval 设置值较小时,Writer 端在写入数据和 Compaction 时的压力较大,需不断调整资源。
  2. 小文件问题:当数仓某一层的数据出现问题,需要通过 Time Travel 重新读取某个快照或时间点的数据修复问题,Snapshot 保留时间越长,生成的小文件越多。
  3. Retry Lookup Join 效率低:目前的 Retry Lookup Join 是有序的,前面一条数据一直 Join 不上时,后面的数据会排队,导致处理效率低。

05 未来规划

  • 完善基于 Apache Paimon(Incubating) 的流式数仓建设;
  • 优化 Presto 查询,提升基于 Paimon 的即席查询能力;
  • 考虑直接在 Paimon 里预聚合结果数据,减少数据处理链路和复杂度;
  • 完善基于 Apache Paimon(Incubating) 的平台服务建设。


end