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

Spark任务最后阶段卡顿求助:多Join与Shuffle异常排查

问题描述

在Databricks上运行Spark作业时,多数阶段执行正常,但某阶段的最后一个任务卡顿。该任务的shuffle读取/处理行数远高于其他任务,尝试多种repartition策略保证数据均匀分布后问题仍存在。

数据背景

  • Table 1(事件日志表):总大小218.2 TiB,已按date分区过滤最近4天数据,但数据量仍较大;
  • Table 2(反向查询哈希表):大小1793.9 GiB,仅包含Key、hash、timestamp、type四列。

业务场景

需对事件日志中的4个哈希字段(adId、ip、ua、uuid)分别与lookup表做left outer join,该操作无法避免但开销极大。曾尝试基于哈希键做repartition,期望相同哈希键进入同一分区由同一executor处理,但仍有单个任务因shuffle read量过大长时间运行。

疑问

  1. 是否误用了repartitioning?
  2. 如何实现迭代广播(ITERATIVE broadcasting),将lookup表拆分为小于8GB的块多次广播后合并?

代码片段

Union操作代码

allIncrementalEvents.as("e")
  .filter(col("e.type") === "authentication")
  .filter(lower(col("e.payload.type")).isin(eventConf.eventTypes:_*))
  .filter(lower(col("e.payload.os.name")).isin(eventConf.osNames:_*))
  .filter(lower(col("e.payload.device.manufacturer")).isin(eventConf.manufacturers:_*))
  .repartition(partitions)

UNION

allIncrementalEvents.as("e")
  .filter(col("e.type") === "session")
  .filter(lower(col("e.payload.type")).isin(eventConf.eventTypes:_*))
  .filter(lower(col("e.payload.os.name")).isin(eventConf.osNames:_*))
  .filter(lower(col("e.payload.device.manufacturer")).isin(eventConf.manufacturers:_*))
  .repartition(partitions)   

UNION

allIncrementalEvents.as("e")
  .filter(col("e.type") === "other")
  .filter(lower(col("e.payload.type")).isin(eventConf.eventTypes:_*))
  .filter(lower(col("e.payload.os.name")).isin(eventConf.osNames:_*))
  .filter(lower(col("e.payload.device.manufacturer")).isin(eventConf.manufacturers:_*))
  .repartition(partitions)

Join操作代码

extractAuthEvents
  .union(extractSubEvents)
  .union(extractOpenEvents)
  .union(extractSessionEvents)
  .join(reverseLookupTableDf.as("adId"),
    col("adId") === col("adId.hashed"),
    "leftouter"
  )
  .join(reverseLookupTableDf.as("ip"),
    col("ae.ip") === col("ip.hashed"),
    "leftouter"
  )
  .join(reverseLookupTableDf.as("ua"),
    col("ae.ua") === col("ua.hashed"),
    "leftouter"
  )
  .join(reverseLookupTableDf.as("uid"),
    col("ae.uuid") === col("uid.hashed"),
    "leftouter"
  )

解决方案

1. repartition误用分析及优化

你当前的repartition操作有两个核心问题:

  • 重复repartition+union:每个分支单独repartition后再union,会导致分区数翻倍(3个分支各partitions个分区,union后总分区数是3*partitions),后续join时需要额外shuffle,反而加剧数据倾斜。应该先union所有过滤后的数据集,再统一做一次repartition;
  • 未指定分区键:当前repartition只指定了分区数,没有基于join用到的哈希字段(adId、ip、ua、uuid)分区,无法让相同哈希值的数据落到同一分区,等于没针对join做优化。

优化后的Union代码:

// 封装过滤逻辑,避免重复代码
def filterEvents(df: DataFrame, eventType: String): DataFrame = {
  df.as("e")
    .filter(col("e.type") === eventType)
    .filter(lower(col("e.payload.type")).isin(eventConf.eventTypes:_*))
    .filter(lower(col("e.payload.os.name")).isin(eventConf.osNames:_*))
    .filter(lower(col("e.payload.device.manufacturer")).isin(eventConf.manufacturers:_*))
}

// 先union再统一repartition,基于join用到的哈希字段分区
val filteredEvents = filterEvents(allIncrementalEvents, "authentication")
  .union(filterEvents(allIncrementalEvents, "session"))
  .union(filterEvents(allIncrementalEvents, "other"))
  .repartition(partitions, col("adId"), col("ip"), col("ua"), col("uuid"))

2. 迭代广播(Iterative Broadcasting)实现

lookup表1.7TB无法直接广播(Spark默认广播阈值10MB,即使调大到8GB也放不下),迭代广播的核心是将lookup表拆分为多个小批次,每个批次广播后与主数据集join,最后合并结果。具体步骤:

步骤1:拆分lookup表

基于hash字段的哈希值拆分lookup表,确保每个子表大小在8GB以内(可根据集群内存调整):

// 计算拆分批次数量,1.7TB≈212个8GB批次
val batchCount = math.ceil(1793.9 / 8).toInt
val lookupBatches = reverseLookupTableDf
  .withColumn("batch_id", hash(col("hash")) % batchCount)
  .groupBy("batch_id")
  .mapGroups((_, iter) => iter.toList) // 按批次拆分数据集
  .collect()

步骤2:迭代广播并join

遍历每个lookup批次,广播后与主数据集join,最后合并所有结果:

var finalResult = filteredEvents

for (batch <- lookupBatches) {
  val lookupBatchDf = spark.createDataFrame(batch, reverseLookupTableDf.schema)
  // 广播当前批次
  val broadcastedLookup = broadcast(lookupBatchDf)
  
  // 依次完成4个字段的join,注意别名避免冲突
  val joined = finalResult
    .join(broadcastedLookup.as("adId"), col("adId") === col("adId.hashed"), "leftouter")
    .join(broadcastedLookup.as("ip"), col("ae.ip") === col("ip.hashed"), "leftouter")
    .join(broadcastedLookup.as("ua"), col("ae.ua") === col("ua.hashed"), "leftouter")
    .join(broadcastedLookup.as("uid"), col("ae.uuid") === col("uid.hashed"), "leftouter")
  
  // 去重合并,避免同一主数据行被多个批次重复匹配
  finalResult = joined.dropDuplicates(filteredEvents.columns ++ Seq("adId.Key", "ip.Key", "ua.Key", "uid.Key"))
}

注意事项

  • 拆分lookup表时,尽量保证每个批次大小均匀,避免某批次过大;
  • 每次join后必须去重,防止同一主数据行被多个批次匹配多次;
  • 可根据集群executor内存调整batchCount,避免出现OOM。

3. 其他优化建议

  • 优化lookup表:先对lookup表的hash字段去重,减少数据量;同时对lookup表按hash字段分区并Z-Order排序,提升join时的读取性能;
  • 启用自适应执行(AQE):在Databricks中开启AQE(spark.sql.adaptive.enabled=true),它会自动处理数据倾斜,动态调整分区数;
  • 分阶段join:不要一次性做4个join,可分阶段执行,每个阶段join后做一次repartition,避免数据倾斜累积;
  • 处理高频倾斜值:检查是否存在空值或高频哈希值导致的数据倾斜,可单独过滤这些值做join后再与主结果合并。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 01:35:57