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

内存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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 12:13:12