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

Databricks数据倾斜问题:三种方案无效,求高效优化方案

倾斜数据Join优化问题及解决方案建议

我创建了倾斜数据以测试加盐(Salting)方案,尝试三种不同解决方法但均未实现显著运行时长优化,集群配置为32GB内存、4核,共8个worker节点。

测试代码

import pyspark.sql.functions as F
# 修正原代码语法错误:drip→drop,补全缺失引号
df1 = spark.range(300_000_000).withColumn(
    'value',
    F.when(F.rand() < 0.6, 1).otherwise((F.rand() * 100).cast('int'))
).drop('id')
df2 = spark.range(200_000_000).withColumn(
    'value',
    F.when(F.rand() < 0.2, 4).otherwise((F.rand() * 100).cast('int'))
).drop('id')
final_df = df1.join(df2, on='value', how='inner')
final_df.write.format('parquet').save(path)

已尝试的三种方案

  • 方案一:启用AQE及skewjoin.enabled=true、coalescePartitions.enabled=True、shuffle.partitions=auto配置,作业运行超20分钟后手动终止。
  • 方案二:采用加盐技术,关闭AQE并设置shuffle.partitions=1000,修改关联键为['value','salt'],作业运行超30分钟后手动终止,代码如下:
# 修正原代码拼写错误:withColum→withColumn
df1 = (df1.withColumn('salt_numbers', F.expr('sequence(0,1)'))
       .withColumn('salt', F.explode('salt_numbers'))
       .drop('salt_numbers'))
df2 = (df2.withColumn('salt_numbers', F.expr('sequence(0,3)'))
       .withColumn('salt', F.explode('salt_numbers'))
       .drop('salt_numbers'))
  • 方案三:采用与方案二相同的加盐代码并启用AQE,作业运行时长相近后手动终止。

优化建议

1. 精准调整加盐策略

原加盐方案对全量数据统一加盐,未针对热点value做差异化处理,导致热点数据拆分后仍存在不均衡问题:

  • 先统计两张表中各value的出现频次,定位热点值(如df1中value=1占60%、df2中value=4占20%);
  • 仅对热点值进行多维度加盐,非热点值保持原key,避免不必要的数据膨胀;
  • 确保同一热点值在两张表中的盐值范围一致,保证join时数据均匀分布。

示例代码:

# 统计热点value
df1_freq = df1.groupBy('value').count().orderBy(F.desc('count')).limit(5).collect()
hot_values_df1 = {row['value'] for row in df1_freq if row['count'] > 10_000_000}
df2_freq = df2.groupBy('value').count().orderBy(F.desc('count')).limit(5).collect()
hot_values_df2 = {row['value'] for row in df2_freq if row['count'] > 10_000_000}
hot_values = hot_values_df1 & hot_values_df2

# 对热点value加盐,非热点值盐值固定为0
df1_salted = df1.withColumn(
    'salt',
    F.when(F.col('value').isin(hot_values), F.floor(F.rand() * 10)).otherwise(F.lit(0))
)
df2_salted = df2.withColumn(
    'salt',
    F.when(F.col('value').isin(hot_values), F.floor(F.rand() * 10)).otherwise(F.lit(0))
)

# 基于加盐后的key关联,最后移除盐字段
final_df = df1_salted.join(df2_salted, on=['value', 'salt'], how='inner').drop('salt')

2. 优化Spark集群配置

  • 调整shuffle分区数:总worker核数为8×4=32,建议设置spark.sql.shuffle.partitions=96(核数的3倍),平衡任务调度开销与单分区数据量;
  • AQE精细化配置:
    spark.sql.adaptive.enabled=true
    spark.sql.adaptive.skewJoin.enabled=true
    spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=128MB  # 单分区超过该值判定为倾斜
    spark.sql.adaptive.advisoryPartitionSizeInBytes=64MB
    spark.sql.shuffle.partitions=96
    
  • 内存分配优化:每个worker分配spark.executor.memory=24GB,spark.executor.cores=4,spark.driver.memory=8GB,预留部分内存给系统缓存。

3. 数据预处理与执行优化

  • 提前过滤无效数据:若业务允许,先过滤掉join后无匹配的value,减少参与计算的数据量;
  • 缓存热点数据集:对热点value对应的子数据集进行缓存(如df1.filter(F.col('value')==1).cache()),避免重复计算;
  • 尝试拆分join逻辑:将热点value与非热点value分开join,再合并结果,减少单任务处理压力。

4. 尝试Broadcast Join(针对热点子集)

若某张表的热点value子集数据量较小(如df2中value=4的4000万条数据),可将该子集广播后与另一表的对应热点数据join,避免shuffle开销:

# 拆分热点与非热点数据
df1_hot = df1.filter(F.col('value').isin(hot_values))
df1_non_hot = df1.filter(~F.col('value').isin(hot_values))
df2_hot = df2.filter(F.col('value').isin(hot_values))
df2_non_hot = df2.filter(~F.col('value').isin(hot_values))

# 热点数据用broadcast join,非热点数据用普通join
join_hot = df1_hot.join(F.broadcast(df2_hot), on='value', how='inner')
join_non_hot = df1_non_hot.join(df2_non_hot, on='value', how='inner')

# 合并结果
final_df = join_hot.union(join_non_hot)

内容的提问来源于stack exchange,提问作者Learn Hadoop

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 07:32:17