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个分桶数。 - 另外还有一种小概率情况:如果你的原始DataFrame
df0本身就有320个分区,且某些操作导致bucketBy没有触发shuffle(这种情况很少见,因为分桶逻辑本身需要shuffle来重新分配数据),也会出现多文件的问题,但核心还是shuffle分区的配置问题。
解决办法
要让分桶表生成正好50个文件,你可以试试这两种方案:
- 直接设置shuffle分区数等于分桶数:在写入分桶表前,添加一行配置代码:
这样shuffle操作就会生成50个分区,每个分区对应一个分桶,最终输出50个part文件。spark.conf.set("spark.sql.shuffle.partitions", 50) - 先按分桶键重分区:在写入前对DataFrame做一次重分区操作,指定分区数为50且按
userid分区:
之后再执行val df0 = df.repartition(50, $"userid")bucketBy写入,就能保证输出文件数和分桶数一致。
补充说明:Spark在后续的2.x版本(比如2.3+)中优化了这个逻辑,
bucketBy会自动将shuffle分区数设置为分桶数,但2.0.1版本还没有这个优化,所以需要手动调整。
内容的提问来源于stack exchange,提问作者astro_asz
相关产品推荐
相关产品推荐

