B 站基于 Iceberg 的流批一体的探索和实践
导读
大家好,我是来自哔哩哔哩的张陈毅,今天给大家分享的 topic 是B 站基于 Iceberg 的流批一体的探索和实践。
本次的分享主要分为五个部分:
1.
海量用户行为数据传输
2.
商业和 AI 的在线训练
3.
DB 数据同步
4.
Iceberg 维表 Join
5.
Q&A
分享嘉宾|
张陈毅
哔哩哔哩
资深开发工程师
编辑整理|
杨维旭
内容校对|李瑶
出品社区|
DataFun
01
海量用户行为数据传输
1.
实时数据传输概览
分享嘉宾 INTRODUCTION
张陈毅
哔哩哔哩
资深开发工程师
专注于 Flink SQL/State 和流批一体的工作,为内部提供 Flink 引擎相关的技术支撑
。
- 实时链路的数据新鲜度高,在毫秒级;而离线链路的数据新鲜度低,一般由 Hive 的分区按天或小时而决定。
- 存储成本上,实时链路的 Kafka 一般会选择高成本存储介质 SSD 或 NVMe;而离线链路的 Hive 存储在 HDFS 上,一般会选择低成本存储介质 HDD。
- 实时链路的 Kafka 一般仅会保留 1~3 天内的数据,而离线链路的 Hive 可以按需保留一年或更久的时间。
- 在 Batch 的性能上,实时链路本身的特性就是,读数据需要将整个 row data 反序列化后,获取指定列的结果,在性能上比使用列式存储的 Hive 会弱一些。
- Kafka 需要指定 timestamp 或者是 offset 的方式去做 Batch 的读,在易用性方面也比较弱。
- 流计算任务目前设置 checkpoint 默认为 5 分钟,频繁的写入,会导致生成大量的小文件,这在一定程度上会影响下游查询性能。小文件合并主要是为了减少小文件数量,通过将多个小文件合并成较大的文件,从而提升查询的效率。大文件在读取数据时可以减少文件系统的开销,改善磁盘的 IO 性能,同时减少元数据的管理负担。
- Iceberg 支持对数据进行排序,用于提高过滤和聚合操作的效率。通过对数据进行合理的排序,可以优化数据读取的顺序,减少数据随机访问的次数,从而提升整体的查询效率。尤其是在处理范围查询的时候,排序可以大幅减少所需扫描的数据量。
- 在 Iceberg 中可以选择合适的分区和数据分布,以确保数据在各个节点中的数据均匀分布,从而实现更快的并行处理。优化后可以减少数据倾斜现象,提高查询的并行效率。
- 通过创建索引,使下游查询效率得到提高,避免了全表扫描。
- 通过创建预计算文件,提前计算和存储预计算的结果,使后续查询速度得到了提升。
- 首先,训练数据存储不统一,实时训练特征计算结果会存在 Kafka,而离线训练结果存储在 Hive 中。
- 其次,使用实时训练和离线训练两套 API 来开发,存在计算层不统一的现象。
- 第三,为了做离线训练,需要将实时流训练样本 Kafka 数据,通过 Flink SQL Dump 到 Hive 表中,而 Dump 任务的 CPU 消耗一般是原特征计算任务的 1/3,原因是整个链路中数据流量本身比较大,为了节省数据传输的网络带宽消耗,写入 Kafka 的数据在多层中都使用了 PB 格式,在流量大的情况下,下游 Dump 任务反序列化 PB 数据以及 Dump 写到 Hive 表中均需要比较多的计算资源。
- 除了发布到线上的训练流程外,算法人员也会启动多个实验链路来优化整体的训练模型,这样的结果会发现展点流 base Kafka 的网络带宽以及磁盘 IO 均是写 IO 的 20 多倍,即使扩展了 Kafka 的分区数,对此 topic 所在 Kafka 集群网络读 IO 也是比较大的。
- Data File 主要是记录具体存储的数据。
- Position Delete Files 主要是记录文件路径以及所需要删除文件的 offset 信息,一般是记录相同 snapshot 下同一批 commit 的 data file 的行级删除数据。
- Equality Delete Files 会记录相同 snapshot 下未写入 data file 需要删除的数据,其 schema 为主键,value 为需要删除的具体值。
- 即时查询, Iceberg 可以作为 MySQL 的镜像表来做实时查询,达到分钟级延迟。
- 在实时消费的场景下,我们希望 Iceberg 具备增量流读 change log 的能力。
- 离线场景下,希望 Iceberg 每个 tag 是一个全量的历史快照数据。
- 首先,tag 想要将数据按照时间整点去做切割的话,需要借助于 Flink checkpoint 的触发,去关联 watermark,也就是 watermark 到了整点的时间点去触发对应的 checkpoint,然后在 snapshot 上面去打一个对应的 tag。
- 其次,为了避免数据漂移问题,比如 T2 的数据可能会落入到 T1 中,为了解决这个问题,我们在 Writer 算子前面增加一个 cache 算子,用来缓存整点时间尚未到达时提前进入 Writer 算子的数据。
- 最后,需要将 tag 打出来的数据独立于 Data File,以避免非分区表 data 目录下面文件数量达到上限的问题。
- 首先,我们内部提供的维表连接器数量是比较多的,每个连接器的私有参数也很多,导致用户的整体理解和使用成本较高,其中有部分维表属于 MySQL 维表中 DB 无法承受那种流计算高 QPS 压力而出现的替代品,如 Redis,其存储成本也会有所提升。
- 另外,双流 join 任务中会挂一个维表 join,其任务并行度比较高,致使 TaskManager 数量比较多,导致其 connector 链接数增加。
- 最后一点是,即使在本地每个 TaskManager 做了一些 cache 的优化,比如针对 MySQL,基于内存做了优化,对于 Hive 和 HDFS 维表我们在本地做了 RockDB 的缓存。但是这些缓存都是无状态的,当任务重启时,所有维表数据都需要从远端重新拉起,整体来看维表使用 DB 压力会比较大。
- 第一,维表的全量数据需要存储在 State 当中,这样任务即使重启,也无需对全量数据做重新加载。
- 第二,State 需要具备无 TTL 过期的属性,也就是数据不会过期。
- 第三,维表需要保有全量数据外,还需要具有增量流读能力,来更新新 join 算子中的 State 数据。
- 第四,增量流读的数据同时也需要具备 change log 的特性,以 MySQL 表做 join 为例,用户作维表 join 的 on 表达式使用的字段,一般不会和主键字段一致。如果没有 change log 特性,一条主键 id=1 的数据在多次更新后,可能会在 join 算子的多个 subtask 中同时存在,最终会影响 join 的结果。
- 第五,维表需要有主键。这是在 keyed State 当中,根据主键去做 State 更新的关键。
- 第一,为了支持流转批场景,使流计算 SQL 能够在基本不改变 SQL 写法的情况下,自动切换至在 Flink Batch 中运行,我们借助了 Calcite 的使用 SQLHint 的方式来支持新 join 语法,降低 SQL 的撰写难度和区分存量 join,只需要在 a join b 后面加一个 SQLHint 和对应的 options 即可。
- 第二,在用户做维表 join 时,我们不希望强制用户使用主键字段去做 join,否则易用性会比较差。而这样带来的一个问题是 Flink SQL 的 Rule 集合优化器,会因为整个 DAG 当中没有使用主键字段,而进行裁剪移除,最后会导致 Flink join 算子中无法获取主键字段,这样也就没办法保证在 State 中根据主键作维表数据更新的诉求。所以我们对 Calcite 和 Flink Project Rule 做了一些调整,支持保留主键字段,至少在 LookupJoin2 Operator 这个算子是有所保留的,后面的就可以不需要了。
- 第三,我们在使用 Flink1.15 的时候,会发现 change log 会有一些优化,Flink 会提供一个 DropUpdateBefore 算子,用以剔除-U 的数据。在当前 join 的场景当中,我们将维表 source 算子后的此算子做了移除,使得 State 数据更新得以准确。
- 最后,新的 Join Operator 在右流初始化全量数据完成前,会将左流数据阻塞。因为对于用户来说,任务启动后,Sink 端如果接收到了数据,就表明已经做完了前面全部的流计算逻辑了。如果说主流在 Right 维表没有做完初始化前不做阻塞的操作,主流 Join 的部分数据可能会存在 Join 结果为 Null 的情况。我们通过实践发现,维表数据即使在亿级别做全量加载,也可以在一分钟左右完成,因为它本身是并行分布式的。
分享嘉宾 INTRODUCTION