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

Spark中.collect()触发Shuffle但.show()不触发?附行计数疑问

Spark排序操作的两个疑问解答

背景

我正在通过《Spark: The Definitive Guide》学习Spark,遇到了一些难以理解的行为。读取了一个包含DEST_COUNTRY_NAME、ORIGIN_COUNTRY_NAME、count三列的CSV文件,数据集共256行。

执行代码

flightData2015 = (spark
                  .read
                  .format('csv')
                  .option('inferSchema', 'true')
                  .option('header', 'true')
                  .option('path', 'dbfs:/FileStore/tables/spark_definitive_guide/data/flight_data/csv/2015_summary.csv')
                  .load())

# 设置shuffle分区数为5
spark.conf.set("spark.sql.shuffle.partitions", "5")

排序后的物理执行计划

flightData2015.sort("count").explain()

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Sort [count#337 ASC NULLS FIRST], true, 0
   +- Exchange rangepartitioning(count#337 ASC NULLS FIRST, 5), ENSURE_REQUIREMENTS, [plan_id=255]
      +- FileScan csv [DEST_COUNTRY_NAME#335,ORIGIN_COUNTRY_NAME#336,count#337] Batched: false, DataFilters: [], Format: CSV, Location: InMemoryFileIndex(1 paths)[dbfs:/FileStore/tables/spark_definitive_guide/data/flight_data/csv/201..., PartitionFilters: [], PushedFilters: [], ReadSchema: struct<DEST_COUNTRY_NAME:string,ORIGIN_COUNTRY_NAME:string,count:int>

疑问1:为何.collect()会触发Shuffle,而.show(1000)不会?

  • 执行flightData2015.sort("count").show(1000)时,Spark UI的DAG中无Shuffle,仅进行数据扫描和展示;
  • 执行flightData2015.sort("count").collect()时,DAG和执行计划中却出现Shuffle。
    两者都获取全部256行,为何存在此差异?

解答

这是Spark针对不同Action操作的优化逻辑差异导致的:

  • show(n)的优化逻辑:当指定的n大于等于数据集总行数时,Spark会采用「本地排序+Driver端合并」的策略:每个分区先独立完成排序,再将排序后的分区数据发送到Driver端,由Driver对所有分区的结果做一次全局合并排序。这个过程不需要跨节点的Shuffle(Exchange操作),因为数据不需要重新分区,只是在每个分区内部排序后汇总到Driver做最终整合。
  • collect()的执行逻辑:collect()没有条数限制,Spark会采用通用的全局排序执行计划:按照spark.sql.shuffle.partitions设置的分区数,对数据做范围分区(range partitioning)的Shuffle,将数据按count值分配到5个分区,每个分区内的数据保持有序,之后再将所有分区的数据拉取到Driver端合并成完整的有序结果。这就是为什么collect()会触发Shuffle,而show(1000)不会。

疑问2:.collect()的DAG中,Scan csv显示返回512行,但实际文件仅有256行,这是为什么?

解答

这个现象和Spark的行数估算机制有关:

  • Spark在执行文件扫描前,会根据文件的元数据(比如文件大小、历史统计的行平均长度)估算返回的行数,而非实时精确计数。如果CSV文件的行长度差异较大,或者文件元数据信息不准确,估算值就会和实际行数出现偏差。
  • 另外,若开启了自适应执行(Adaptive Execution),Spark可能会在执行过程中动态调整统计信息,但UI上展示的是初始阶段的估算行数,所以会出现和实际行数不一致的情况。需要注意的是,这个估算值不会影响最终处理结果的正确性,Spark最终会读取到准确的256行数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 13:57:02