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

Spark中saveAsTable搭配bucketBy设50桶却生成大量文件的原因?

问题:Spark分桶表生成文件数远超预期的原因分析

我在Spark 2.0.1的cluster模式下提交任务,代码如下:
首先初始化SparkSession:

val spark = SparkSession.builder
  .appName("myApp")
  .config("hive.metastore.uris", "thrift://XXX.XXX.net:9083")
  .config("spark.sql.sources.bucketing.enabled", true)
  .enableHiveSupport()
  .getOrCreate()

然后从HDFS读取Parquet文件:

val df = spark.read
  .format("parquet")
  .load("hdfs://XXX.XX.X.XX/myParquetFile")

最后将DataFrame按userid分50桶保存为Hive表:

df0.write
  .bucketBy(50, "userid")
  .saveAsTable("myHiveTable")

但查看HDFS上的Hive仓库目录/user/hive/warehouse/myHiveTable时,发现有320个part-*.parquet文件,远超过预期的50个,这是为什么?


原因分析

我来给你拆解下这个问题,核心原因是Spark 2.0.1版本中bucketBy的shuffle行为和你预期的不一致,具体细节如下:

  • shuffle分区数不由分桶数决定:在Spark 2.0.1这个较早的版本里,当你使用bucketBy时,Spark不会自动将shuffle分区数设置为分桶数(也就是50),而是沿用spark.sql.shuffle.partitions这个配置的数值。如果你的集群中这个配置被设为320(或者你没有修改,集群默认是这个值),那么shuffle操作会生成320个数据分区。
  • 每个shuffle分区对应一个输出文件:每个shuffle分区的任务会根据userid的哈希值,将数据写入对应的分桶文件中,但每个shuffle分区只会生成一个part文件。最终的文件总数就等于shuffle分区数——也就是320个,而非你设置的50个分桶数。
  • 另外还有一种小概率情况:如果你的原始DataFramedf0本身就有320个分区,且某些操作导致bucketBy没有触发shuffle(这种情况很少见,因为分桶逻辑本身需要shuffle来重新分配数据),也会出现多文件的问题,但核心还是shuffle分区的配置问题。

解决办法

要让分桶表生成正好50个文件,你可以试试这两种方案:

  • 直接设置shuffle分区数等于分桶数:在写入分桶表前,添加一行配置代码:
    spark.conf.set("spark.sql.shuffle.partitions", 50)
    
    这样shuffle操作就会生成50个分区,每个分区对应一个分桶,最终输出50个part文件。
  • 先按分桶键重分区:在写入前对DataFrame做一次重分区操作,指定分区数为50且按userid分区:
    val df0 = df.repartition(50, $"userid")
    
    之后再执行bucketBy写入,就能保证输出文件数和分桶数一致。

补充说明:Spark在后续的2.x版本(比如2.3+)中优化了这个逻辑,bucketBy会自动将shuffle分区数设置为分桶数,但2.0.1版本还没有这个优化,所以需要手动调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 03:35:50