Spark作业卡在最后2个任务(共100个),求排查调优方案
嘿,作为Spark新手碰到这种卡任务的情况太正常了!结合你的描述,我大概率能帮你定位问题——那个groupBy片段就是数据倾斜的罪魁祸首,咱们一步步来解决:
一、先揪出核心问题:数据倾斜导致任务卡住
你提到的这段代码:
Dataset<Row> duplicatePrefetchPrerenderHashDS = hashedPageViewDS .select(hashedPageViewDS.col(PREFETCH_PRERENDER_HASH)) .groupBy(hashedPageViewDS.col(PREFETCH_PRERENDER_HASH)) .count() .withColumnRenamed("count", "cnt") .where("cnt>1");
当你按PREFETCH_PRERENDER_HASH分组统计时,如果某些哈希值的出现次数特别多(比如几十万甚至上百万次),就会导致少数几个任务要处理远超其他任务的数据量。最后你看到的98/100卡住,就是剩下的2个任务在啃这些超大分区,完全跑不动。
二、用Broadcast优化:小表广播替代Shuffle-heavy的GroupBy
既然你关注到了负载不均的问题,那咱们就把broadcast用起来!核心思路是:先提取出重复的哈希值小集合,广播到所有Executor节点,再通过关联操作替代groupBy过滤逻辑,彻底避免大规模shuffle:
优化后的代码示例
// 1. 先提取重复的哈希值,只保留需要的列,得到一个轻量的小Dataset Dataset<Row> duplicateHashesDS = hashedPageViewDS .select(PREFETCH_PRERENDER_HASH) .groupBy(PREFETCH_PRERENDER_HASH) .count() .where("cnt > 1") .select(PREFETCH_PRERENDER_HASH); // 只保留哈希列,进一步缩小数据集 // 2. 手动广播这个小Dataset(Spark会自动判断是否适合广播,手动声明更保险) Broadcast<Dataset<Row>> broadcastDuplicateHashes = spark.sparkContext().broadcast(duplicateHashesDS); // 3. 用广播关联替代原逻辑(这里假设你要保留重复的记录,可根据实际业务调整关联类型) Dataset<Row> filteredDS = hashedPageViewDS.join( broadcast(broadcastDuplicateHashes.value()), hashedPageViewDS.col(PREFETCH_PRERENDER_HASH).equalTo(broadcastDuplicateHashes.value().col(PREFETCH_PRERENDER_HASH)), "inner" );
这样一来,原本需要shuffle整个4亿条数据的groupBy操作,变成了只广播几千/几万条的小数据集,任务负载会均匀很多!
三、其他辅助优化点
1. 调整分区缩减的方式
你用了coalesce(100),它默认不触发shuffle,会直接合并现有分区。如果原来的1000分区本身数据分布不均(比如前面的倾斜导致),合并后的100分区也会继承这种不均,进而拖慢写入阶段。
如果集群资源充足,可以改成repartition(100)(会触发shuffle,重新均匀分配数据);如果担心shuffle开销,那就先解决前面的数据倾斜问题,再用coalesce。
2. S3写入调优
卡在写入阶段也可能和S3的特性有关,试试这些参数优化:
- 开启高效提交算法:
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=2,大幅减少S3的commit时间 - 若使用AWS环境,启用S3Guard:避免S3最终一致性导致的文件读写问题
- 检查
partitionBy("year", "month")的基数:如果某个年月的记录量特别大,会生成超大文件拖慢写入,可考虑新增day分区分散压力
总结
优先解决数据倾斜(用broadcast优化groupBy逻辑),再调整分区策略,最后优化S3写入参数,应该就能解决卡在98/100的问题啦!
内容的提问来源于stack exchange,提问作者Geoff L.

