内存repartition与磁盘partitionBy对比及相关技术疑问
Spark中repartition()与partitionBy()的常见疑问解答
问题1:如何确认使用repartition()时内存中存在多个分区?
有几种直接的验证方式:
- 调用API方法:直接用DataFrame/RDD的
getNumPartitions()方法获取当前内存分区数。比如Python中执行articles.repartition(5).rdd.getNumPartitions(),返回结果为5就说明内存中存在5个分区。 - 查看Spark UI:提交作业后进入Stages页面,对应Stage的
Input Size / Records区域会列出每个内存分区的信息,统计条目数量即可得到分区数。 - 用glom()直观验证:将DataFrame转成RDD后调用
glom(),该方法会把每个分区的所有元素打包成一个列表,执行collect()后统计返回的列表数量,就是当前的内存分区数。示例代码:articles.repartition(3).rdd.glom().collect(),若返回3个列表则对应3个内存分区。
问题2:repartition仅在内存中创建分区,为何执行articles.repartition(1).write.saveAsTable("articles_table", format = 'orc', mode = 'overwrite')仅生成一个文件?该操作与partitionBy()有何区别?
首先纠正一个误解:repartition()并非只用于读取操作,它的核心是调整内存中DataFrame/RDD的分区数量,而Spark写入数据的默认逻辑是每个内存分区对应生成一个输出文件。所以当你用repartition(1)将内存分区数调整为1时,写入自然只会生成一个文件——因为所有数据都在同一个内存分区里,最终由这个分区的数据写出单个文件。
它和partitionBy()的核心区别如下:
- 作用维度不同:
- repartition()是对内存中的数据做数量上的分区拆分/合并,完全基于并行度调整,和数据内容无关,本质是通过shuffle重新分配数据到指定数量的内存分区中。
- partitionBy()是在磁盘上做业务字段维度的目录分区,比如
partitionBy("category")会在输出目录下生成category=科技/、category=娱乐/这类子目录,数据按字段值进入对应目录。
- 输出结构不同:
- repartition(N)写入后,所有文件平级放在根输出目录下,文件数量等于N(无额外小文件合并逻辑时)。
- partitionBy(col)写入后,文件会分散在各个字段值对应的子目录中,每个子目录下的文件数量等于当前的内存分区数(未额外调整repartition的情况下)。
- 使用场景不同:
- repartition()适合调整作业并行度、减少小文件数量,或在shuffle后平衡分区数据量。
- partitionBy()适合按业务维度做分区存储,方便后续读取时只加载指定分区的数据(比如仅查询某类别的数据,无需扫描全量文件)。
内容的提问来源于stack exchange,提问作者Blue Clouds
相关产品推荐
相关产品推荐

