Delta Lake是Databricks开源的存储框架,在数据湖之上提供ACID事务、Schema演进、时间旅行等功能,这几年很火,是数据湖三剑客之一(另外两个是Apache Iceberg和Apache Hudi)。

我们团队两年前开始用Delta Lake做数据湖架构,从最开始的调研测试,到后来的生产环境大规模使用,踩了很多坑,也总结了很多经验。今天分享实战过程中踩过的坑和解决方案,包括小文件问题、性能问题、Schema演进的坑、和其他系统集成的问题、运维的坑等,希望能给打算用Delta Lake的朋友一些参考。

我们的环境:Delta Lake 1.0,Spark 3.1,存储是S3兼容的对象存储,数据量几十TB,每天增量几百GB,几百张Delta表,跑在Kubernetes集群上,用Airflow调度。

一、为什么选Delta Lake

在说踩坑之前,先说说为什么选Delta Lake。我们之前用传统数据仓库,Hive加Parquet,遇到很多问题:没有ACID事务,写入过程中读会读到不完整数据,写入失败会留脏数据;Schema演进困难,改表结构很麻烦;没有时间旅行,数据写错了想回滚很难;upsert和删除困难,只能全量重写表。

为了解决这些问题,我们调研了数据湖三剑客,最后选了Delta Lake,主要原因:和Spark集成最好,API最简单,团队Spark经验丰富;功能完善,ACID、Schema演进、时间旅行、upsert都有;社区活跃,资料多;Databricks背书,发展有保障。

二、坑一:小文件问题严重,性能差

第一个大坑就是小文件问题,这也是用Delta Lake最常见的坑。

最开始我们按普通Spark写Parquet的方式写Delta表,每个任务写一次就生成一批文件。跑了几个月,有些表有几十万甚至上百万个小文件,每个文件才几MB甚至几KB。查询特别慢,因为Spark要读大量小文件的元数据,开销很大,并行度也不好控制。Delta的事务日志也会记录每个文件信息,文件太多日志也大,读元数据也慢。

解决方案:

第一,控制写入的文件数量和大小。 写数据前用repartition或coalesce控制输出文件数,让每个文件128MB到1GB左右。repartition会触发shuffle有开销,coalesce不会shuffle但可能数据倾斜,根据实际情况选。

第二,定期做OPTIMIZE合并小文件。 Delta Lake提供OPTIMIZE命令,自动合并小文件。我们每天凌晨对核心表跑一次OPTIMIZE,效果很明显。OPTIMIZE还支持ZORDER排序,把经常一起查询的列放在一起,提升过滤查询性能。

第三,配置自动优化参数。 Delta有一些参数能自动优化小文件,比如自动合并小文件、目标文件大小等。我们配置了这些参数,写入时就能控制文件大小。不过有些自动优化是Databricks运行时特有的,开源版支持有限,我们主要靠手动OPTIMIZE。

第四,合理设置分区。 分区字段基数太高会产生大量分区,小文件问题更严重。比如按用户ID分区会有几百万分区。我们有张表最开始按用户ID分区,小文件爆炸,后来改成按日期分区,用户ID做ZORDER,问题就解决了。

三、坑二:Schema演进的坑

Delta Lake支持Schema演进,能方便地加字段、改注释,但用的过程中也踩了坑。

坑一:加字段默认值不回填。 有次给表加了个字段设了默认值,以为历史数据会自动填默认值,结果查询历史数据发现这个字段是null。原来Delta的默认值只对新写入的数据生效,历史数据不会自动回填。解决方法是手动跑UPDATE把历史数据的字段填默认值,或者重写历史数据。

坑二:改字段类型不支持。 Delta不支持改字段类型,比如int改long。有次数据量大了int不够用想改long,直接报错。解决方案是用overwriteSchema重写表,把数据读出来转成新Schema再覆盖写回去。但这样慢,期间不能写数据。建议建表时字段类型尽量选大一点,比如能用long就不用int。

坑三:删字段改字段名不支持。 这些不兼容的Schema变更也需要重写表,用overwriteSchema覆盖。建议建表时字段设计谨慎,减少变更。

四、坑三:和其他系统集成的问题

我们的数据架构不是只有Spark和Delta,还有Hive、Presto、Flink等,集成过程中踩了不少坑。

坑一:Hive读Delta表的问题。 最开始以为Delta表就是Parquet加事务日志,Hive直接读Parquet就行,结果不行。因为Delta表可能有多个版本文件,直接读会读到已删除的旧数据,也读不到最新Schema。解决方案是用Delta的Hive连接器,把Delta表注册成Hive表,用Delta的InputFormat读。但Hive读Delta是只读的,性能也比Spark差一些。

坑二:Presto读Delta表的问题。 Presto最开始不支持Delta,后来加了连接器但支持不完善,时间旅行等功能用不了。我们升级了Presto版本,只用基础查询功能,比较稳定。复杂查询还是用Spark。

坑三:Flink写Delta表的问题。 Flink-Delta连接器还比较初级,不支持upsert只支持append,稳定性也一般。我们测试遇到bug没敢在生产用。现在实时数据先写Kafka,再用Spark Structured Streaming读Kafka写Delta,虽然多一步但稳定可靠。

五、坑四:时间旅行和数据回滚的坑

Delta的时间旅行功能很好用,但也踩了坑。

坑一:历史版本保留时间。 Delta默认保留30天历史版本,超过会被清理。有次想回滚到两个月前的版本,发现已经被清理了。解决方案是根据需求用delta.logRetentionDuration调整保留时间,核心表保留180天。但保留越长日志越大,占用存储越多,要权衡。

坑二:VACUUM误删数据。 VACUUM命令清理不用的旧文件节省存储,但保留时间设太短会把还需要的历史版本文件删掉。有次测试VACUUM保留0小时,把所有旧版本文件都删了,只剩最新版本,还好是测试环境。生产环境VACUUM保留时间至少7天,操作前确认不需要历史版本了再做。

六、坑五:并发写入的问题

Delta支持并发写入,有乐观并发控制,但并发高时还是会冲突。

有张表多个任务同时写,经常报并发冲突错误。原因是多个任务同时修改同一个分区,乐观控制检测到冲突就让其中一个失败重试。

解决方案: 第一,尽量让不同任务写不同分区,减少冲突。比如按日期分区,每个任务写不同日期。第二,配置自动重试次数,大部分冲突重试就能成功。第三,高并发upsert场景按业务线分区,每个业务线写自己的分区,冲突少很多。

七、性能优化的经验

除了小文件,我们还做了一些性能优化。

第一,分区和ZORDER结合。 把经常过滤的低基数字段做分区(日期、地区),高基数字段做ZORDER(用户ID、订单ID)。我们有张订单表按日期分区、用户ID ZORDER,查某个用户的订单性能提升了10倍以上。

第二,数据跳过(Data Skipping)。 Delta会收集文件的统计信息,查询时跳过不符合条件的文件。要让过滤字段在前32列,这样能收集统计信息享受数据跳过。我们建表时把经常过滤的字段放前面。

第三,缓存热点数据。 对经常查询的热点表或分区,用Spark cache缓存到内存,查询性能提升明显。我们对维度表和热点分区做了缓存。

第四,合理设置Spark参数。 executor内存、core数量、并行度、shuffle分区数等,根据集群资源和任务情况合理设置。不同任务不同参数,大任务多给资源,小任务少给,平衡性能和资源利用率。

八、运维的坑

最后说说运维的坑。

坑一:事务日志太大。 Delta的事务日志存在deltalog目录,提交频繁或保留时间长的话日志会很大,占用存储,读元数据也慢。解决方案是定期清理旧日志,配置日志保留时间自动清理。我们日志保留30天,超过自动清理。

坑二:元数据管理。 Delta的元数据在事务日志里,和Hive元数据分开。用Hive读Delta需要把表注册到Hive,Schema变更时要同步更新Hive元数据。我们用脚本每天自动同步Delta Schema到Hive。

坑三:监控和告警。 要监控表的文件数量、小文件比例、查询性能、写入延迟等指标。我们用Prometheus加Grafana做监控大盘,关键指标设告警。小文件比例超阈值就提醒跑OPTIMIZE,查询延迟升高就排查。

九、写在最后

以上就是我们用Delta Lake两年多踩过的主要的坑和解决方案。总结一下:Delta Lake确实很好用,解决了传统数据湖的很多问题,ACID、Schema演进、时间旅行、upsert都很实用,值得用。但小文件问题是最常见的坑,一定要注意控制写入文件大小,定期OPTIMIZE。Schema演进要注意默认值不回填、改类型不支持。和其他系统集成要注意兼容性,尽量用Spark读写最稳定。时间旅行和VACUUM要小心,避免误删数据。并发写入要注意冲突,合理分区配置重试。运维要做好监控告警。

现在数据湖三剑客都在快速发展,各有优缺点。我们因为Spark用得多选了Delta,用下来整体满意。如果在考虑用Delta Lake,希望这篇文章能帮你少踩坑,顺利搭建数据湖架构。

如果你用Delta也踩过其他坑,欢迎在评论区留言分享,大家一起交流学习。祝大家的数据湖项目都能顺利上线,稳定运行。