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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 16:53:19