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

