Delta表变更引发现有Spark DataFrame数据丢失问题排查
问题原因与解决方案
核心问题根源
Spark采用惰性求值机制:你定义的df3并没有立即计算出结果,只是生成了逻辑执行计划。虽然df3.count()触发了一次计算,此时数据正常,但执行DELETE操作删除源表foo中的目标记录后,df3.write会重新执行整个df3的生成逻辑——再次从foo表读取数据,而此时需要转换的记录已经被删除,导致df3变为空。
解决方案
提前将转换后的df3数据物化(持久化),避免后续写入时重新读取已被修改的源表。以下是两种可行方案:
方案1:缓存DataFrame
通过缓存将df3的数据持久化到内存/磁盘,后续写入时直接使用缓存数据,不再重新查询源表:
from pyspark.sql import functions as F df1 = spark.table("foo") \ .filter(F.col("baz") == "A") df2 = df1.filter(F.expr("position('|', discount)") > 1) df3 = df2 \ .select(F.explode(F.split("discount", "\\|")).alias("discount_denorm"), "*") \ .withColumn("new_discount", 1 - F.get(F.regexp_extract_all(F.col("discount_denorm"), r'(\d+)%', 1), 0) / 100) \ .select( "id", F.col("new_discount").alias("discount"), "baz", F.col("price") * F.col("new_discount").cast("decimal(10,2)").alias("price") ) # 缓存df3并触发计算,确保数据被持久化 df3.cache() df3.count() # 触发缓存加载 df2.select("id").createOrReplaceTempView("ids_to_delete") # 删除源表中的目标记录 spark.sql(""" DELETE FROM foo WHERE EXISTS (SELECT 1 FROM ids_to_delete WHERE foo.id = ids_to_delete.id) """) # 写入缓存中的数据,不会重新读取源表 df3.write.format("delta").mode("append").saveAsTable("foo") # 操作完成后释放缓存,避免占用资源 df3.unpersist()
方案2:写入临时Delta表
将转换后的df3写入临时表,后续从临时表读取数据写入目标表:
from pyspark.sql import functions as F df1 = spark.table("foo") \ .filter(F.col("baz") == "A") df2 = df1.filter(F.expr("position('|', discount)") > 1) df3 = df2 \ .select(F.explode(F.split("discount", "\\|")).alias("discount_denorm"), "*") \ .withColumn("new_discount", 1 - F.get(F.regexp_extract_all(F.col("discount_denorm"), r'(\d+)%', 1), 0) / 100) \ .select( "id", F.col("new_discount").alias("discount"), "baz", F.col("price") * F.col("new_discount").cast("decimal(10,2)").alias("price") ) # 将转换后的数据写入临时Delta表 df3.write.format("delta").mode("overwrite").saveAsTable("temp_transformed_data") df2.select("id").createOrReplaceTempView("ids_to_delete") # 删除源表中的目标记录 spark.sql(""" DELETE FROM foo WHERE EXISTS (SELECT 1 FROM ids_to_delete WHERE foo.id = ids_to_delete.id) """) # 从临时表读取数据写入目标表 spark.table("temp_transformed_data").write.format("delta").mode("append").saveAsTable("foo") # 可选:清理临时表 spark.sql("DROP TABLE temp_transformed_data")
注意事项
- 在Azure Databricks Serverless模式下,缓存的生命周期与当前会话绑定,确保整个操作流程在同一会话内完成
- 临时表默认仅在当前会话可见,若需跨会话使用可改用全局临时表(
createOrReplaceGlobalTempView) - 若数据量较大,推荐使用
persist()指定存储级别(如StorageLevel.DISK_ONLY),避免内存不足
内容的提问来源于stack exchange,提问作者Adam
相关产品推荐
相关产品推荐

