阿里巴巴瓴羊基于Flink实时计算的优化和实践
以下文章来源于 Apache Flink,作者王柳焮@阿里云

Apache Flink:Apache Flink 中文社区唯一官微,由 Flink PMC 维护。
摘要:本文整理自阿里云智能集团技术专家王柳焮在 Flink Forward Asia 2023 中平台建设专场的分享。内容主要为以下四部分:
- 阿里巴巴瓴羊基于 Flink 实时计算的平台演进
- Flink 能力优化与建设
- 基于 Flink 的最佳实践
- 未来规划
Tips:点击 阅读原文 在线观看 FFA 2023 会后资料~
01 阿里巴巴瓴羊基于 Flink 实时计算的平台演进
1.1 关于瓴羊

瓴羊是阿里云智能集团的重要业务,致力于将阿里巴巴沉淀十余年的数字化服务经验,系统化、产品化地全面对外输出给千行百业。2012 年开始,阿里提出数据建设方法论,内部形成多款标杆数据产品,如生意参谋、双十一数据大屏等。2015 年阿里巴巴启动中台战略,进行全域数据建设,强化数据能力在业务端的价值显现,让数字化能力广泛服务于各个业务。2018 年应商家日益增长的需求,阿里巴巴集团中台能力进行对外输出,推出 Dataphin、QuickBI 等产品,帮助企业通过数字化技术驱动创新和增长,在生态内外产生显著价值。2021 年瓴羊成立,成为阿里巴巴动物园中的一员。
1.2 Dataphin 平台实时业务规模

在集团内部,平台约承载有 15000+ 的实时计算任务,1000+ 的流批一体任务,覆盖约 50+ 个 BU。在云上,客户分布在电商、金融、交通、零售、制造等行业,应用场景覆盖诸如实时 ETL、实时大屏、实时集成等实时主要场景,以及像金融行业的特征计算以及风控场景等。
1.3 Dataphin 平台实时架构大图

Dataphin 平台的实时架构大图如下:
- 部署形态:支持多云输出,提供公共云、私有云、专有云、混合云等多种输出形态。
- 调度系统:支持基于不同调度系统的作业运行能力。
- Flink 发行版:提供多版本多集群管理能力,以及多数据源、计算源的元数据管理体系。
- 上层模块:
- 数据集成:提供无代码化全增量一体的数据集成能力。
- 研发中心:提供编译优化、任务模版、流批一体等能力的实时研发体验。
- 运维中心:提供任务全方位的监控告警体系,保障线上任务的稳定性。
- 资产:支持表、字段、函数、血缘等数据资产的盘点与展示、标准定义与管理、质量评估及保障、分类分级与脱敏等能力。
02 Flink 能力优化与建设

2.1 CDC 能力提升

平台除了支持社区 Flink CDC 的能力外,还提供了自研的 Flink Connector 去满足更多样化的数据同步场景需求。平台目前已经支持的提供增强 CDC 去做数据同步的部分输入输出,更多的数据源还在陆续增加中。
- 自动化感知来源库/表变更的能力
- 支持配置多样的规则满足不同数据源的入湖入仓场景
- 提供无代码化的整库实时同步能力
- 凭借元数据管理体系支撑丰富多样的源端到源端的数据流通能力
2.2 元数据管理

面对多达几十上百种不同的数据源,用户如何有效应对?如何有效管理不同的数据源连接信息,做好元数据管理显得尤为重要。平台构建了一套元数据的管理体系,分别是 Source 层、Meta 层和 Job 层。
- Source 层:按照物理数据源进行管理,与引擎解耦,管理基础的连接信息。
- Meta 层:基于数据源的构建,可以是按照 Connector 类型区分的 Flink DDL,也可以是来自数据源或者计算源中的物理表。
- Job 层:可以重复引用 Meta 层的表,避免反复构建 DDL,使得 DDL 的变更可以被所有引用的作业所感知。
2.3 流批一体建设

平台的解法是提供流批一体化的开发模式。面向开发人员,只需维护一套代码,由平台根据流批不同的模式翻译为对应的流批可执行 SQL。在存储层面提供面向流批一体逻辑镜像表,无论底层是统一存储或流批不同存储,面向 SQL 侧看来是一张表,这张表可对外提供一致的数据服务。
2.4 CDC 运维体系

整个运维体系支持多调度模式、多集群环境、多引擎版本。平台通过 Metric Reporter 和 Log Agent 分别采集任务的运行指标和日志,指标会落到时间序列数据库中进行持久化存储,日志会落到文件存储中,平台提供统一的监控看板供用户查看。
03 基于 Flink 的最佳实践
3.1 特征计算

特征计算的架构经历了从明细数据方案到预计算方案,再到全计算方案的演进。预计算方案通过 Flink 依赖算子拆分日、时、分维度计算并存储于 Hbase,配置不同滑窗的特征批量生成,基于 Hbase 查询及内存计算,计算压力大大减少。全计算方案则直接服务端进行划窗聚合,写入 Hbase,集群的存储数据量更低,服务端的聚合查询性能更好。
3.2 湖仓一体

湖仓一体的架构体系基于 Flink 或 DataX 的集成方案写入数据湖中,提供 OpenAPI 的方式将数据写入湖的底层存储中。湖的上层提供诸如 Hive、Spark、Presto 等的 Adhoc 即席查询能力,湖中的数据可通过数据服务对外提供服务。
04 未来规划

- 行业解决方案:优先将湖仓一体的场景深化落地,在产品上做成完整的解决方案面向客户。
- 平台功能完备:继续丰富 Flink CDC 支持的数据源,建设全域元数据中心。
- 持续进行体验优化:帮助用户定位数据流向位置,方便纠错排查;基于典型异常场景,通过智能诊断,做到预警。
end
