爱奇艺数据湖实战
01
什么是数据湖?
公有云数据湖
非公有云数据湖
爱奇艺数据湖
统一存储:支持灵活的存储底座(公有云/私有云、HDD/SSD/缓存),具备集中的、足够大的存储空间
通用数据抽象/组织层:支持结构化、半结构化、非结构化等不同数据类型,并抽象出一层统一的数据组织形式,Hudi、Iceberg、Delta Lake等目前仅统一了结构化数据格式
支持批处理、流计算、机器学习等不同计算类型
统一的数据管理(元数据中心、生命周期、数据治理等),避免形成数据孤岛
02
为什么需要数据湖?
时效性好:数据湖里的数据近实时( 1-5 分钟)可见,比 Hive 离线延时优势明显,能满足大部分业务的要求 规模大:数据湖存储使用的是 HDFS,写入吞吐几乎无瓶颈 成本低:数据湖无需单独服务器,机器成本低,运维成本低 查询支持:数据湖支持明细数据查询,也支持各种复杂的联合分析 数据分享:数据湖支持 Spark、Trino、Flink 等各类引擎分析,而数据进入 ElasticSearch、Kudu、Druid、ClickHouse 后都只能使用专用的查询引擎分析,数据分享需额外导出为 Parquet 等格式
场景二 - 变更数据分析
离线导出:缺点是同步延迟大,通常是 T+1;另一方面是同步代价大,每次需全量导出,给数据源 MySQL/HBase 带来较大读压力;
实时同步进 Kudu:面临规模小、成本大等痛点
场景三 - 流批一体
离线通路时效性差、而实时通路容量低
维护两套逻辑,开发效率低
容易出现数据不一致
维护两套服务机器,成本高
支持海量数据近实时更新
一套代码,避免重复开发
避免数据不一致
服务存储和计算成本更低
满足上述场景产品特性总结
规模大,成本低:能支持PB级别数据规模
支持更新:包括历史分区新增数据、行级更新等
增量拉取:将表的变更转成流数据用于构建下游表
时效性:近实时(5分钟)
查询快:交互级查询速度
03
数据湖选型及原理说明
Iceberg表格式
表格式定义
Iceberg 不是存储引擎:其支持 HDFS、S3 对象存储作为底层存储引擎 Iceberg 不是文件格式:使用 Parquet 存储数据文件 Iceberg 不是查询引擎:可通过 Spark/Flink/Trino/Hive 查询
Hive 表格式
MySQL Metastore 存储元数据,包括库、表和分区信息,不包含文件信息,最小的原子操作是分区级替换
用目录树组织数据文件,通过 LIST 目录接口获取分区下的数据文件列表,可实现分区级过滤
制定执行计划慢:假设一个表以小时分区且每小时有 100 个子分区,则 7 天范围共 16.8K 个分区,一个简单地扫描任务需执行 O(N)次 NameNode RPC 调用,N=扫描分区数,假设一次 RPC 调用需 2ms,制定执行计划需耗时 33.6 秒;
无法应用文件级过滤:例如表存储的是广告点击记录,且写入时按照广告主 ID 排序,此时我们查询特定广告主 ID 的记录,每个分区下仅命中少量文件,但 Hive 并无相关信息用于过滤掉其他文件;
不支持修改:分区任何修改都需执行分区级覆盖,如历史分区新增一列;不支持行级修改;
不支持增量:若有一个任务消费表 A 更新表 B,假设表 A 新增了迟到的增量文件,无法获取表 A 的增量更新部分触发计算更新表 B。业务要么选择丢弃增量部分不往下游传递,要么对整个分区进行重算;
依赖文件系统重命名:对于对象存储不友好;
Iceberg表格式
读写互不干扰:读写可操作不同的快照,写在提交前不可见; 支持并行写入:采用乐观锁的机制,写的过程不加锁,提交前检查是否冲突,无冲突则提交成功,包含冲突内容则放弃提交稍后重试;
执行计划快:如前文所述 Hive 制定执行计划耗时和查询涉及的分区数正相关,而 Iceberg 直接读取元数据文件即可获得文件列表,制定执行计划耗时大幅缩短; 文件过滤加速执行:Iceberg 记录了文件的统计信息,不同的执行引擎可基于统计信息(MinMax 值、字典、布隆过滤器等)过滤掉无关的文件,大幅减少实际读取的文件数加速执行;
新增数据:Iceberg 支持往已有的表/分区中添加少量文件,无需分区级覆盖; 获取增量:Iceberg 支持获取 2 次快照间的文件变化,并支持流式地读取变更,从而实现增量更新下游表;
新定义 DeleteFile:格式上仍然是 DataFile,记录本次提交删除的行; Merge on Read:读取时将 DataFile 和 DeleteFile 内容合并,得到准确的结果
图3-4 Iceberg 行级更新一个例子
技术小结 - Iceberg如何实现设计目标
04
数据湖业务落地
Venus 日志采集平台
业务痛点
大部分业务配置的是 0 副本:因 ES 写入成本高,所有业务配置 1 副本写入成本需翻倍;0 副本导致任意硬盘/结点/集群故障都会影响部分业务写入;
业务隔离:给高优业务以独立 ElasticSearch 集群,低优业务共用公共集群,避免低优业务流量增长影响高优业务,但无法解决高优业务自身增长的问题;
流量调度:单个集群流量到瓶颈时,将部分流量调度到其他空闲集群;单个集群故障时,将业务流量调度到其他 ElasticSearch 集群;
写入失败多:业务排查时经常遇到日志延迟半小时以上,甚至写入失败,日志丢失等情况
排障压力大:由于 0 Replica 很容易导致写入失败,每天需处理10+的运维请求;
成本高:ElasticSearch 设计上是牺牲写入时性能以换取查询性能,而日志类特点是写入 QPS 大,查询 QPS 低,Venus 机器经常磁盘达到瓶颈,而 CPU 和内存大量浪费; 新架构
数据延迟低:日志采集到查询需要分钟级的延时;
查询速度快:交互式排障需要查询在秒级返回;
写入带宽高:峰值 QPS 千万/秒,总数据量在 PB 级;
落地效果
成本优化:Iceberg 存储复用的 HDFS,查询所有业务共用一个 Trino 集群,无需部署独立的集群,节省大量机器成本;
写入稳定:由于 Iceberg 存储是 HDFS 3 副本,单个硬盘/结点故障不影响写入,且 Iceberg 写入带宽近乎无限,几乎不再发生达到写入瓶颈、存储容量不足、日志丢失的情况;
排障减少:Venus 团队统计入湖后运维量降低 80%,节省一个运维人力;
审核数据
业务痛点
MongoDB:存储全量审核数据,规模在百亿行,仅对 ID 构建索引,无法对其他列开启索引;
ElasticSearch:存储用于检索的列,因数据量限制不存储原始消息;线上服务查询某个关键字的记录时,先通过 ElasticSearch 服务筛选命中的 ID 列表,再对 ID 列表逐一查询 MongoDB 获得原始记录;
MySQL:针对一些报表需求,通过定时任务查询 ES 并将聚合结果存储在 MySQL;
Hive:业务原计划将 MongoDB 全量导出为 Hive 用于离线分析;
开发成本高:每新增一个报表需求,需开发一个 ES 定时查询任务,将结果记录为一个 MySQL 表,并在报表页面进行适配,无法满足快速变化的分析需求;
ES 查询瓶颈:当定时查询任务较多时,给 ES 服务造成较大的压力,影响线上通路性能和稳定性;
数据质量:当历史数据发生变更,如曾经审核通过的记录当前审核不通过,并不会更新已算好的统计值(如审核通过率),从而报表数据质量会逐渐下降;
存储容量:ES 容量有限,当前 MongoDB 诸多大表不在 ES 通路;
Hive 通路:业务初步调研后发现不可行,一方面全量导出耗时很久,执行一天仍未完成,另一方面导出过程给 MongoDB 造成较大的压力; 新架构
行级更新:审核的记录会一直变化,如审核状态、修改时间等;
高效查询:支持基于不同列的高效过滤分析,支持和其他表联合分析;
容量大:支持百亿量级,且未来还有更多场景接入;
开发成本低:撰写 SQL 即可;
查询可扩展:SparkSQL 算力可水平扩展,且不影响线上通路;
数据质量高:行级更新保证 Iceberg 数据和 MongoDB 完全一致;
存储容量大:Iceberg 存储是 HDFS,支持 PB 以上的规模;
时效性好:数据延迟在 5 分钟,近实时地反映数据变更; 落地效果
审核团队在 Iceberg 表落地后赋予业务了一系列新的可能性,审核团队基于 Iceberg 表拓展了一系列从无到有的场景,其中部分场景如下:
数据统计:审核团队人效统计、风险监控实时报警;
基于关键字下线:原先需对 ElasticSearch 表全量扫描,影响线上稳定性,现在批量扫描 Iceberg 表即可;
导出数据:由于 MongoDB 无法做 Group By 分析,需导出到 CSV 后再用 Shell 脚本处理,需十几个小时;当前大幅降低工作量,执行 SQL 语句即可,耗时缩短到 5 小时;
降低风险:对数由原先 16 小时缩短到 5 小时,降低漏审/误审带来的内容安全风险。
Pingback流批一体
业务痛点
离线通路有小时以上的延时,无法对最新数据做分析;
实时通路不支持数据明细查询,全量分析;
为了同时支持全量分析和低延时数据可见性,构建 Lambda 架构,而同时维护两个开发链路导致开发维护成本高、实时离线数据不一致等问题;
新架构
生产 ODS 层表:使用 Flink 增量消费 Kafka 中的全量数据,解析并按投递规则拆分生成 ODS 层表,一个 Flink 任务会拆分生成数百张 Iceberg 表;
生产 DWD 层表:通过 Flink 增量消费 Iceberg 表,进行维度扩展、标准化等加工生成 DWD 层表
下游 Pipeline:下游业务可通过 RCP 实时计算平台、Babel 离线计算平台继续构建 Pipeline,也可在 RAP 实时分析平台、魔镜离线分析平台进行查询分析
相比离线通路:数据可见性延时降低到5分钟以内,并且支持增量读取、版本回退等新特性;
相比实时通路:如果可接受5分钟的延时,具备实时通路增量读取的特性,并能支持全量读取、明细数据查询,具有更好的容错性;
相比Lambda架构,能做到存储计算的流批一体,避免开发维护两套代码及实时离线数据不一致的问题;
成本收益:预期近实时通路成本和离线通路接近,同时节省大量实时通路资源; 落地效果
播放 Pingback;已生产播放 Pingback 峰值 QPS 百万级的数据,并使用增量读取数据湖的方式构建了爱奇艺的点播、直播报表,数据与已有离线数据一致,延时在 1 分钟左右,相比实时通路成本下降 90%
QOS Pingback:QOS Pingback 是监控 APP 运行状态的埋点信息,用于监控和排障。相比通过离线明细数据进行故障定位,使用近实时通路,在发现问题后,可立即查询明细数据定位故障,将大幅缩短故障定位时间。当前已稳定生产了 QOS Pingback 的 600 多张表,正在推动业务迁移到近实时通路。
会员订单
业务痛点
数据时延大:当前导出是天级,业务只能分析一天前的数据;
MySQL 压力大:每天全量导出数据量非常大,容易打满 MySQL CPU;
Kudu 压力大:订单表消耗了 Kudu TB 级的写内存,一方面机器成本高,经常需运维集群,另一方面未来扩展性差,难以承接其他 MySQL 场景;
写任务运维:Kudu 集群写入性能有波动,会造成消费 CDC 变更流写 Kudu 任务堆积,需运维处理;
Spark 任务失败:风控业务定期扫描分析 Kudu 表,一旦 Kudu 表 Tablet 有迁移会造成任务失败;
新架构
延时低:近实时延迟,低至5分钟/1分钟;
查询快:通过 SparkSQL 查询,结合文件合并等优化,性能和 Kudu 方案接近;
成本低:Iceberg 无需单独集群,机器成本非常低;
运维低:不会给 MySQL 造成巨大压力,无需特殊运维;
落地效果
05
总结及规划
引用
1、Dixon, James (14 October 2010). "Pentaho, Hadoop, and Data Lakes". James Dixon’s Blog.
2、AWS. What is a data lake
3、Google Cloud. What is a data lake
4、数据湖 | 一文读懂Data Lake的概念、特征、架构与案例]
5、Uber’s case for incremental processing on Hadoop
6、Iceberg: A modern table format for big data
7、Apache Iceberg: An Architectural Look Under the Covers
8、Iceberg Table Spec
也许你还想看