尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

spark中涉及到哪些shuffle

spark中涉及到哪些shuffle Spark 中“Shuffle 有几种”要分两个角度回答按 RDD 依赖分类宽依赖产生 Shuffle按 Shuffle 的实现方式分类Hash Shuffle、Sort Shuffle 等。实际开发中通常重点理解第二种。一、按计算依赖分类1. Narrow Dependency窄依赖父 RDD 的一个分区只被一个子 RDD 分区使用不需要跨节点重新分发数据。常见算子map filter flatMap mapPartitions例如rdd.map(xx*2)数据可以在当前分区内直接处理不发生 Shuffle。2. Shuffle Dependency宽依赖父 RDD 的一个分区可能被多个子 RDD 分区使用需要根据 Key 重新分区和网络传输。常见算子reduceByKey groupByKey join distinct sortByKey repartition例如rdd.reduceByKey(__)相同 Key 的数据可能来自不同节点需要把它们发送到同一个下游分区这就是 Shuffle。二、按实现方式分类1. Hash Shuffle早期 Spark 主要使用 Hash Shuffle。Map Task 会根据目标分区和 Key 进行 Hash然后把数据写入不同的临时文件。例如有M 个 Map Task R 个 Reduce 分区可能产生接近M × R个文件。特点实现简单根据 Hash 分区文件数量可能非常多容易造成小文件和文件句柄压力现代 Spark 中已经不是主流实现。2. Sort Shuffle现代 Spark 默认主要使用 Sort Shuffle。Map Task 先把数据写入内存缓冲区内存不足时溢写到磁盘最后将多个溢写文件归并并按照分区信息组织输出。Reduce Task 再读取对应分区的数据。特点文件数量相对可控支持内存溢写通常比旧 Hash Shuffle 更适合大规模数据是现代 Spark 的主流 Shuffle 实现。执行过程大致是Map Task - 内存缓冲 - Spill 到磁盘 - 多个 Spill 文件归并 - 生成 Shuffle 输出 - Reduce Task 拉取三、Sort Shuffle 内部的几种写法严格来说Sort Shuffle 还可以根据数据规模和配置走不同路径。1. Unsafe Shuffle Writer针对 UnsafeRow 等二进制格式优化利用 Tungsten 内存管理和高效排序。通常适用于 Spark SQL、DataFrame、Dataset 场景。特点二进制处理减少对象创建内存利用率较高性能通常较好。2. Serialized Shuffle Writer数据序列化后在内存中排序减少 Java 对象开销。3. BypassMergeSortShuffleWriter当满足特定条件时使用例如没有 map-side combine下游分区数不超过spark.shuffle.sort.bypassMergeThreshold默认通常是 200。它会分别写入各个分区文件最后再合并这些文件。优点不需要对记录进行排序某些场景下更快。缺点仍然可能产生较多临时文件不支持 map-side combine。4. SerializedShuffleWriter用于序列化数据的 Shuffle 写入路径通常在不需要特殊排序逻辑时使用。实际使用哪一种由 Spark 的执行计划、数据格式、是否需要聚合以及配置共同决定。四、Shuffle 的两个重要阶段一次 Shuffle 通常包含两个阶段。Map 阶段Map Task读取上游数据根据分区器计算目标分区可能执行 map-side combine写入内存或磁盘生成 Shuffle 文件向 Driver 汇报输出位置。Reduce 阶段Reduce Task从各个 Executor 拉取自己负责的分区数据进行归并、聚合或排序输出最终结果。可以理解为Map 端写 Reduce 端拉五、哪些操作会触发 Shuffle通常会触发rdd.groupByKey()rdd.reduceByKey(__)rdd.aggregateByKey(...)rdd.join(other)rdd.distinct()rdd.sortByKey()rdd.repartition(n)DataFrame / SQL 中常见的 Shuffle 来源GROUPBYJOINORDERBYDISTINCTDISTRIBUTEBYREPARTITION例如df.groupBy(user_id).count()需要把相同user_id的数据发送到同一个分区通常会产生 Shuffle。通常不会触发map filter flatMap mapValues filterByRangeDataFrame 中常见的窄依赖操作df.select(...)df.filter(...)df.withColumn(...)但要注意具体是否产生 Shuffle最终应以物理执行计划为准。六、reduceByKey和groupByKey的区别两者都会 Shuffle但效率通常不同。groupByKeyrdd.groupByKey()先把同一个 Key 的所有 Value 发送到一起再聚合。reduceByKeyrdd.reduceByKey(__)可以在 Map 端先进行局部聚合也就是 map-side combine减少网络传输量。例如原始数据(a, 1) (a, 1) (a, 1)Map 端可以先变成(a, 3)再发送到 Reduce 端。所以需要聚合时通常优先使用reduceByKey、aggregateByKey或combineByKey而不是先groupByKey。七、DataFrame 中常见的 Shuffle 类型Spark SQL 中经常看到以下分区方式1. HashPartitioning按照 Key 的 Hash 值分区hash(key) % numPartitions常用于GROUP BY等值 JoindropDuplicatesrepartition(key)。2. RangePartitioning按照范围分区常用于排序相关操作df.repartitionByRange(order_id)或执行全局排序时使用。3. RoundRobinPartitioning轮询分发数据常见于df.repartition(10)不指定 Key 时通常按照轮询方式重新分区。八、Broadcast Join 不一定需要 Shuffle如果一张表很小可以广播到各个 Executorfrompyspark.sql.functionsimportbroadcast resultlarge_df.join(broadcast(small_df),user_id)这种情况下小表被广播大表通常不需要为了 Join Key 做全量 Shuffle可以避免大规模网络重分区。但广播表需要能放进 Executor 内存不能盲目使用。九、如何查看是否发生 ShuffleDataFrame / SQLdf.explain(formatted)重点观察是否出现Exchange例如Exchange hashpartitioning(user_id, 200)通常表示发生了 Shuffle。RDD 可以查看依赖rdd.toDebugString如果出现ShuffledRDD说明存在 Shuffle 依赖。Spark UI 中也可以查看Shuffle ReadShuffle WriteFetch Wait TimeRecords Read / WriteSpill MemorySpill Disk。总结如果按最常用的实现方式回答Spark Shuffle 主要可以说1. Hash Shuffle历史实现 2. Sort Shuffle现代主流实现如果进一步细分 Sort Shuffle 的写入器还包括- Unsafe Shuffle Writer - Serialized Shuffle Writer - BypassMergeSortShuffleWriter而从 Spark 计算模型看真正决定是否发生 Shuffle 的核心是窄依赖不需要跨分区重组 宽依赖需要跨分区重组会产生 Shuffle在实际排查性能时最值得关注的是是否发生了不必要的ExchangeShuffle 分区数是否合理是否存在数据倾斜是否产生了磁盘 Spill是否可以使用 map-side combine是否可以使用 Broadcast Join。
返回列表