Delta Lake,是Databricks开源的。存储层框架,在Parquet之上。提供了ACID事务、Schema演进、时间旅行、UPSERT/DELETE/MERGE等功能,解决了数据湖的很多痛点。最近几年,非常火。越来越多的公司,开始用Delta Lake。构建数据湖。或者,Lakehouse架构。
我在项目中,用Delta Lake。已经有一年多了。从。最开始的,调研。试用。到。后来的,大规模。生产环境。使用。踩了很多坑。也总结了很多经验。和,最佳实践。
之前,我写过一篇。Delta Lake的踩坑总结。分享了。我们在使用过程中,遇到的。一些问题。和。解决方法。今天。这篇文章。分享。Delta Lake的最佳实践。包括,表设计、数据写入、数据读取、性能优化、运维管理、常见问题等方面。希望。能给。正在,使用。或者。打算。使用Delta Lake的朋友,一些参考。
先说明一下。这些最佳实践。是。基于。我们的,项目经验。和。场景。总结的,不一定。适合。所有的。场景,大家。要。根据。自己的。实际情况。灵活,参考。不要。生搬硬套。而且。Delta Lake。发展,很快。版本。更新。也,很快。我。写这篇文章的时候。用的,是。Delta Lake 1.0.x。版本。不同的,版本。可能。会有。一些。差异,大家。要。注意。版本。的,差异。
一、Delta Lake简介
在,分享。最佳实践。之前。先。简单。介绍一下。Delta Lake。是什么,帮助。不了解的朋友。有个。概念。
Delta Lake,是。一个。开源的。存储层,框架。它。不是。一个。新的,存储。格式。而是。在。Parquet。之上,加了。一层。事务。日志。也就是。deltalog,目录。通过。这个。事务。日志。来,实现。ACID事务。Schema演进。时间旅行,UPSERT/DELETE/MERGE等。功能。
简单来说,Delta Lake。就是。给。Parquet。加了。事务。和。更多。功能。让,数据湖。也能。像。数据仓库。一样。支持。事务。支持,更新。删除。支持。Schema演进。支持。时间旅行。解决了。传统,数据湖的。很多。痛点。比如。数据。不一致。更新。删除。困难,Schema。演进。困难。等。
Delta Lake。和。Spark,集成。得。很好。支持,Spark的。批处理。和。流处理。也。支持。其他。计算。引擎。比如。Presto,Hive。Flink等。当然。支持。最好的。还是,Spark。
Delta Lake,的。核心。特性。包括:
- ACID事务:支持,原子性。一致性。隔离性。持久性。多个。写入。不会。互相。干扰,读者。不会。读到。不完整。的。数据。
- 可扩展的,元数据。处理:元数据。存储。在。事务。日志。中。支持,处理。PB级。的。数据。和。数十亿。的。文件。
- Schema演进:支持,Schema。的。增加。修改。删除。不需要。重写。数据,能。自动。处理。Schema。的。变化。
- 时间旅行:支持,查询。历史。版本。的。数据。可以。查看。数据。的,历史。变化。也。可以。回滚。到。历史。版本。
- UPSERT/DELETE/MERGE:支持,更新。删除。合并。数据。解决了。传统。数据湖。只能,追加。不能。更新。删除。的。痛点。
- 统一,批处理。和。流处理:Delta Lake。既。支持。批处理。写入。和。读取。也。支持,流处理。写入。和。读取。一套。存储。支持。两种。处理,方式。
- 数据,质量。保障:支持。Schema。校验。约束。等。能。保障。写入,数据。的。质量。
以上。就是。Delta Lake,的。简单。介绍。和,核心。特性。接下来。分享,Delta Lake。的。最佳实践。
二、表设计最佳实践
表设计,是。使用。Delta Lake。的。第一步。也是。最重要的。一步。表。设计得,好。后面的。性能。运维。都会。顺利。很多。表。设计得。不好,后面。会。遇到。很多。性能。问题。和。运维。问题。
1. 选择,合适的。分区。列
分区,是。Delta Lake。最重要的。性能。优化。手段。之一。选择,合适的。分区。列。能。大大。提高。查询。性能,减少。扫描。的。数据量。
选择,分区。列。的。原则:
- 选择,经常。用于。查询。过滤。的。列。作为。分区,列。这样。查询的时候。能。跳过。不需要的。分区。减少。扫描,的。数据量。
- 选择,基数。比较低。的。列。作为。分区。列。也就是。 distinct 值。比较少。的,列。比如。日期。地区。部门。业务。类型。等。不要。选择。基数,很高。的。列。比如。用户ID。订单ID。作为。分区。列。因为。这样,会。产生。大量的。分区。每个。分区。的。文件,很小。反而。会。影响。性能。
- 分区,的。数量。不要。太多。也。不要。太少。一般。建议。每个。表。的。分区,数量。在。几千。到。几万。之间。不要。超过。十万,太多的。分区。会。导致。元数据。过大。查询。元数据,的。时间。变长。也。会。导致。小文件。问题。
- 最常用的,分区。列。是。日期。比如。dt。或者。date。按。天,分区。或者。按。月。分区。这是。最常见。也。最推荐的。分区。方式。因为,大部分。查询。都会。按。时间。过滤。而且。日期。的。基数。也。比较,合适。
我们的,实践。是。大部分。表。都。按。dt。日期。分区。按,天。分区。对于。数据量。特别大。的。表。可以。按。小时。分区。对于,数据量。比较小。的。表。可以。按。月。分区。同时。对于,经常。按。地区。或者。业务。类型。查询。的。表。可以,加。二级。分区。比如。dt。和。region。或者。dt。和。biz_type。但是。要。注意,分区。的。数量。不要。太多。
2. 选择,合适的。排序。列。和。Z-Order
除了,分区。Delta Lake。还。支持。排序。和。Z-Order。来。优化。查询,性能。
排序。就是。把,数据。按照。某个。列。排序。这样。查询。的时候。按照。这个,列。过滤。就能。跳过。很多。不需要的。数据。提高。查询,性能。
Z-Order,是。一种。多维。排序。技术。能。同时。按照。多个。列。排序,让。多个。列。的。过滤。都。能。受益,提高。多维度。查询。的。性能。
选择,排序。列。和。Z-Order。列。的。原则:
- 选择,经常。用于。查询。过滤。的。列。作为。排序,列。或者。Z-Order。列。
- 如果。只有,一个。列。经常。用于。过滤。用。排序。就。够了。如果,有。多个。列。经常。用于。过滤。用。Z-Order,效果。更好。
- 不要,选择。基数。太高。的。列。作为。排序。列。或者,Z-Order。列。比如。用户ID。订单ID。因为。这样。排序。的。效果。不好。而且,会。增加。写入。的。成本。
- 排序。和,Z-Order。会。增加。写入。的。成本。因为。写入的时候。需要。排序。所以,要。权衡。写入。性能。和。查询。性能。对于。写入。频繁,查询。少。的。表。可以。不做。排序。或者。Z-Order。对于。写入,少。查询。多。的。表。建议。做。排序。或者。Z-Order。
我们的,实践。是。对于。经常。查询。的。表。选择。1-3个。经常,用于。过滤。的。列。做。Z-Order。比如。userid。productid,region。等。但是。要。注意。不要。选。太多。列。一般。1-3个。就,够了。太多。效果。不好。而且。写入。成本。高。同时。定期,对。表。做。OPTIMIZE。和。Z-Order。合并。小文件。优化,排序。
3. 选择,合适的。文件。大小
文件,大小。也是。影响。Delta Lake。性能。的。重要。因素,文件。太小。会。导致。小文件。问题。查询。的时候。需要,打开。很多。文件。元数据。开销。大。性能。差;文件,太大。会。导致。查询。的时候。无法。并行。或者。并行度。不够。性能。也。差。而且,更新。删除。的时候。需要。重写。整个。大文件。成本,高。
Delta Lake,推荐的。文件。大小。是,128MB。到。1GB。之间。一般。建议,128MB。或者。256MB。这是。比较。合适的。文件,大小。
为了,控制。文件。大小。Delta Lake。提供了。OPTIMIZE。命令。能,自动。合并。小文件。把。小文件。合并。成。大文件,达到。合适的。大小。
我们的,实践。是。定期。对,表。做。OPTIMIZE。合并,小文件。一般。每天。或者。每周。做一次。根据。数据量。和。写入,频率。决定。同时。在。写入,的时候。也。可以。通过。调整。Spark的。参数。比如,spark.sql.files.maxRecordsPerFile。来。控制。每个,文件。的。记录数。从而,控制。文件。大小。避免,产生。太多。小文件。
4. Schema设计,要。合理。避免。频繁,的。Schema演进
虽然,Delta Lake。支持。Schema演进。但是。频繁。的。Schema演进。还是。会。带来。一些,问题。比如。历史。版本。的。查询。性能。差。元数据,过大。等。所以。在。表。设计。的时候。要。尽量,把。Schema。设计。合理。避免。频繁。的,Schema演进。
Schema设计,的。原则:
- 字段,的。命名。要。规范。统一。有。意义。不要。用,缩写。或者。拼音。除非。是。通用的。缩写。
- 字段,的。类型。要。合适。不要,用。太大。的。类型。比如,能用。int。就。不要。用。bigint,能用。string。就。不要。用。varchar。因为,Delta Lake。的。string。是,可变。长度。的。和。varchar,一样。性能。也。一样。
- 尽量。不要,用。嵌套。太复杂。的。结构。比如。多层。的。struct,array。map。因为。这样。会。增加。查询。的。复杂度。和。性能,开销。如果。确实。需要。嵌套。结构。尽量。不要。超过。2层。
- 预留,一些。扩展。字段。或者。用。一个。扩展。的。map。或者,string。字段。来。存储。未来。可能。增加。的。字段,避免。频繁。的。Schema演进。当然。这。要。根据。实际。情况。决定。不要,为了。避免。Schema演进。而。设计。得。太。冗余。
我们的,实践。是。在。表。设计。的时候。充分。考虑,未来。的。业务。发展。预留。一些。可能。会。用到,的。字段。同时。对于。可能。频繁。变化。的。扩展。信息。用。一个,ext。的。string。字段。存储。JSON。格式。的,数据。这样。增加。扩展。字段。的时候。不需要。修改,表。的。Schema。只。需要。在。JSON。里。加,字段。就行。当然。这样。查询。的时候。会。稍微。麻烦。一点。需要。用,from_json。函数。解析。但是。能。避免。频繁。的。Schema演进。还是,值得的。
三、数据写入最佳实践
数据写入,是。Delta Lake。最。常用的。操作。之一。写入。的,方式。和。参数。直接。影响。写入。的。性能。和。后续,查询。的。性能。
1. 用,Spark。写入。是。最佳。选择
虽然,Delta Lake。支持。多种。计算。引擎。但是。和。Spark。的,集成。是。最好的。功能。支持。最。完善。性能。也。最好。所以,用。Spark。写入。Delta Lake。是。最佳。选择。
用,Spark。写入。Delta Lake。很简单。和。写入,Parquet。差不多。只。需要。把,format。改成。delta。就行:
df.write.format("delta").mode("append").save("/path/to/table")或者,用。表名:
df.write.format("delta").mode("append").saveAsTable("db.table")支持,的。写入。模式。有,append。overwrite。ignore。errorifexists。和。Parquet,一样。
2. 批量,写入。避免。频繁。小批量。写入
Delta Lake。虽然。支持,流处理。写入。但是。频繁。的。小批量。写入。会。导致,产生。大量的。小文件。和。大量的。事务。日志。版本,影响。查询。性能。和。元数据。性能。所以。尽量。批量,写入。避免。频繁。小批量。写入。
对于,流处理。写入。建议。调整。触发。间隔。不要。太。频繁。比如。5分钟。或者,10分钟。触发。一次。不要。几秒钟。就。触发。一次。这样。能。减少,小文件。的。产生。同时。定期。做。OPTIMIZE。合并,小文件。
对于,批处理。写入。建议。尽量。一次。写入。足够。的。数据量。不要,分。很多次。小批量。写入。比如。一天。的。数据。一次,写入。不要。分成。很多次。写入。
我们的,实践。是。批处理。任务。尽量。一天。一次。或者。几小时,一次。写入。不要。太。频繁;流处理。任务。触发。间隔。设置,为。5-10分钟。不要。太。短;同时。每天。对。写入。频繁,的。表。做。一次。OPTIMIZE。合并。小文件。
3. 用,MERGE INTO。做。UPSERT。和。更新。删除
Delta Lake,支持。MERGE INTO。语法。能。实现。UPSERT。也就是。存在。就。更新。不存在。就,插入。也。能。实现。更新。和。删除。这是,Delta Lake。非常。重要的。功能。解决了。传统。数据湖。无法,更新。删除。的。痛点。
MERGE INTO,的。语法。和。SQL。的。MERGE。差不多:
MERGE INTO target_table t
USING source_table s
ON t.id = s.id
WHEN MATCHED THEN
UPDATE SET t.name = s.name, t.age = s.age
WHEN NOT MATCHED THEN
INSERT (id, name, age) VALUES (s.id, s.name, s.age)也,支持。DELETE:
MERGE INTO target_table t
USING source_table s
ON t.id = s.id
WHEN MATCHED AND s.is_deleted = true THEN
DELETE
WHEN MATCHED THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *用,MERGE INTO。的。注意事项:
- MERGE INTO,的。条件。要。尽量。用。分区。列。和。排序,列。这样。能。减少。扫描。的。数据量。提高。性能。
- 源表,的。数据量。不要。太大。如果。源表。数据量。很大。建议。先。过滤。或者,分批次。MERGE。避免。性能。太差。
- MERGE INTO,会。重写。涉及。的。文件。所以。如果。更新。删除。的。数据量,很大。会。产生。很多。重写。成本。建议。尽量。减少,大范围。的。更新。删除。
- 不要,用。MERGE INTO。做。全量。覆盖。全量。覆盖。用。overwrite。模式,更。高效。
我们的,实践。是。对于。需要。更新。删除。的。场景。比如。订单。状态。更新,用户。信息。更新。等。都。用。MERGE INTO。来。实现。性能。还。不错,比。传统的。先。删除。再。插入。的。方式,高效。很多。而且。是。原子的。不会。出现。中间。状态。
4. 写入,的时候。注意。Schema。的。一致性
虽然,Delta Lake。支持。Schema演进。但是。写入。的时候。还是。要。注意。Schema,的。一致性。避免。因为。Schema。不匹配。导致。写入,失败。或者。产生。意外。的。Schema演进。
Delta Lake,默认。是。严格。Schema。校验。的。如果。写入。的。DataFrame,的。Schema。和。表。的。Schema。不匹配。会。报错。写入,失败。如果。想要。自动。合并。Schema。需要。开启。mergeSchema,选项:
df.write.format("delta").mode("append").option("mergeSchema", "true").save("/path/to/table")或者,开启。全局。参数:
spark.conf.set("spark.databricks.delta.schema.autoMerge.enabled", "true")但是,不建议。全局。开启。自动。Schema合并。因为。这样。可能。会。因为。一些。意外。的,字段。导致。表。的。Schema。被。意外。修改。建议,在。确实。需要。Schema演进。的时候。才。临时。开启。mergeSchema。选项。并且,提前。确认。Schema。的。变化。是。预期,的。
我们的,实践。是。默认。不。开启。自动。Schema合并。写入,的时候。严格。校验。Schema。如果。需要。增加。字段。先。手动。修改,表。的。Schema。用。ALTER TABLE。ADD COLUMN。然后。再。写入。这样。更,安全。更。可控。避免。意外。的。Schema演进。
四、数据读取最佳实践
数据读取。也是。Delta Lake,最。常用的。操作。之一。读取。的。方式。和,参数。直接。影响。查询。的。性能。
1. 尽量,用。分区。列。过滤。减少。扫描。的。数据量
这是,最基本。也。最。有效。的。性能。优化。手段,查询。的时候。尽量。用。分区。列。过滤。比如。按,dt。日期。过滤。这样。能。跳过。不需要的。分区,只。扫描。需要的。分区。的。数据。大大。提高。查询,性能。
比如。不要。这样:
SELECT * FROM table WHERE user_id = '123'要。这样:
SELECT * FROM table WHERE dt = '2021-11-09' AND user_id = '123'加上,分区。列。的。过滤。能。大大。减少。扫描,的。数据量。提高。查询。性能。
2. 只,查询。需要。的。列。不要。SELECT *
查询,的时候。只。查询。需要。的。列。不要。SELECT *。因为。Delta Lake。是。列存,格式。Parquet。只。读取。需要。的。列。能。大大,减少。IO。提高。查询。性能。
比如。不要。这样:
SELECT * FROM table WHERE dt = '2021-11-09'要。这样:
SELECT id, name, age FROM table WHERE dt = '2021-11-09'只,查询。需要。的。列。能。大大。提高。查询。性能。
3. 用,时间旅行。查询。历史。版本。的。数据
Delta Lake,支持。时间旅行。能。查询。历史。版本。的。数据。这是。非常,有用的。功能。能。用来。查看。数据。的。历史,变化。回滚。错误。的。写入。做。数据。审计,等。
查询,历史。版本。的。方式。有。两种。一种。是。按。版本号。查询。一种。是。按。时间,查询。
按,版本号。查询:
SELECT * FROM table VERSION AS OF 10或者:
spark.read.format("delta").option("versionAsOf", 10).load("/path/to/table")按,时间。查询:
SELECT * FROM table TIMESTAMP AS OF '2021-11-09 10:00:00'或者:
spark.read.format("delta").option("timestampAsOf", "2021-11-09 10:00:00").load("/path/to/table")用,时间旅行。的。注意事项:
- 历史,版本。的。保留。时间。是。可配置的。默认。是,30天。超过。保留。时间。的。历史。版本。会,被。VACUUM。删除。就。无法。查询。了。所以。如果。需要。查询。更久,的。历史。需要。调整。日志。保留。时间。
- 查询,历史。版本。的。性能。会。比。查询。最新,版本。差。一些。因为。需要。读取。历史。的。元数据。和。文件。所以。不要,频繁。查询。太久远。的。历史。版本。
- 时间旅行,主要。用于。数据。审计。错误。回滚。历史。数据,对比。等。场景。不要。用于。常规。的。查询。常规。查询,用。最新。版本。就行。
我们的,实践。是。对于。重要。的。表。保留。30天。的。历史,版本。用于。数据。审计。和。错误。回滚。遇到。错误。的,写入。用。时间旅行。查询。历史。版本。然后。恢复。数据,非常。方便。比。传统的。从。备份。恢复。高效,很多。
4. 用,Delta Cache。加速。查询
如果,用的。是。Databricks。的。运行时。支持。Delta Cache。能,缓存。Delta Lake。的。数据。在。本地。磁盘。和,内存。中。加速。查询。性能。对于。重复。查询,相同。数据。的。场景。性能。提升。很,明显。
开启,Delta Cache。很简单。只。需要。设置。参数:
spark.conf.set("spark.databricks.io.cache.enabled", "true")然后,查询。的时候。就。会。自动。缓存。数据。下次。查询。相同,数据。的。时候。就。会。从。缓存。读取,性能。更。快。
如果,用的。是。开源。的,Spark。没有。Delta Cache。也。可以。用。Spark。自带的,cache。或者。persist。来。缓存,DataFrame。但是。需要。手动。管理。缓存,的。生命周期。没有。Delta Cache,方便。
五、性能优化最佳实践
性能优化,是。使用。Delta Lake。的。永恒。话题。这里。分享。一些。常用的,性能。优化。最佳实践。
1. 定期,做。OPTIMIZE。合并。小文件。优化。排序
OPTIMIZE,是。Delta Lake。提供的。优化。命令。能。合并。小文件,优化。排序。提高。查询。性能。是。最。常用。也。最,有效。的。性能。优化。手段。之一。
OPTIMIZE,的。语法:
OPTIMIZE table或者,指定。Z-Order。列:
OPTIMIZE table ZORDER BY (col1, col2)也。可以,指定。分区:
OPTIMIZE table WHERE dt = '2021-11-09'OPTIMIZE,的。作用:
- 合并,小文件。把。小文件。合并。成。大文件。达到。合适的,大小。减少。文件。数量。提高。查询。性能。
- 优化,排序。如果。指定。了。Z-Order。列。会。按照。Z-Order。重新,排序。数据。提高。这些。列。的。过滤。性能。
OPTIMIZE,的。注意事项:
- OPTIMIZE,是。重写。数据。的。操作。会。产生。新的,版本。和。新的。文件。旧的。文件。会。保留。直到,VACUUM。删除。所以。OPTIMIZE。会。增加。存储。成本。需要。定期,做。VACUUM。清理。旧。文件。
- OPTIMIZE,会。消耗。计算。资源。所以。建议。在。业务。低峰。期,做。比如。凌晨。或者。周末。不要。在。业务。高峰。期。做,影响。正常。业务。
- 不要,太。频繁。做。OPTIMIZE。一般。每天。或者。每周。做。一次。就,够了。太。频繁。会。浪费。计算。资源。和。存储,资源。
- OPTIMIZE,的。文件。大小。是。可配置的。默认。是。1GB。可以。根据,自己的。场景。调整。一般。建议。128MB。到。1GB。之间。
我们的,实践。是。每天。凌晨。对。所有。的。Delta表,做。一次。OPTIMIZE。对于。写入。频繁。的。表,做。Z-Order。对于。写入。少。的。表。只。合并。小文件,不。做。Z-ORDER。这样。能。保证。查询。性能。同时,控制。计算。成本。
2. 定期,做。VACUUM。清理。旧。文件。节省。存储
VACUUM,是。Delta Lake。提供的。清理。命令。能。删除。旧的,版本。的。文件。和。不再。使用。的。文件,节省。存储。成本。
VACUUM,的。语法:
VACUUM table或者,指定。保留。时间:
VACUUM table RETAIN 168 HOURSVACUUM,的。作用:
- 删除,旧的。版本。的。文件。超过。保留。时间。的。历史。版本。的。文件。会。被。删除。节省,存储。成本。
- 删除,不再。使用。的。文件。比如。被。OPTIMIZE。合并。的,旧。小文件。被。MERGE。重写。的。旧。文件。等。这些,文件。不再。被。最新。版本。使用。但是。还。占着。存储,VACUUM。会。删除。它们。节省。存储。
VACUUM,的。注意事项:
- VACUUM,删除。的。文件。是。无法。恢复。的。所以。做,VACUUM。之前。一定要。确认。保留。时间。设置。合理。不要,把。还。需要。的。历史。版本。删除。了。
- 默认,的。保留。时间。是。168小时。也就是。7天。Delta Lake。有,安全。检查。不允许。保留。时间。小于。7天。除非。手动,关闭。安全。检查。不建议。这么。做。太。危险。
- VACUUM。不会,删除。事务。日志。也就是。deltalog。目录。的。内容。事务,日志。的。保留。时间。是。单独。配置的。默认,是。30天。超过。的。会。自动。删除。不需要,手动。VACUUM。
- VACUUM,会。消耗。一定的。计算。资源。但是。比。OPTIMIZE。少,很多。建议。定期。做。比如。每周。一次。或者。每月。一次。根据。存储,增长。情况。决定。
我们的,实践。是。每周。对。所有。的。Delta表。做,一次。VACUUM。保留。7天。的。历史。版本。对于,重要。的。表。保留。30天。这样。能。节省,大量的。存储。成本。同时。保证。有。足够的。历史。版本,用于。审计。和。回滚。
3. 调整,Spark。的。参数。优化。性能
除了,Delta Lake。本身。的。优化。调整。Spark。的。参数。也。能,大大。优化。Delta Lake。的。性能。
常用的,Spark。参数。优化:
- spark.sql.shuffle.partitions:shuffle,的。分区。数。默认,是。200。对于。大数据量。的,任务。可以。调大。比如。1000。或者。2000,提高。并行度。对于。小数据量。的。任务。可以,调小。比如。50。或者。100。减少,调度。开销。
- spark.default.parallelism:默认。并行度。和,shuffle.partitions。类似。根据,数据量。调整。
- spark.sql.files.maxPartitionBytes:每个,分区。读取。的。最大,字节数。默认。是。128MB。可以。根据,文件。大小。调整。提高,读取。并行度。
- spark.sql.files.openCostInBytes:打开,一个。文件。的。成本,默认。是。4MB。可以。调整,影响。文件。的。合并。和,分区。
- spark.databricks.delta.optimize.maxFileSize:OPTIMIZE,的。最大,文件。大小,默认。是,1GB。可以。调整。
- spark.databricks.delta.retentionDurationCheck.enabled:VACUUM,保留。时间,检查。默认,是。true,不建议。关闭。
- spark.databricks.delta.schema.autoMerge.enabled:自动,Schema合并。默认,是。false,不建议。开启。
这些,参数。要。根据。自己的。数据量。集群。规模。任务。类型。灵活,调整。没有。万能。的。参数。需要。不断。测试。调优。
4. 用,分区。裁剪。和。谓词。下推。提高。查询。性能
Delta Lake,支持。分区。裁剪。和。谓词。下推。查询。的时候,用。分区。列。过滤。能。触发。分区。裁剪,跳过。不需要的。分区。用。排序列。或者。Z-Order列。过滤。能,触发。谓词。下推。和。数据。跳过。减少。扫描,的。数据量。
所以,查询。的时候。尽量。用。分区。列。排序。列,Z-Order列。过滤。能。大大。提高。查询。性能。
同时,尽量。避免。在。过滤。列。上。用。函数。比如。不要。这样:
SELECT * FROM table WHERE date(dt) = '2021-11-09'要。这样:
SELECT * FROM table WHERE dt = '2021-11-09'因为,在。列。上。用。函数。会。导致。无法。触发。分区。裁剪。和。谓词。下推。影响。查询,性能。
六、运维管理最佳实践
Delta Lake,的。运维。管理。也。很重要。好的。运维。管理,能。保证。Delta Lake。的。稳定。运行。和。性能。
1. 监控,表。的。状态。和。性能
要,定期。监控。Delta表。的。状态。和。性能。包括,表。的。大小。文件。数量。分区。数量。小文件,比例。查询。性能。写入。性能。等。及时。发现。问题。解决,问题。
可以,用。Delta Lake。提供的。工具。来。查看。表。的,状态。比如:
- DESCRIBE HISTORY table:查看,表。的。历史。版本。和,提交。记录。
- DESCRIBE DETAIL table:查看,表。的。详细。信息。包括,位置。格式。分区。列,表。大小。文件。数量。等。
- DESC FORMATTED table:查看,表。的。格式化。信息,更。详细。
也。可以,用。系统。表。来。监控。比如。在。Databricks。里,有。delta_log。的。系统。表。能。查询。表,的。历史。操作。和。性能。指标。
我们的,实践。是。每天。监控。所有。Delta表。的。大小,文件。数量。小文件。比例。查询。性能。写入。性能,等。指标。发现。异常。及时。处理。比如。小文件。太多。就。做。OPTIMIZE。存储,增长。太快。就。做。VACUUM。查询。性能。下降。就。检查。分区。排序。和,参数。
2. 做好,数据。质量。校验
Delta Lake。虽然。支持,Schema校验。但是。还是。要。做好。数据。质量。校验。保证,写入。数据。的。质量。避免。脏数据。进入,数据湖。
可以,在。写入。之前。对。数据。做。校验。比如。非空,校验。格式。校验。范围。校验。一致性。校验。等。对于。不符合,要求。的。数据。过滤。掉。或者。写入。到。异常,表。人工。处理。
也。可以,用。Delta Lake。的。约束。功能。比如。NOT NULL。约束。CHECK,约束。等。在。表。层面。保证。数据。质量。但是,Delta Lake。的。约束。功能。还。在。发展。中。支持。的。还。不是。很,完善。所以。建议。还是。在。写入。之前。做。数据。质量,校验。更。可靠。
我们的,实践。是。在。数据。写入。Delta Lake。之前。都,会。做。数据。质量。校验。包括。Schema校验。非空,校验。格式。校验。范围。校验。等。对于。异常。数据。记录,日志。告警。人工。处理。保证。写入。Delta Lake。的,数据。都是。高质量。的。
3. 做好,备份。和。灾难。恢复
虽然,Delta Lake。本身。有。时间旅行。能。恢复。历史。版本,的。数据。但是。还是。要。做好。备份。和。灾难。恢复。因为。时间旅行。只能,恢复。保留。时间。内。的。历史。版本。超过,保留。时间。的。就。无法。恢复。了。而且。如果,整个。存储。都。损坏。了。时间旅行。也。没用。
所以,要。定期。对。重要。的。Delta表。做。备份。比如,备份。到。另一个。存储。系统。或者。另一个。地域。保证,在。灾难。发生。的。时候。能。恢复,数据。
备份,的。方式。可以。是。全量。备份。也。可以。是。增量,备份。可以。用。Delta Lake。的。CLONE。功能。克隆。表,到。另一个。位置。作为。备份。也。可以。用。DistCp,等。工具。复制。表。的。文件。和。事务。日志。到,另一个。位置。
我们的,实践。是。对于。重要。的。表。每天。做。一次。全量,备份。到。另一个。存储。系统。保留。30天。的,备份。用于。灾难。恢复。同时。Delta Lake。本身。的。时间旅行,保留。7天。用于。日常。的。错误。回滚。这样,双重。保障。数据。安全。
4. 管理,好。事务。日志。和。元数据
Delta Lake,的。事务。日志。也就是。deltalog。目录。是。Delta Lake,的。核心。存储。了。表。的。所有。元数据。和。历史,版本。所以。要。管理。好。事务。日志。保证。事务,日志。的。安全。和。性能。
事务,日志。的。注意事项:
- 事务,日志。的。文件。是。JSON。格式。的。每个,提交。产生。一个。JSON。文件。不要。手动。修改。或者。删除,事务。日志。的。文件。否则。会。导致。表。损坏,无法。读取。
- 事务,日志。的。保留。时间。默认。是。30天。超过,的。会。自动。删除。不需要。手动。管理。如果。需要。保留,更久。的。历史。版本。可以。调整。日志。保留。时间。
- 事务,日志。的。性能。很重要。如果。事务。日志。太大。或者。文件,太多。会。影响。查询。的。启动。性能。因为。查询,的时候。需要。读取。事务。日志。所以。要。定期。做。CHECKPOINT,把。事务。日志。合并。成。Parquet。格式。的,checkpoint。文件。提高。读取。性能。Delta Lake。默认。每,10次。提交。做。一次。checkpoint。不需要。手动。管理。
- 事务,日志。要。和。数据。文件。存储。在。一起。不要。分开,存储。否则。会。导致。表。无法。读取。或者。性能。差。
我们的,实践。是。事务。日志。和。数据。文件。存储。在。同一个,目录。下。不。手动。管理。事务。日志。让,Delta Lake。自动。管理。定期。监控。事务。日志。的,大小。和。文件。数量。如果。有。异常。及时。处理。
七、常见问题,和,解决方法
最后,分享。一些。我们。在。使用。Delta Lake。过程中。遇到的,常见。问题。和。解决。方法。希望。能。帮到。大家。
1. 小文件,问题
问题:表,的。文件。太多。太小。导致。查询。性能。差。元数据。开销。大。
解决:定期,做。OPTIMIZE。合并。小文件;写入。的时候。调整。Spark。参数,控制。文件。大小。避免。产生。太多。小文件;流处理。写入,调整。触发。间隔。不要。太。频繁。
2. 查询,性能。差
问题:查询,慢。性能。差。
解决:检查,是否。用。分区。列。过滤;检查。是否。只。查询。需要,的。列;检查。表。是否。做了。OPTIMIZE。和。Z-Order;检查。Spark。参数,是否。合理;检查。文件。大小。是否。合适。有没有。小文件,问题。
3. 写入,性能。差
问题:写入,慢。性能。差。
解决:检查,Spark。参数。是否。合理。比如。shuffle.partitions。并行度。等;检查。是否,有。太多。的。小文件。写入。导致。元数据。开销,大;检查。是否。做了。太多。的。排序。或者。Z-Order。增加,写入。成本;检查。存储。的。IO。性能。是否,足够。
4. MERGE INTO,性能。差
问题:MERGE INTO,慢。性能。差。
解决:MERGE,的。条件。尽量。用。分区。列。和。排序。列。减少,扫描。的。数据量;源表。数据量。太大。的话。先。过滤。或者,分批次。MERGE;避免。大范围。的。更新。删除。尽量。只,更新。需要。的。分区;定期。做。OPTIMIZE。保证。文件。大小,合适。
5. Schema演进,问题
问题:写入,的时候。Schema。不匹配。报错。或者,意外。的。Schema演进。
解决:写入,之前。确认。Schema。的。一致性;需要。增加。字段。先,手动。ALTER TABLE。ADD COLUMN。再。写入;不要。全局。开启。自动,Schema合并。避免。意外。的。Schema演进;定期。检查。表。的,Schema。是否。符合。预期。
6. 时间旅行,查询。不到。历史。版本
问题:查询,历史。版本。的时候。报错。说。版本。不存在。或者。文件,不存在。
解决:检查,历史。版本。是否。超过。了。保留。时间。被,VACUUM。删除。了;检查。事务。日志。是否。完整。有没有,被。手动。删除。或者。损坏;检查。存储。的。文件。是否,完整。有没有。被。误删。
7. 存储,成本。高
问题:表,的。存储。太大。成本。高。
解决:定期,做。VACUUM。清理。旧。文件。和。历史。版本;调整,历史。版本。的。保留。时间。不要。保留。太久;检查。是否,有。不必要。的。数据。可以。清理。或者。归档。到。更。便宜。的,存储;压缩。数据。Parquet。本身。就是。压缩。的。可以。调整。压缩。算法,提高。压缩率。
八、写在最后
以上。就是。我,总结。的。Delta Lake。的。最佳实践。包括。表设计、数据写入、数据读取、性能优化、运维管理、常见问题。等,方面。都是。我们。在。实际。项目。中。踩坑,总结。出来的。经验。希望。能。给。正在。使用。或者,打算。使用Delta Lake。的。朋友。一些。参考。
Delta Lake,是。一个。非常。优秀。的。存储层。框架。解决了,传统。数据湖。的。很多。痛点。比如。ACID事务。更新。删除,Schema演进。时间旅行。等。让。数据湖。也能。像。数据仓库,一样。好用。是。构建。Lakehouse。架构。的。重要,基础。最近几年。发展。很快。越来越。成熟。越来越。多的,公司。开始。使用。
但是,Delta Lake。也。不是。银弹。不是。用了。就。万事大吉。了。还是。需要。好的。表,设计。好的。写入。读取。方式。好的。性能。优化,好的。运维。管理。才能。发挥。Delta Lake。的。最大,价值。否则。也。会。遇到。很多。性能。问题。和,运维。问题。
而且,Delta Lake。发展。很快。版本。更新。也。很快。新的,功能。不断。推出。旧的。问题。不断。修复。所以。大家,要。关注。Delta Lake。的。最新。发展。学习。新的,功能。和。最佳实践。不断。优化。自己的。使用,方式。
最后,希望。这篇文章。能。帮到。大家。也。欢迎。大家,交流。讨论。分享。你们。的。Delta Lake。使用。经验。和。最佳实践,一起。学习。一起。进步。
如果你。也。在,使用。Delta Lake。或者。打算。使用。有。什么。问题。或者。经验,欢迎。在。评论区。留言。交流。
祝,大家。都。能。用好,Delta Lake。构建。高性能。高可靠,的。数据湖。和。Lakehouse。架构。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录