Delta Lake(OSS)Merge操作无法完成或耗时过长问题排查
问题描述
- 存储方案:使用Delta Lake作为每日更新表格的主存储,通过
MERGE INTO实现Upsert操作,目标表为Delta格式,更新数据以Parquet文件存储。 - 表特征:包含10000列,多数为二进制类型,使用字符串哈希的行标识列作为Merge匹配条件。
- 测试异常:主表仅5000行(单10MB Parquet文件),但
MERGE INTO操作极慢甚至卡住,未做表分区处理。 - 环境信息:Spark 3.5.1、Delta Lake(delta-spark)3.1.0、EC2实例r5n.16xlarge;执行后输出
DAGScheduler广播大任务二进制的警告,随后Spark会话挂起。 - 复现代码:
import pyspark.sql.functions as F from pyspark.sql import SparkSession from delta import * def get_spark(): builder = SparkSession.builder.master("local[4]").appName('SparkDelta') \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .config("spark.driver.memory", "32g") \ .config("spark.jars.packages", "io.delta:delta-spark_2.12:3.1.0," "io.delta:delta-storage:3.1.0") \ spark = builder.getOrCreate() return spark spark = get_spark() ### GENERATE SOME FAKE DATA import pandas as pd import numpy as np num_rows = 5000 num_cols = 10000 array = np.random.rand(num_rows, num_cols) df = pd.DataFrame(array) df = df.reset_index() df.to_csv('features.csv') ### DATA TO DELTA LAKE TABLE df_spark = spark.read.format('csv').load('features.csv') df_spark.write.format('delta').mode("overwrite").save('deltalake_features') # Creating data to merge - just first 100 rows df_slice = df.iloc[:100, 0:2] df_slice.to_csv('features_slice_100.csv') df_spark_slice = spark.read.format('csv').load('features_slice_100.csv') ### MERGE INTO code target_delta_table = DeltaTable.forPath(spark, 'deltalake_features') df_updates = df_spark_slice df_target = target_delta_table.toDF() ( target_delta_table.alias('target') .merge( df_updates.alias('updates'), ( (F.col('target._c1') == F.col('updates._c1')) ) ) .whenMatchedUpdate(set={ '_c2': F.lit(42) }) .execute() )
解决方案
1. 削减广播数据量(核心优化)
10000列的表在Merge时,默认广播更新数据集会导致数据量过载,触发警告并卡住:
- 只保留更新数据的匹配列(行标识)和需要更新的列,避免加载全量列:
df_updates = df_spark_slice.select(F.col("_c1"), F.col("_c2")) - 强制使用Shuffle Join替代广播,关闭自动广播:
df_updates = df_updates.hint("SHUFFLE_HASH") - 若必须广播,调大允许的广播阈值(默认10MB),在SparkSession配置中添加:
.config("spark.sql.autoBroadcastJoinThreshold", "100m")
2. 优化Delta表操作
- 删除冗余代码:无需调用
target_delta_table.toDF()加载全量目标表到内存,直接用DeltaTable对象执行Merge即可。 - 开启数据跳过与索引优化:
# 写入Delta表时开启数据跳过 df_spark.write.format('delta').option("dataSkippingNumIndexedCols", "1").mode("overwrite").save('deltalake_features') # 对匹配列创建Z-Order索引,加速匹配查询 target_delta_table.optimize().executeZOrderBy("_c1")
3. 调整Spark资源配置
- 利用实例全部CPU:将
master("local[4]")改为master("local[*]"),充分使用r5n.16xlarge的64核CPU。 - 集群模式下补充Executor配置(若使用集群):
.config("spark.executor.instances", "8") .config("spark.executor.cores", "8") .config("spark.executor.memory", "64g")
4. 优化数据生成与读取流程
避免使用CSV作为中间格式(10000列的CSV读写效率极低),直接从Pandas DataFrame写入Spark:
# 跳过CSV写入,直接将Pandas数据转为Spark DataFrame df_spark = spark.createDataFrame(df) df_spark.write.format('delta').mode("overwrite").save('deltalake_features')
优化后完整代码示例
import pyspark.sql.functions as F from pyspark.sql import SparkSession from delta import * def get_spark(): builder = SparkSession.builder.master("local[*]").appName('SparkDelta') \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .config("spark.driver.memory", "32g") \ .config("spark.sql.autoBroadcastJoinThreshold", "100m") \ .config("spark.jars.packages", "io.delta:delta-spark_2.12:3.1.0," "io.delta:delta-storage:3.1.0") \ spark = builder.getOrCreate() return spark spark = get_spark() ### GENERATE SOME FAKE DATA import pandas as pd import numpy as np num_rows = 5000 num_cols = 10000 array = np.random.rand(num_rows, num_cols) df = pd.DataFrame(array) df = df.reset_index() ### DATA TO DELTA LAKE TABLE df_spark = spark.createDataFrame(df) df_spark.write.format('delta') \ .option("dataSkippingNumIndexedCols", "1") \ .mode("overwrite") \ .save('deltalake_features') # 创建更新数据,仅保留需要的列 df_slice = df.iloc[:100, 0:2] df_spark_slice = spark.createDataFrame(df_slice) ### MERGE INTO code target_delta_table = DeltaTable.forPath(spark, 'deltalake_features') # 筛选必要列+禁用广播 df_updates = df_spark_slice.select(F.col("_1"), F.col("_0")).hint("SHUFFLE_HASH") ( target_delta_table.alias('target') .merge( df_updates.alias('updates'), (F.col('target._1') == F.col('updates._1')) ) .whenMatchedUpdate(set={ '_2': F.lit(42) }) .execute() ) # 可选:对匹配列做Z-Order优化 target_delta_table.optimize().executeZOrderBy("_1")
内容的提问来源于stack exchange,提问作者Anton Kretov
相关产品推荐
相关产品推荐

