Spark DataFrame左反连接后按分区设文件数的优化方案咨询
问题根源:你当前用coalesce(numPartitions)是给整个DF设置全局分区数,导致所有date分区共享这10个文件;而直接repartition("date")会触发全局shuffle,把所有数据重新按date分区,代价极高。
给你几个实用的解决办法,按需选择:
办法1:按日期分区逐个处理(无全局shuffle,精确控文件数)
既然目标表本身就是按date分区的Hive表,直接逐个读取每个date分区的数据处理就行。这种方式完全避免全局shuffle,仅在单个date内部做必要的小量shuffle,还能精确保证每个date生成10个文件。
代码示例:
// 先定义好参数 val numFilesPerDate = 10 val inputTable = "input_table" val referenceTable = "reference_table" val customerIDCol = "account_id" // 替换成你的Customer_ID变量值 val outputPath = "hdfs_path" // 先拿到目标表所有的date分区列表 val datePartitions = spark.sql(s"SHOW PARTITIONS $inputTable") .map(row => row.getString(0).split("=")(1)) // 从"date=20220101"里提取出日期值 .collect() // 缓存reference表,避免每次循环都重新读表 val referenceDF = spark.table(referenceTable).cache() // 循环处理每个日期分区 datePartitions.foreach { date => // 只读取当前date分区的数据 val targetDF = spark.sql(s"SELECT * FROM $inputTable WHERE date = '$date'") // 左反连接,用broadcast广播reference表,避免shuffle目标表的大数据 val purgeDF = targetDF.join(broadcast(referenceDF), Seq(customerIDCol), "left_anti") // 重分区为10个,保证这个date输出10个文件 purgeDF.repartition(numFilesPerDate) .write .mode("append") // 逐个分区追加写入 .partitionBy("date") .parquet(outputPath) } // 用完释放缓存 referenceDF.unpersist()
为啥好用:
- 提前缓存reference表,省了重复读表的开销
- 用broadcast广播小表,左反连接时根本不会shuffle目标表的数据,性能拉满
- 单个日期分区单独处理,资源占用更均衡,不会出现全局shuffle那种资源挤爆的情况
办法2:自定义分区器(批量处理,无跨日期shuffle)
要是不想循环处理每个date,可以用自定义分区器,把同一个date的数据分配到该date专属的10个分区里,不同date的数据绝不混在一起,写入时每个date分区自动生成10个文件。
代码示例:
// 定义参数 val numFilesPerDate = 10 val inputTable = "input_table" val referenceTable = "reference_table" val customerIDCol = "account_id" val outputPath = "hdfs_path" // 获取所有date分区列表 val datePartitions = spark.sql(s"SHOW PARTITIONS $inputTable") .map(row => row.getString(0).split("=")(1)) .collect() val totalPartitions = datePartitions.size * numFilesPerDate // 自定义分区器:每个date对应10个专属分区 class DatePartitioner(dateList: Array[String], filesPerDate: Int) extends org.apache.spark.Partitioner { override def numPartitions: Int = dateList.length * filesPerDate override def getPartition(key: Any): Int = { val (date, accountId) = key.asInstanceOf[(String, String)] val dateIndex = dateList.indexOf(date) // 用account_id的hash值把同一个date的数据分散到10个分区 dateIndex * filesPerDate + Math.abs(accountId.hashCode() % filesPerDate) } } // 执行左反连接,还是用broadcast优化 val purgeDF = spark.sql( s"""SELECT /*+ BROADCASTJOIN(ref) */ target.* |FROM $inputTable target |LEFT ANTI JOIN $referenceTable ref |ON target.$customerIDCol = ref.$customerIDCol""".stripMargin ) // 转成RDD,用(date, account_id)当key,自定义分区器重新分区 val partitionedRDD = purgeDF.rdd .map(row => ((row.getAs[String]("date"), row.getAs[String](customerIDCol)), row)) .partitionBy(new DatePartitioner(datePartitions, numFilesPerDate)) .values // 转回DataFrame写入 spark.createDataFrame(partitionedRDD) .write .partitionBy("date") .mode("overwrite") .parquet(outputPath)
优势:
- 自定义分区器保证同一date的数据只在内部shuffle,不会跨date混排,shuffle代价比全局repartition小太多
- 一次性处理所有date,不用循环提交作业
办法3:用文件记录数参数(近似控制,零shuffle)
要是对文件数的精确性要求不高,直接设置spark.sql.files.maxRecordsPerFile参数,让Spark自动按记录数分割文件,近似达到每个date分区10个文件的效果。这种方式完全没shuffle,效率最高,但文件数可能因为每个date的数据量不同略有偏差。
代码示例:
// 估算每个date分区的记录数,比如每个date有100万条,就设10万,这样每个date大概生成10个文件 spark.conf.set("spark.sql.files.maxRecordsPerFile", "100000") val df_purge = spark.sql(s"SELECT /*+ BROADCASTJOIN(ref) */ target.* FROM input_table target LEFT ANTI JOIN ${reference_table} ref ON target.${Customer_ID} = ref.${Customer_ID}") // 不用手动设分区,Spark自动按记录数切文件 df_purge.write.partitionBy("date").mode("overwrite").parquet("hdfs_path")
总结:优先选办法1,既能精确控文件数,又能避免全局shuffle,性能最优;想要批量处理就选办法2;对精度要求不高的话,办法3最省事。
内容的提问来源于stack exchange,提问作者Arvinth
相关产品推荐
相关产品推荐

