哔哩哔哩技术

内容速达 | 哔哩哔哩如何基于Trino+Iceberg打造高效湖仓一体平台

哔哩哔哩如何基于Trino+Iceberg打造高效湖仓一体平台

前言: 哔哩哔哩 OLAP平台负责人-李呈祥在2022年9月24日的Hadoop meetup 2022上海站活动分享了题为《哔哩哔哩如何基于Trino+Iceberg打造高效湖仓一体平台》的主题内容,以下是具体分享。

大家好,我是来自哔哩哔哩的李呈祥,非常高兴来参加Hadoop Meetup,能和大家线下沟通交流,这次我主要分享的内容是哔哩哔哩如何基于Trino+Iceberg打造高效湖仓一体平台。

首先我介绍下哔哩哔哩的数据处理流程,我们的数据主要来自于三个方面:来自APP/网页端等的埋点数据,服务端的日志数据以及业务系统数据库中的数据,经过我们数据采集端的工具和服务,采集到大数据平台中,包括离线的HDFS数据和实时的Kafka数据。

进入大数据平台后,我们的数据开发人员按照数据需求对数据进行加工处理,分层数据建模。主要是用Hive和Spark进行离线的数据处理,使用Flink进行实时的数据处理。用户可以使用Spark或者Trino直接访问数据开发同学建模好的数据表,但是很多时候,为了更好地支持业务在数据探索/BI报表/数据服务/全文检索等方面的实际需求,我们还需要将数据从HDFS上的Hive表导出到外部存储,从而更高效地满足用户需求,比如ClickHouse/Redis/ES等。

当前的这套架构可以基本满足我们对于数据分析的需求,但是也存在许多不足之处,主要问题在于:

第一,从Hive表导出到外部存储,需要额外的数据同步和数据存储成本,整个数据加工的链路也变长了,可靠性降低。

第二,各个存储之间实际上是数据孤岛,跨源查询的成本高,效率低,基本不太可行。

第三,我们内部基于Hadoop/Hive生态构建的数据质量/数据血缘/安全/元数据管理等平台服务工具链很难完整覆盖各个外部存储计算引擎。

所以我们引入了基于Trino+Iceberg的湖仓一体架构,希望能够简化当前的数据处理流程,核心的目标有两个:一是之前存在大量的通过Trino/Spark访问Hive的情况,我们希望通过湖仓一体加速这部分的查询效率,在提升用户体验的同时,降低查询的机器资源成本。

第二是对于之前很多需要同步到外部存储的场景,在对性能要求没有特别高的情况下,可以不用同步到外部存储,直接使用Iceberg表响应,简化业务的数据开发流程,提升开发效率。

最后基于Trino+Iceberg的湖仓一体的架构可以基本上完全兼容我们大数据开发平台的所有工具链,适配的成本非常低。

我们希望能够基于Trino和Iceberg实现高效的OLAP引擎,提供秒级的查询响应能力,支持大部分的交互式分析需求,少部分特殊需求通过ClickHouse/ElasticSearch以及KV类存储满足。业界对于湖仓一体有多个方向的实践,比如解决数据upsert场景,或者实时数据可见性等,B站在这些方向也都有一些探索和落地实践,不过我今天要介绍的重点是我们在湖仓一体中对于仓的查询性能方向的探索。


我们湖仓一体的整体架构如图所示,基于Spark/Flink的ETL任务读写Iceberg表,外部服务通过Trino引擎查询Iceberg表数据。

Magnus是我们自研的Iceberg智能数据管理服务,Spark/Flink每次向Iceberg commit新的文件时会向Magnus通知commit信息,Magnus根据表的commit信息以及相关的policy异步调度Spark任务用于对Iceberg表中的数据组织进行优化,比如小文件合并,数据排序等等。

我们的目标是基于Trino和Iceberg构建秒级响应的湖仓一体平台,这个目标的设定主要基于以下的几个我们观察到的事实:

一是我们主要目标业务场景,比如报表/数据产品,他们的表都是经过我们数据开发同学ETL后的强Schema规范化数据,查询的场景也主要是投影/过滤/关联/聚合这几种基本算子的组合,像是两个大表关联,或者复杂嵌套子查询,当然可以执行,但不是我们主要的目标场景。对于这种SPJA查询,一般来说结果集都是非常小的。

二是我们可以通过对Iceberg和Trino进行增强,支持排序/索引/预计算等OLAP高级特性,使得查询时只扫描SQL逻辑上需要的数据,不需要的数据都Skip掉,同时控制需要扫描的数据量在一定的范围内。

第三,Iceberg的事务支持是的我们可以安全地对数据进行合理的重新组织,这个基础是我们能够通过Magnus服务异步进行数据排序/索引/预计算的基础。

同时我们也并不追求向ClickHouse那样的毫秒级响应的能力,主要是不同于ClickHouse的存算一体架构,Iceberg数据存储在HDFS分布式文件系统上,引入了额外网络和文件系统开销,Iceberg主要是在文件级别进行元数据管理,文件一般在256M左右大小级别,粒度相比ClickHouse更粗。

此外,开放的查询引擎在计算侧相比于其他基于native语言开发的,充分利用向量化能力的OLAP引擎也是有不少差距的。基于可预期的数据扫描量和可控的SPJA计算复杂度,我们可以有一个可预期的查询响应时间,那问题就在于:我们如何基于Trino和Iceberg做到在执行查询时,尽量只访问查询逻辑上需要的数据?

对于典型的多位分析场景SPJA算子,我们针对每个算子类型进行分析。首先是投影,Iceberg表实际的数据存储类型是ORC列存格式,所以查询中投影的字段下推到TableScan层,ORC Reader只会读取需要的字段而不是所有字段,这是一个已经解决的问题。

对于过滤,我们需要考虑不同的过滤类型,找到合适的解决方案。过滤一般可以分为两种过滤条件,等值和范围过滤,过滤字段本身也可以根据字段基数的不同分为高基数字段和低基数字段。

在当前Trino和Iceberg的社区版本中,已经有实现了一些针对过滤条件的Data Skipping相关技术。在引擎侧,Trino通过FiterPushDown相关优化器规则会尽可能把过滤条件下推到最底下的TableScan层。

Iceberg首先支持Partition Prunning,分区的文件存储在分区目录中,过滤条件中包含分区时可以skip不相干的分区目录,同时会在表的metadata中记录每个文件所有字段的minmax值信息,在生成InputSplit时,如果某个字段的过滤条件和该字段的minmax值匹配判断是否需要扫描这个文件。

此外,Iceberg还支持用户定义排序字段,比如示例中,age是常用的过滤字段,用户可以定义age为排序字段,Magnus会拉起异步的Spark任务,将该Iceberg表的数据按照用户的定义将数据按照age字段排序。数据排序的好处是重新调整数据的聚集性,让他们按照age字段聚集,比如对于age=16的查询,排序前需要扫描3个文件,排序后只需要扫描1个文件就行了。

上面的例子解释了数据的排序分布对于索引和查询性能的影响,我们进一步拓展了Iceberg对于数据排序分布的能力。我们首先将Iceberg表数据的排序分布分为两类:文件间的数据组织和文件内的数据组织,两者相互独立,可以分别单独配置。我们主要扩展的是文件间的数据组织,总共支持了Hash/Range/Zorder和HibertCurve四种数据组织方式。Hash和Range大家都比较熟悉,这里主要简单介绍下Zorder和HibertCurve两种Distribution。

如果对于某个Iceberg表,有多个常用的过滤字段,我们使用Order By a,b,c对多个字段进行排序后,数据的聚集性对于a,b,c依次降低,data skipping效果也会依次下降,尤其是a的基数比较高的时候,很可能对只有b或c过滤条件的查询无法skip任何一个文件。

Zorder的做法就是将多个字段值的多维数据依照规则映射成一维数据,我们按照映射成的一维数据排序组织,这个一维数据按照大小顺序连接起来是一个嵌套的Z字型,所以被成为Zorder排序。

Zorder可以保证映射后的一维数据的顺序可以同时保证原始各个维度的聚集性,从而保证对于各个参与Zorder排序的过滤字段,都有比较好的Data Skipping效果。

针对不同数据类型和数据分布,Zorder的实现也是一个比较有挑战的事情,有兴趣的同学可以参考我们之前的一篇文章:https://zhuanlan.zhihu.com/p/354334895。

我们是在文件间排序阶段支持Zorder,所以实际上我们需要这个嵌套Z字型的数据分布切成很多段,每一段对应的数据存储在一个文件中。

可以看到,Zorder有些连接线的跨度比较大,如果跨度大的连接线连结的两个点的数据被切分到了一个文件的话,这个文件在对应字段上的minmax值的范围就会很大,对应字段过滤条件很可能就没法跳过这个文件,导致data skipping概率降低。希伯特曲线和Zorder类似,好处是它不存在跨度很大的连接线,所以是比Zorder更优的一种多维字段排序方式。

这是我们一个具体的测试场景,用了star schema 1TB的数据集,总共1000个文件,可以看到,按照Zorder排序后,针对三个参与Zorder排序字段的等值过滤,都只需要扫描一百多个文件,可以skip掉80%以上的文件,而希伯特曲线排序后,需要扫描的文件数量有进一步的降低。

除了数据排序分布,我们也在Iceberg支持的索引方面进行了增强,支持了多种索引,以应对不同的过滤条件和字段类型。

BitMap索引可以支持范围过滤,并且多个过滤条件的bitmap可以求并,增加skip概率。但是bitmap的主要问题有两个,一个是对于每个基数值都存储一个bitmap代价太大了,二是范围查询时需要读取大量bitmap计算交并差,这大大限制了bitmap索引的应用场景,使用索引可能导致性能的逆优化。我们在这块也有一些探索,感兴趣的话可以参考:https://zhuanlan.zhihu.com/p/433622640。

BloomRangeFilter是我们参考公开论文实现的一个类似BloomFilter但是支持Range过滤场景的索引,有False Positive的可能,但是需要的存储空间相比于BitMap这种精确索引大大减少,在我们实际测试中,一般可以达到我们优化后的Bitmap索引的十分之一大小。

基于Iceberg支持了丰富的索引类型 ,以及通过数据排序分布提升数据聚集性,保证索引的效果,那么Trino是怎么使用Iceberg的索引的呢?

可以分成两个阶段,第一个阶段是Coordinator在获取InputSplit时,这个阶段使用存储在Iceberg表metadata文件中的相关索引信息,比如每个文件各字段的minmax值,skip掉的文件不会生成InputSplit。

第二个阶段是在Trino Worker接收到分配的task,处理Input Split中的数据时,首先根据文件读取文件对应的索引文件数据,判读是否可以skip当前文件。

我们生成的索引文件和数据文件是一一对应的,当索引大小 小于某个阈值时保存在表的metadata中,在阶段一时使用,当大于阈值时,保存在独立文件中,阶段二使用。

在有Join的查询中,如何有效地Skip不需要访问的数据是一个很难解决的问题。对于典型的星型模型场景,影响性能的关键是扫描事实表的数据量,但是过滤条件一般是根据维度表中的维度字段过滤,Trino是没办法使用维度字段的过滤条件去skip事实表的文件的。

我们支持了在Iceberg表上定义虚拟的关联列,关联列相当于把维度表的维度字段打宽到事实表上,当然实际上不会真的存储,只是一个逻辑上的定义,然后用户就可以像对待原始的列一样对待关联列,可以基于关联列定义数据排序组织,可以基于关联列定义索引。关联列要求事实表和维度表满足一定的约束关系,也就是事实表和维度表Join后的结果相当于对事实表的打宽,事实表的行数没有增加也没有减少,称为Record-Preserved Join。

一般满足这种条件的是:事实表Left Join维度表,且维度表的Join key满足Unique Key的约束,或者,两表的join key满足PKFK的约束,那么事实表和维度表LEFT JOIN或者Inner Join都可以保证Record-Preserved Join。

我们能够根据关联列定义数据排序组织和索引后,针对维度字段的过滤条件,通过添加一个Trino的优化器Rule,把符合条件的过滤条件就可以从维度表的TableScan中抽取下推到事实表的TableScan中,利用定义在事实表上该字段的索引数据判断是否可以Skip当前的事实表数据文件,从而使得星型模型的Data Skipping效果可以达到和大款表类似的效果,对于我们支持星型模型的业务场景,是一个非常大的提升。

通过对于数据排序组织,索引和关联列的支持,Trino + Iceberg可以基本上做到文件级的只实际扫描SQL逻辑上需要的数据到引擎中参与计算,但是对于部分包含聚合算子的查询场景,可以SQL逻辑上就需要计算大量的数据,聚合成少量结果集返回给用户,这样的场景主要是需要通过预计算解决性能上的问题,通过预计算的结果直接响应查询,从而避免实际扫描计算大量的数据。

我们当前支持了直接通过Iceberg metadata中的数据直接响应用户表/分区级别的count/min/max聚合查询。对于更通用的预计算方案,还在开发过程中,如何实现高效的文件级别预计算存储和查询,如何利用部分文件预计算结果加速查询,如何解决预计算cube维度爆炸问题等,这是一个非常有意思而且有挑战的方向,我们后面有实际成果的时候到时会和大家在分享在这个方面的工作。


B站的湖仓一体平台目前处在快速发展的阶段,这里和大家分享下当前的一些关键指标,我们的Trino集群大概是5376个core,每天有7万的查询量,总接入的数据量目前是2PB,通过数据排序/索引等广泛的应用,平均查询只需要扫描2GB的数据,总体P90的响应时间在2s以内,基本上达到了我们建设秒级响应的湖仓一体平台的目标。