Spark如何实现按列分区时每个分区组生成固定数目的子分区
问题原因
你当前使用的repartition(4, col("country_code"), col("record_date"))逻辑不符合预期:Spark哈希分区规则是对指定列的取值计算哈希值后,对第一个参数的总分区数取模,相同列值组合的哈希结果固定,所以同一个(country_code, record_date)的所有数据只会被分到同一个Spark RDD分区,最终写盘时每个Hive分区自然只能生成1个文件。
解决方法
要实现每个Hive分区对应4个文件,需要给同个Hive分区的数据补充随机分片字段,让同个Hive分区的数据可以分散到4个不同的Spark分区中,示例代码如下:
// 导入依赖函数 import org.apache.spark.sql.functions.{rand, col} val processedDF = df // 给每条数据生成0-3的随机分片标记,同Hive分区内会均匀分成4组 .withColumn("random_shard", (rand() * 4).cast("integer")) // 按Hive分区列+随机分片列做重分区 .repartition(col("country_code"), col("record_date"), col("random_shard")) // 写完丢弃辅助用的分片列 .drop("random_shard") // 后续按正常流程写入Hive分区表即可 processedDF.write .partitionBy("country_code", "record_date") .mode("append") .saveAsTable("目标表名")
如果不想新增辅助列,也可以通过写入参数控制单个文件的最大记录数间接实现目标,适合你能提前估算出128MB对应单文件记录数的场景:
df.write .partitionBy("country_code", "record_date") .option("maxRecordsPerFile", 你估算的单文件记录数) .mode("append") .saveAsTable("目标表名")
注意事项
- 如果你的Hive分区总数较多,可以把
repartition的第一个参数设置为「Hive分区总数 * 4」,避免多个不同Hive分区的数据落到同一个Spark分区,导致最终文件大小不均匀。 - 不需要额外设置其他文件合并参数,避免覆盖你控制的分片逻辑。
内容的提问来源于stack exchange,提问作者Udaya
相关产品推荐
相关产品推荐

