Spark SQL 中的数据倾斜是分布式计算中的常见问题,指数据在分区之间分布不均匀,导致某些任务处理时间远超其他任务。以下是详细的处理策略和优化方法。

数据倾斜的识别和诊断

症状识别

-- 通过Spark UI观察任务执行时间
-- 倾斜特征:少数任务执行时间极长,多数任务快速完成

-- 检查数据分布
SELECT key, COUNT(*) as count 
FROM table 
GROUP BY key 
ORDER BY count DESC 
LIMIT 10;

监控指标

  • 任务执行时间差异: 最大/最小任务时间比 > 3:1
  • 数据量差异: 最大分区数据量 > 平均值的2-3倍
  • Shuffle读写量: 个别executor的shuffle量异常高

数据倾斜的优化策略

1. 聚合操作倾斜优化

两阶段聚合(Local-Global Aggregation)
-- 原始查询(可能倾斜)
SELECT user_id, COUNT(*) as action_count
FROM user_actions 
GROUP BY user_id;

-- 优化:添加随机前缀进行局部聚合
SELECT 
    user_id,
    SUM(partial_count) as total_count
FROM (
    SELECT 
        CONCAT(user_id, '_', CAST(rand()*10 as int)) as salted_key,
        COUNT(*) as partial_count
    FROM user_actions 
    GROUP BY CONCAT(user_id, '_', CAST(rand()*10 as int))
) tmp
GROUP BY user_id;
DataFrame API实现
import org.apache.spark.sql.functions._

// 第一阶段:添加随机前缀局部聚合
val localAgg = df
  .withColumn("salted_key", concat($"user_id", lit("_"), (rand() * 10).cast("int")))
  .groupBy("salted_key")
  .agg(count("*").as("partial_count"))

// 第二阶段:去除前缀全局聚合
val globalAgg = localAgg
  .withColumn("user_id", split($"salted_key", "_")(0))
  .groupBy("user_id")
  .agg(sum("partial_count").as("total_count"))

2. Join操作倾斜优化

广播连接(Broadcast Join)
-- 小表广播避免shuffle
SELECT /*+ BROADCAST(small_table) */ *
FROM large_table 
JOIN small_table ON large_table.key = small_table.key;
倾斜键分离处理
-- 分离倾斜键和非倾斜键分别处理
WITH skewed_keys AS (
    -- 识别倾斜键(如null或特定值)
    SELECT key FROM table WHERE key IS NULL OR key = 'hot_value'
),
non_skewed_data AS (
    SELECT * FROM large_table 
    WHERE key NOT IN (SELECT key FROM skewed_keys)
),
skewed_data AS (
    SELECT * FROM large_table 
    WHERE key IN (SELECT key FROM skewed_keys)
)

-- 非倾斜数据正常join
SELECT * FROM non_skewed_data JOIN small_table USING (key)

UNION ALL

-- 倾斜数据添加随机后缀扩大并行度
SELECT * FROM (
    SELECT *, CONCAT(key, '_', CAST(rand()*100 as int)) as expanded_key 
    FROM skewed_data
) expanded 
JOIN (
    SELECT *, CONCAT(key, '_', CAST(rand()*100 as int)) as expanded_key 
    FROM small_table 
    WHERE key IN (SELECT key FROM skewed_keys)
) expanded_small 
ON expanded.expanded_key = expanded_small.expanded_key;

3. 分区策略优化

自定义分区器
import org.apache.spark.Partitioner

class SkewAwarePartitioner(numPartitions: Int, hotKeys: Set[String]) 
  extends Partitioner {
  
  override def numPartitions: Int = numPartitions
  
  override def getPartition(key: Any): Int = {
    val keyStr = key.toString
    if (hotKeys.contains(keyStr)) {
      // 热键分散到多个分区
      (keyStr.hashCode.abs + System.currentTimeMillis().toInt) % numPartitions
    } else {
      // 普通键均匀分布
      keyStr.hashCode.abs % numPartitions
    }
  }
}

// 使用自定义分区器
val partitionedRDD = rdd.partitionBy(new SkewAwarePartitioner(200, hotKeys))
增加分区数
// 增加shuffle分区数
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
spark.conf.set("spark.sql.shuffle.partitions", "200")  // 默认200,可增加到500-1000

4. 自适应查询执行(AQE)

Spark 3.0+ 的AQE自动处理倾斜:

# 启用AQE
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.skewJoin.enabled=true

# 倾斜join优化参数
spark.sql.adaptive.skewJoin.skewedPartitionFactor=5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256MB

5. 数据预处理策略

过滤异常值
-- 过滤导致倾斜的异常值单独处理
WITH normal_data AS (
    SELECT * FROM source_table 
    WHERE key NOT IN ('异常值1', '异常值2', ...)
),
skewed_data AS (
    SELECT * FROM source_table 
    WHERE key IN ('异常值1', '异常值2', ...)
)

-- 分别处理后再合并
数据采样和重分布
// 采样分析数据分布
val sample = df.sample(0.1).cache()
val keyDistribution = sample.groupBy("key").count().orderBy(desc("count"))

// 基于采样结果动态调整分区策略
val hotKeys = keyDistribution.filter($"count" > threshold).select("key").collect()

具体场景解决方案

场景1:用户行为分析中的热门用户

// 识别热门用户
val hotUsers = df.groupBy("user_id").count()
  .filter($"count" > 10000)  // 定义热门阈值
  .select("user_id")
  .collect()
  .map(_.getString(0))
  .toSet

// 分离处理
val normalData = df.filter(!$"user_id".isin(hotUsers.toSeq: _*))
val hotData = df.filter($"user_id".isin(hotUsers.toSeq: _*))

// 热门用户数据添加随机后缀
val expandedHotData = hotData
  .withColumn("salted_key", concat($"user_id", lit("_"), (rand() * 100).cast("int")))

// 分别聚合后合并

场景2:电商订单关联用户信息

-- 大表(订单)关联小表(用户),但某些用户订单量巨大
WITH skewed_users AS (
    SELECT user_id FROM orders 
    GROUP BY user_id HAVING COUNT(*) > 10000
),
normal_orders AS (
    SELECT o.* FROM orders o 
    LEFT JOIN skewed_users s ON o.user_id = s.user_id 
    WHERE s.user_id IS NULL
),
skewed_orders AS (
    SELECT o.* FROM orders o 
    JOIN skewed_users s ON o.user_id = s.user_id
)

SELECT * FROM normal_orders JOIN users USING (user_id)
UNION ALL
SELECT * FROM (
    SELECT *, CONCAT(user_id, '_', CAST(rand()*10 as int)) as expanded_key 
    FROM skewed_orders
) o JOIN (
    SELECT *, CONCAT(user_id, '_', CAST(rand()*10 as int)) as expanded_key 
    FROM users WHERE user_id IN (SELECT user_id FROM skewed_users)
) u ON o.expanded_key = u.expanded_key;

监控和调优工具

Spark UI分析

  • Stages标签页: 查看任务执行时间分布
  • Storage标签页: 检查缓存使用情况
  • SQL标签页: 分析查询计划和执行统计

性能监控配置

# 详细日志记录
spark.sql.adaptive.logLevel=DEBUG
spark.sql.execution.sort.spill.numElementsForceSpillThreshold=1000000

# 内存调优
spark.memory.fraction=0.6
spark.memory.storageFraction=0.5

最佳实践总结

  1. 预防为主: 在设计阶段考虑数据分布
  2. 监控预警: 建立倾斜检测机制
  3. 分层处理: 分离倾斜数据特殊处理
  4. 动态调整: 利用AQE自动优化
  5. 资源预留: 为可能的倾斜预留额外资源

通过综合运用这些策略,可以有效解决Spark SQL中的数据倾斜问题,提升作业执行的稳定性和效率。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐