Delta Lake是构建在数据湖之上的开源存储层,提供ACID事务、Schema管理等能力。但是用不好的话,性能可能比直接用Parquet还慢。本文是我在实际项目中对Delta Lake进行性能优化的实战经验,包括小文件合并、数据分区、Z-Order优化、数据跳过、Vacuum清理等多个方面。每一项优化都有具体的操作步骤和效果对比,希望能帮到正在使用Delta Lake的同学。

一、背景介绍

先说说我们的项目背景。

我们团队用Delta Lake做数据湖已经有一段时间了。数据量大概在几十TB级别,每天有增量数据写入,也有大量的查询需求。刚开始用的时候,觉得Delta Lake确实好用,ACID事务、时间旅行、Schema演进这些特性都很实用。

但是用了一段时间之后,性能问题开始显现。查询越来越慢,尤其是一些复杂的查询,跑几十分钟甚至几个小时都出不来结果。写入也越来越慢,每天的增量数据写入要花好几个小时。

我们排查了一下,发现主要有几个问题:小文件太多、分区不合理、没有利用数据跳过、过期版本没有清理。针对这些问题,我们做了一系列优化,效果很明显,查询速度提升了好几倍,写入速度也提升了不少。

本文就把这些优化经验分享出来,每一项都有具体的操作和效果对比。如果你也在用Delta Lake,或者准备用,希望这篇文章能帮你少走一些弯路。

二、小文件合并

第一个优化是小文件合并,这也是最基础、效果最明显的一项优化。

问题分析

Delta Lake在写入数据的时候,默认会生成很多小文件。尤其是流式写入或者频繁小批量写入的场景,每个批次都会生成一批新文件。时间长了,一个分区下面可能有成千上万个小文件。

小文件多了之后,查询的时候Spark需要打开大量的文件,每个文件都有打开和读取的开销,性能会非常差。而且文件多了之后,元数据管理也会变慢,Delta Lake的事务日志会变得很大。

我们的情况就是这样,有些分区下面有几万个小文件,查询的时候光列文件就要花好几分钟,更别说读取数据了。

优化方法

Delta Lake提供了OPTIMIZE命令,可以自动合并小文件。用法很简单:

OPTIMIZE delta.`/path/to/table`

这个命令会把小文件合并成更大的文件,默认目标文件大小是1GB。你也可以通过配置调整目标文件大小:

SET spark.databricks.delta.optimize.maxFileSize = 1073741824;

如果表有分区,还可以指定分区进行优化:

OPTIMIZE delta.`/path/to/table` WHERE date >= '2021-01-01'

我们的做法是每天凌晨跑一个定时任务,对前一天写入的分区执行OPTIMIZE。这样既能控制小文件数量,又不会因为优化全表而花费太多时间。

效果对比

优化之前,我们一个查询平均要跑15分钟。优化之后,同样的查询平均只需要3分钟,提升了5倍。而且随着小文件数量的减少,集群的压力也小了很多,资源利用率提高了。

注意事项

OPTIMIZE虽然好用,但是也有一些需要注意的地方:

  • 不要太频繁地执行,否则会产生很多不必要的重写
  • 优化的时候会占用集群资源,最好在业务低峰期执行
  • 优化之后旧的文件还在,需要用VACUUM清理

三、合理分区

第二个优化是合理分区。

问题分析

分区是大数据中常用的优化手段,但是分区不合理反而会影响性能。我们刚开始的时候,按照日期分区,每天一个分区。这个看起来没问题,但是我们的数据量每天只有几个GB,每个分区下面只有很少的数据,导致分区数很多,但是每个分区的数据量很小。

查询的时候,如果查询范围跨越很多分区,Spark需要扫描很多分区,每个分区的数据量又很小,开销很大。而且分区多了之后,元数据管理也会变慢。

优化方法

我们重新设计了分区策略。对于数据量小的表,改为按月分区,减少分区数量。对于数据量大的表,保持按日分区,但是增加了二级分区,比如按地区或者业务线分区,让每个分区的数据量更均匀。

另外,我们还注意了分区列的选择。分区列应该是查询中经常用作过滤条件的列,而且基数不能太高也不能太低。基数太高的话分区太多,基数太低的话每个分区数据量太大。一般来说,分区列的基数在几百到几千比较合适。

我们还调整了分区的目录结构。以前是date=2021-01-01这种Hive风格的分区,现在改成了更紧凑的格式,减少了目录层级。

效果对比

重新分区之后,查询的分区扫描数量减少了很多。以前一个查询可能要扫描几百个分区,现在只需要扫描几十个。查询速度又提升了大概30%。

注意事项

分区调整是一个比较大的改动,需要重新写数据。我们是在业务低峰期,用INSERT OVERWRITE的方式重新写入了全量数据。这个过程花了一些时间,但是效果值得。

四、Z-Order优化

第三个优化是Z-Order,这是Delta Lake提供的一个高级优化特性。

问题分析

即使做了小文件合并和合理分区,查询的时候还是要扫描很多数据。因为数据在文件中的存储顺序是按照写入顺序来的,和查询的过滤条件没有关系。比如查询的时候经常按user_id过滤,但是数据在文件中是按时间顺序存储的,那么每个文件中都可能包含符合条件的数据,Spark需要扫描所有文件。

优化方法

Z-Order可以把相关的列在空间上放在一起,这样查询的时候可以利用数据跳过,只读取包含目标数据的文件。用法也很简单:

OPTIMIZE delta.`/path/to/table` ZORDER BY (user_id, product_id)

这个命令会按照指定的列进行Z-Order排序,让相关的数据在物理上聚集在一起。这样当查询用这些列做过滤条件的时候,Delta Lake可以通过统计信息判断哪些文件包含目标数据,跳过不包含的文件。

Z-Order的原理是Z-order曲线,也叫莫顿曲线。它可以把多维数据映射到一维,同时保持多维空间的局部性。简单来说,就是在多个维度上相近的数据,在一维存储中也相近。

我们选择了查询中最常用的几个过滤列作为Z-Order列,包括userid和productid。

效果对比

做了Z-Order之后,按userid和productid过滤的查询,数据跳过率达到了80%以上,也就是说只需要扫描20%的数据。查询速度又提升了2到3倍。

注意事项

Z-Order也有一些限制:

  • Z-Order列不要太多,一般2到4个比较合适,太多的话效果会下降
  • Z-Order需要和OPTIMIZE一起执行,每次OPTIMIZE的时候都会重新排序
  • Z-Order对高基数列效果好,对低基数列效果不明显
  • Z-Order会增加OPTIMIZE的时间,因为需要排序

五、数据跳过和统计信息

第四个优化是充分利用数据跳过和统计信息。

问题分析

Delta Lake在写入数据的时候,会自动收集每个文件的统计信息,包括每列的最小值、最大值、空值数量等。查询的时候,Delta Lake可以利用这些统计信息判断哪些文件可能包含目标数据,跳过不包含的文件。

但是我们发现,很多查询并没有利用数据跳过。原因有几个:一是统计信息没有正确收集,二是查询的写法导致无法利用统计信息,三是没有开启相关的配置。

优化方法

首先,确保统计信息被正确收集。Delta Lake默认会收集统计信息,但是如果数据是通过外部方式写入的,可能没有统计信息。可以用下面的命令重新收集:

ANALYZE TABLE delta.`/path/to/table` COMPUTE STATISTICS FOR COLUMNS user_id, product_id

其次,优化查询的写法。尽量在WHERE子句中使用分区列和Z-Order列作为过滤条件,而且避免在过滤列上使用函数,因为函数会导致无法利用统计信息。比如:

-- 不好的写法,无法利用统计信息
SELECT * FROM table WHERE date_format(created_at, 'yyyy-MM') = '2021-01'

-- 好的写法,可以利用统计信息
SELECT * FROM table WHERE created_at >= '2021-01-01' AND created_at < '2021-02-01'

第三,开启数据跳过的相关配置:

SET spark.databricks.delta.stats.collect = true;
SET spark.databricks.delta.stats.localLimit = 10000000;

我们还在查询中加入了分区过滤,确保每次查询都指定分区范围,避免全表扫描。

效果对比

优化之后,大部分查询都能利用数据跳过,平均扫描的数据量减少了60%,查询速度又提升了不少。

六、Vacuum清理过期文件

第五个优化是用VACUUM清理过期文件。

问题分析

Delta Lake支持时间旅行,可以查询历史版本的数据。但是这也意味着,即使数据被删除或者更新了,旧的文件还会保留,占用存储空间。时间长了,存储成本会越来越高,而且文件太多也会影响性能。

我们的情况是,保留了所有历史版本,存储用量是实际数据量的3倍多,而且还在不断增长。

优化方法

Delta Lake提供了VACUUM命令,可以清理过期的文件。用法:

VACUUM delta.`/path/to/table` RETAIN 168 HOURS

这个命令会删除超过保留期的旧文件,只保留最近7天(168小时)的版本。保留期可以根据需要调整,但是不能太短,否则可能会影响正在运行的查询。

我们的做法是每周执行一次VACUUM,保留7天的历史版本。对于一些需要长期保留历史的表,保留30天。

效果对比

执行VACUUM之后,存储用量减少了60%,节省了大量的存储成本。而且文件数量减少了,查询和写入的性能也有一定提升。

注意事项

VACUUM是一个危险操作,删除的文件无法恢复。执行之前一定要确认保留期设置合理,而且没有正在使用旧版本的查询。另外,VACUUM不会删除最新版本的数据,只会删除过期的旧文件。

七、其他优化

除了上面几个主要的优化,我们还做了一些其他的优化。

1. 调整Spark配置

我们调整了一些Spark的配置,比如增加executor内存、调整并行度、开启AQE(自适应查询执行)等。这些调整虽然不是Delta Lake特有的,但是对性能提升也有帮助。

2. 使用正确的文件格式

Delta Lake底层用的是Parquet格式,我们确保所有数据都是用Parquet存储的,而且开启了Snappy压缩。Parquet的列式存储和压缩,对查询性能有很大帮助。

3. 避免过度更新

Delta Lake的更新操作是通过写新文件实现的,频繁更新会产生很多小文件。我们尽量把更新操作合并成批量更新,减少写入次数。

4. 监控和告警

我们搭建了监控系统,监控Delta Lake表的文件数量、存储用量、查询性能等指标。发现异常及时处理,避免问题积累。

八、优化效果总结

经过这一系列优化,我们的Delta Lake性能有了质的飞跃。

查询方面,平均查询时间从15分钟降到了2分钟左右,提升了7倍多。一些复杂的查询,从几个小时降到了十几分钟。

写入方面,每天的增量写入时间从3小时降到了40分钟,提升了4倍多。

存储方面,通过VACUUM清理,存储用量减少了60%,节省了大量成本。

当然,这些优化不是一劳永逸的。随着数据量的增长和业务的变化,还需要持续监控和调整。但是至少目前来看,效果是非常好的。

九、写在最后

Delta Lake是一个很优秀的数据湖存储层,但是要用好它,需要了解它的特性,做针对性的优化。小文件合并、合理分区、Z-Order、数据跳过、VACUUM清理,这些都是非常有效的优化手段。

如果你也在用Delta Lake,而且遇到了性能问题,不妨试试这些优化方法。每一项都不难操作,但是组合起来效果会非常明显。

当然,每个项目的情况不同,优化的重点也不同。建议先分析自己的瓶颈在哪里,然后针对性地做优化,不要盲目照搬。

最后用一句话结束本文:"性能优化是一个持续的过程,没有最好,只有更好。"希望这篇文章能帮到正在使用Delta Lake的你,也欢迎大家在评论区交流自己的优化经验。