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

强制Spark使用BaseShuffleHandle后查询结果异常求助

强制Spark使用BaseShuffleHandle后查询结果异常的问题

Spark包含三种ShuffleHandle实现类:SerializedShuffleHandle、BypassMergeSortShuffleHandle、BaseShuffleHandle。我希望强制Spark仅使用BaseShuffleHandle:通过调整bypassMergeThreshold参数可以轻松阻止BypassMergeSortShuffleHandle的使用,但要禁用SerializedShuffleHandle则比较棘手。

为此我修改了SortShuffleManager的registerShuffle方法,移除原有分支逻辑,直接返回BaseShuffleHandle,修改前后的代码如下:

原方法代码

override def registerShuffle[K, V, C](
    shuffleId: Int,
    dependency: ShuffleDependency[K, V, C]): ShuffleHandle = {
  if (SortShuffleWriter.shouldBypassMergeSort(conf, dependency)) {
    // If there are fewer than spark.shuffle.sort.bypassMergeThreshold partitions and we don't
    // need map-side aggregation, then write numPartitions files directly and just concatenate
    // them at the end. This avoids doing serialization and deserialization twice to merge
    // together the spilled files, which would happen with the normal code path. The downside is
    // having multiple files open at a time and thus more memory allocated to buffers.
    new BypassMergeSortShuffleHandle[K, V](
      shuffleId, dependency.asInstanceOf[ShuffleDependency[K, V, V]])
  } else if (SortShuffleManager.canUseSerializedShuffle(dependency)) {
    // Otherwise, try to buffer map outputs in a serialized form, since this is more efficient:
    new SerializedShuffleHandle[K, V](
      shuffleId, dependency.asInstanceOf[ShuffleDependency[K, V, V]])
  } else {
    // Otherwise, buffer map outputs in a deserialized form:
    new BaseShuffleHandle(shuffleId, dependency) 
  }
}

修改后的方法代码

override def registerShuffle[K, V, C](
    shuffleId: Int,
    dependency: ShuffleDependency[K, V, C]): ShuffleHandle = {
    
  new BaseShuffleHandle(shuffleId, dependency)
}

我在1GB规模的TPC-DS数据集上运行了如下SQL查询(query03),Spark历史服务器显示所有阶段执行成功,但最终查询结果错误:

select  dt.d_year 
     ,item.i_brand_id brand_id 
     ,item.i_brand brand
     ,sum(ss_ext_sales_price) sum_agg
from  date_dim dt 
      ,store_sales
      ,item
where dt.d_date_sk = store_sales.ss_sold_date_sk
  and store_sales.ss_item_sk = item.i_item_sk
  and item.i_manufact_id = 128
  and dt.d_moy=11
group by dt.d_year
      ,item.i_brand
      ,item.i_brand_id
order by dt.d_year
         ,sum_agg desc
         ,brand_id
limit 100;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 23:57:32