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
相关产品推荐
相关产品推荐

