最近在项目中用Delta Lake做数据湖,踩了不少坑,也积累了一些实战经验。
Delta Lake是Databricks开源的一个数据湖存储层,基于Parquet格式,提供了ACID事务、Schema演进、时间旅行、数据校验等企业级特性。它的目标是解决数据湖的痛点,让数据湖也能像数据仓库一样可靠、高效、易用。
我们项目一开始用的是纯Parquet格式的数据湖,后来因为数据质量问题、并发写入问题、历史数据回溯问题,引入了Delta Lake。用了之后,确实解决了很多问题。但也踩了不少坑,有性能上的,有使用上的,有运维上的。
本文总结Delta Lake的核心概念、常见坑点、性能优化、最佳实践等内容,结合实际项目中的经验,给打算用Delta Lake的朋友一些参考。
一、Delta Lake核心概念
在说踩坑之前,先简单介绍一下Delta Lake的核心概念,帮助理解后面的内容。
1. ACID事务
Delta Lake最核心的特性,就是ACID事务。传统的数据湖,基于Parquet或ORC格式,只支持文件级别的原子性,不支持跨文件的事务。如果一个作业写入了一半失败了,就会留下脏数据,读者可能会读到不完整的数据。
Delta Lake通过事务日志(deltalog目录),实现了跨文件的ACID事务。所有的写入操作,都会先写数据文件,然后写事务日志,记录这次写入的元数据。读者读取的时候,会根据事务日志,找到当前版本的所有数据文件,保证读到的是一致的、完整的数据。
2. 事务日志
事务日志是Delta Lake的核心,它记录了表的所有变更历史,包括添加了哪些文件、删除了哪些文件、Schema变更、元数据变更等。事务日志是一系列的JSON文件,每个文件对应一个事务版本,按序号命名,比如000000.json、000001.json。
读取数据的时候,Delta Lake会读取事务日志,找到当前版本的所有数据文件。然后读取这些文件。写入数据的时候,Delta Lake会先写数据文件。然后追加一个新的事务日志文件,记录这次写入。
事务日志是Delta Lake实现ACID事务、时间旅行、Schema演进的基础。
3. 时间旅行
因为事务日志记录了所有的变更历史,所以Delta Lake支持时间旅行,可以查询历史任意版本的数据。你可以查询某个时间点的数据,也可以查询某个版本号的数据,非常方便。
时间旅行在很多场景下都很有用,比如数据出错了,可以回溯到之前的版本;比如需要对比历史数据,看数据的变化;比如审计需要,查看某个时间点的数据状态。
4. Schema演进
Delta Lake支持Schema演进,可以在不重写数据的情况下,修改表的Schema,比如添加列、删除列、修改列类型、重命名列等。
传统的Parquet格式,Schema是写死的,如果要修改Schema,就需要重写所有数据,成本很高。Delta Lake通过事务日志记录Schema的变更,读取的时候根据最新的Schema来解析数据,实现了Schema演进,不需要重写数据。
5. 数据校验
Delta Lake支持数据校验。可以在写入的时候,对数据进行校验。比如非空校验、范围校验、枚举校验等。如果数据不符合校验规则,就会写入失败,保证数据质量。
数据校验在数据湖场景下非常重要,因为数据湖的数据来源多。格式杂,质量参差不齐。很容易有脏数据。通过数据校验。可以在写入的时候就拦截脏数据,保证数据湖的数据质量。
二、踩坑总结
介绍完核心概念,下面说说我们在实际使用中踩的坑。
1. 小文件问题
这是Delta Lake最常见的坑,也是最影响性能的问题。
Delta Lake每次写入,都会生成一批新的数据文件。如果写入频率很高,每次写入的数据量很小,就会生成大量的小文件。小文件多了之后,读取性能会急剧下降,因为读取的时候需要打开大量的文件,每个文件都有IO开销。而且元数据也会变得很大。
我们项目一开始,是流式写入。每5分钟写一次,每次写几万条数据。结果一个月下来,生成了几十万个小文件。查询一张表需要好几分钟,非常慢。
解决方法是定期做OPTIMIZE,把小文件合并成大文件。Delta Lake提供了OPTIMIZE命令,可以把小文件合并成大文件,提高读取性能。我们设置了每天凌晨做一次OPTIMIZE,合并前一天的小文件,性能提升了很多,查询时间从几分钟降到了几秒。
另外,也可以调整写入的批次大小。减少写入频率,每次多写一些数据,从源头上减少小文件的产生。
2. 事务日志膨胀问题
因为Delta Lake的所有变更都记录在事务日志里。所以如果表的变更很频繁,事务日志会变得越来越大,文件越来越多。读取的时候。需要读取所有的事务日志文件,才能知道当前版本的数据文件,这会影响读取性能。
我们项目有一张表,每天写入几十次。半年下来,事务日志有几万个文件。每次读取光解析事务日志就要十几秒,非常慢。
解决方法是定期做VACUUM,清理旧的事务日志和旧的数据文件。VACUUM会删除超过保留期的事务日志和数据文件,只保留最近的版本。我们设置了保留期30天,每周做一次VACUUM,清理30天之前的事务日志和数据文件,性能提升了很多。
但要注意,VACUUM之后,就不能再时间旅行到被删除的版本了。所以保留期要根据业务需求来设置。不能太短。
3. 并发写入冲突
Delta Lake支持并发写入,但也有冲突的可能。如果多个作业同时写入同一张表,并且修改了相同的文件,就会发生冲突,导致其中一个作业失败。
Delta Lake使用乐观并发控制,写入的时候不会加锁,而是在提交的时候检查是否有冲突。如果有冲突,就会失败。需要重试。
我们项目一开始,有多个作业同时写入同一张表。经常发生冲突,作业失败。需要手动重试,很麻烦。
解决方法是尽量避免多个作业同时写入同一张表,或者把写入操作串行化。我们后来调整了作业调度,让写入同一张表的作业错开时间,不要同时运行,冲突就少了很多。
另外,Delta Lake也支持自动重试。可以设置重试次数,发生冲突的时候自动重试,不需要手动干预。
4. Schema演进的坑
Delta Lake支持Schema演进,但也有一些坑。
比如,添加列的时候,如果新列有默认值。那么历史数据的这个列会是null,而不是默认值。因为历史数据文件里没有这个列,读取的时候会用null填充,而不是默认值。默认值只对新写入的数据生效。
再比如,修改列类型的时候。不是所有类型都能互相转换。比如string转int,如果历史数据有不能转成int的字符串,就会报错。而且,修改列类型需要重写数据。不是元数据操作,性能比较差。
还有,删除列的时候。Delta Lake实际上不是真的删除了列,而是在元数据里标记为删除。数据文件里还是有这个列,只是读取的时候不读了。这会导致数据文件里有冗余的列,浪费存储空间。
我们项目就踩过添加列默认值的坑,添加了一个有默认值的列。以为历史数据也会有默认值,结果查询的时候历史数据都是null。导致业务逻辑出错。后来才知道,默认值只对新数据生效,历史数据需要手动回填。
所以,使用Schema演进的时候,一定要仔细看文档,了解每个操作的行为。不要想当然。
5. 与Spark版本的兼容性
Delta Lake和Spark的版本,有严格的对应关系,不是所有版本都能兼容。如果Delta Lake的版本和Spark的版本不匹配,就会出现各种奇怪的问题,比如读取失败、写入失败、性能下降、数据损坏等。
我们项目一开始,用的是Spark 3.0。Delta Lake 0.7,后来升级Spark到3.1。Delta Lake没升级,结果出现了读取数据报错的问题,查了很久才发现是版本不兼容。
解决方法是升级Delta Lake到兼容的版本,或者降级Spark到兼容的版本。Databricks的官方文档里,有详细的版本对应关系,升级之前一定要查清楚,不要盲目升级。
另外,Delta Lake的不同版本之间。也有兼容性问题。比如旧版本写的数据,新版本能不能读。新版本写的数据,旧版本能不能读。都要注意。一般来说,新版本是向后兼容的,能读旧版本的数据。但旧版本不一定能读新版本的数据。尤其是用了新特性的时候。
6. 分区列的选择
Delta Lake支持分区,和Hive表类似。可以按某个列分区。提高查询性能。但分区列的选择,很重要,选不好会影响性能。
如果分区列的基数太高。比如按用户ID分区,就会有大量的分区。每个分区的数据量很小,导致小文件问题,查询性能反而下降。
如果分区列的基数太低。比如按性别分区,只有两个分区。每个分区的数据量很大,查询的时候还是要扫描大量数据,分区的效果不明显。
我们项目一开始,按日期分区。每天一个分区,效果还可以。后来有一张表。数据量很大,我们又加了一个地区列做二级分区。结果查询性能反而下降了,因为地区列的基数不高不低。导致分区数太多,小文件太多。后来去掉了二级分区。只按日期分区,性能反而好了。
所以,分区列的选择,要根据数据量和查询模式来定,一般选择基数适中、经常作为查询条件的列做分区,比如日期、地区、部门等。不要选择基数太高或太低的列做分区。
7. 数据校验的性能影响
Delta Lake支持数据校验,可以在写入的时候校验数据,但数据校验会影响写入性能,尤其是校验规则复杂的时候。
我们项目有一张表,设置了十几个校验规则。结果写入性能下降了30%,因为每条数据都要做校验,开销很大。
解决方法是只在关键的列上设置校验规则。不要在所有列上都设置。而且,校验规则尽量简单。不要用复杂的表达式,复杂的校验可以在写入之前,用Spark作业提前做。不要依赖Delta Lake的校验。
另外,也可以只在重要的表上设置校验,不重要的表可以不设置,平衡数据质量和性能。
三、性能优化
说完踩坑,再说说性能优化。Delta Lake的性能优化,主要包括读取优化和写入优化。
1. 读取优化
读取优化,主要是减少扫描的数据量,提高查询效率。
第一,分区裁剪。查询的时候。尽量带上分区列的条件,这样Delta Lake只会扫描相关的分区。不会扫描全表,性能会好很多。
第二,数据跳过。Delta Lake支持数据跳过,会在事务日志里记录每个文件的最小最大值,查询的时候,如果查询条件不在某个文件的最小最大值范围内,就会跳过这个文件,不读取。数据跳过对数值型、日期型的列效果很好。可以大大减少扫描的数据量。
第三,Z-Order。Delta Lake支持Z-Order排序。可以把相关的列放在一起,提高数据跳过的效果。比如,经常一起查询的列。可以做Z-Order,这样查询的时候,数据跳过的效果更好,扫描的数据量更少。
第四,缓存。对于经常查询的表。可以缓存到内存里,提高查询速度。Delta Lake和Spark集成得很好。可以用Spark的缓存机制,把表缓存到内存里。
2. 写入优化
写入优化,主要是提高写入速度,减少小文件的产生。
第一,调整批次大小。写入的时候。尽量每次多写一些数据,减少写入次数。从源头上减少小文件的产生。流式写入的话。可以调整触发间隔,延长间隔时间,每次多写一些数据。
第二,调整并行度。写入的时候,并行度不要太高。否则每个并行任务写的数据量很小,会产生很多小文件。可以根据数据量,调整并行度,让每个任务写的数据量在合适的范围。一般每个文件128MB到1GB比较合适。
第三,OPTIMIZE。定期做OPTIMIZE。把小文件合并成大文件,提高读取性能。OPTIMIZE可以按分区做。也可以全表做,根据需要选择。
第四,自动优化。Delta Lake支持自动优化。可以在写入的时候,自动合并小文件,不需要手动做OPTIMIZE。可以通过配置参数开启自动优化,适合写入频繁的场景。
四、最佳实践
最后,总结一些Delta Lake的最佳实践。
- 定期做OPTIMIZE和VACUUM,OPTIMIZE合并小文件,VACUUM清理旧数据,保持表的健康。
- 合理选择分区列,选择基数适中、经常查询的列做分区,不要选择基数太高或太低的列。
- 控制小文件数量,从写入端控制批次大小和并行度,定期做OPTIMIZE合并。
- 谨慎使用Schema演进,了解每个操作的行为,不要想当然,重要的Schema变更要先测试。
- 注意版本兼容性,Delta Lake和Spark的版本要匹配,升级之前查清楚版本对应关系。
- 合理设置数据校验,只在关键列设置简单的校验规则,复杂的校验提前做,不要影响写入性能。
- 做好监控,监控表的文件数量、数据量、查询性能、写入性能,发现问题及时处理。
- 重要的表,做好备份,虽然Delta Lake有时间旅行,但还是要定期备份,防止意外情况。
五、写在最后
Delta Lake是一个很好的数据湖存储层,它解决了传统数据湖的很多痛点,提供了ACID事务、时间旅行、Schema演进、数据校验等企业级特性,让数据湖也能像数据仓库一样可靠、高效、易用。
但它也不是银弹。不是用了就万事大吉了,它也有很多坑。需要我们去了解。去避免。只有深入理解它的原理,掌握它的最佳实践。才能用好它,发挥它的价值。
我们项目用了Delta Lake之后,数据质量提高了。并发问题解决了,历史数据回溯方便了。整体来说,收益还是很大的。虽然踩了不少坑。但这些坑都是可以通过学习和实践避免的。
如果你也在做数据湖,也在为数据质量、并发写入、历史回溯等问题头疼,不妨试试Delta Lake,它可能会给你带来惊喜。
希望这篇文章能给打算用Delta Lake的朋友一些参考,也欢迎大家交流讨论,一起学习进步。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录