You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 10:02:31