Delta Lake 快速入门:30分钟搭出带事务的数据湖
2026/9/19 19:23:41 网站建设 项目流程

Delta Lake 快速入门:30分钟搭出带事务的数据湖

【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/delta

Delta Lake 是一个开源存储框架:它在数据湖的普通文件(Parquet)之上铺一层"事务日志",让数据湖具备 ACID 事务(原子性、一致性、隔离性、持久性)、时间旅行和模式校验能力。读完这篇,你能在本地建出第一张 Delta 表并把数据读回来,还会搞懂三个生产环境必碰的核心机制:事务日志、优化写入、流式事件时间排序。

先说痛点

想象一个常见场景:数据以 Parquet 文件形式放在 S3 或 HDFS 上,20 个 ETL 任务共用一个目录。

某天两个任务同时写:一个写完了,另一个写到一半,目录里留下半截文件,下一次读取直接读到"脏"数据。更惨的是凌晨任务把数据覆盖错了,白天跑批才发现结果不对——想回滚?晚了。

你可以每次读取前检查一遍文件列表,但那只是降低出错的概率,挡不住出错。基于文件的数据湖缺三样东西:原子提交、历史版本、模式校验。Delta Lake 补的就是这三样。

它到底是什么、凭什么能用

一句话定位:Delta Lake = 文件 + 事务日志

每次写入都会在_delta_log目录追加一个小 JSON 文件,记录"这次新增了哪些数据文件、删掉了哪些"。读取时不是直接扫目录,而是先回放日志、重建当前快照。ACID 和时间旅行,就是从这套日志机制里长出来的。

图中是三个角色和五步数据流:Spark Driver是接 SQL 的"管家",只管发请求收结果;Delta Kernel Connector是"翻译官",把 Driver 的 schema 请求和过滤器(静态+动态)下推给内核层,并取回要扫描的文件清单;Delta Kernel是"懂 Delta 日志"的核心,负责解析日志、判定该读哪些文件;最后真正的数据扫描由 Spark 自带的 Parquet reader 完成。这样内核只处理日志逻辑,数据走原有高性能读取器,两头的好处都拿到了。

从0到1跑通 🧪

环境要求:3条清单

  • Java 8 / 11 / 17 任一(java -version确认)
  • 内存 4GB 以上即可,本地模式能跑通
  • Python 3.9+(用下面示例时需要)

安装:一条命令

pip install delta-spark==4.0.0

想从源码构建的话,克隆 https://gitcode.com/GitHub_Trending/del/delta 后执行build/sbt package即可;初次跑通用现成包更快。

最小闭环:一段脚本搞定

脚本干三件事:建会话、写一张表、读回来。

import shutil from pyspark.sql import SparkSession from delta import configure_spark_with_delta_pip shutil.rmtree("/tmp/delta-table", ignore_errors=True) # 清掉上次运行残留 spark = configure_spark_with_delta_pip( SparkSession.builder.appName("quickstart").master("local[*]") ).getOrCreate() # 1. 建表:DataFrame 以 delta 格式写出 spark.range(0, 5).write.format("delta").save("/tmp/delta-table") # 2. 读取:按路径读回,无需额外配置 spark.read.format("delta").load("/tmp/delta-table").show()

configure_spark_with_delta_pip会把 Delta 包自动装配进 Spark,不用手写任何配置。

怎么判断跑通了

三个可验证的信号:

  1. 控制台打印出 0~4 共 5 行数据;
  2. /tmp/delta-table目录里既有 Parquet 数据文件,也有_delta_log目录——日志目录是 Delta 表的"身份证";
  3. 日志目录里能看到00000000000000000000.json,这就是第 0 版的提交记录。

往深挖:3个值得了解的特性 🧩

小文件问题怎么解:优化写入

流式作业和批处理每几分钟写一小批,跑几天表里就是几十万个碎文件,查询越来越慢。

左边是传统写入:多个 executor 各写各的,往分区目录里堆小文件,越积越多;右边 Optimized Writes 先把同一分区的小文件归拢重写,落成少量大文件。文件数下来了,查询性能回升;再配合 VACUUM(清理过期旧版本文件)把存储成本也控住。

流式事件时间排序:迟到数据不丢

流数据经常"乱序":事件发生在 10:00,数据 10:05 才到系统,处理不处理?

最上一行是初始快照,3 个文件里的记录按事件时间有序;中间一行关闭了事件时间排序,Batch 2 里晚到的记录"2"被当成迟到事件直接丢弃(红块);最下一行开启排序后,引擎按事件时间重排,"2" 被正确放到"3"的后面。生产上建议开启排序,让迟到数据落在对的位置,而不是被悄悄扔掉。

时间旅行与 ACID

每次提交都记一个版本号。表被覆盖之后,用VERSION AS OF 0还能读回旧版本——审计、回滚、可复现的模型训练都靠它。两个并发写也不会互相踩坏:提交走乐观锁,后提交者发现日志变了就自动重试,隔离级别是可串行化的。

谁在用、用来干嘛

  • 批处理 ETL 团队:替换 Hive/数仓表,让湖上的数据也能 UPDATE、DELETE、MERGE(有则更新、无则插入),不再只会追加。
  • 流式团队:Kafka 写进同一张 Delta 表,批式回填和实时查询共用一份数据,不用维护两条链路。
  • 分析(BI)团队:PrestoDB、Trino、Hive 直接查 Delta 表,不换引擎、不搬数据。
  • 算法团队:用时间旅行版本锁定数据快照,复现实验时按同一份数据重跑训练。

Delta Lake 适合已经在用数据湖、或正准备搭"湖仓一体"的团队:你要的是湖的成本加上仓库的一致性,它正好卡在中间。如果看到这里,不妨直接在自己机器上把上面那几条命令敲一遍——十分钟亲眼看到_delta_log目录,比看十篇文章都踏实。

延伸阅读:官方文档首页、快速开始指南、事务协议规范;核心源码在 spark/src/main/scala/org/apache/spark/sql/delta/,内核实现在 kernel/kernel-api/。

【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/delta

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询