Spark DataFrame distinct去重写入S3 输出体积暴涨近10倍问题咨询
问题背景
- 业务场景下通过Spark SQL处理数据后使用DataFrame写入S3,两张输入表现状:
input_table_1:数据大小13GB,包含account、membership字段,其中membership为存储member_id列表的数组类型字段input_table_2:数据大小1MB,包含role、id、membership_id、quantity、start_date字段
- 执行查询的SQL如下,对小表添加了广播Join提示:
SELECT /*+ BROADCASTJOIN(input_table_2) */ account, role, id, quantity, cast(start_date AS string) AS start_date FROM input_table_1 INNER JOIN input_table_2 ON array_contains(input_table_1.membership, input_table_2.membership_id)
- 直接写入查询结果时,S3输出数据约1.1TiB,总记录数约7000亿条,排查发现结果存在重复数据。
- 调用
dataframe.distinct.write.parquet("s3path")执行去重写入后,记录数降至原规模的1/3(约2000亿条),但S3输出体积反而增长至17.2TiB。 - 当前Spark配置为
spark.sql.shuffle.partitions=20000,曾尝试调用coalesce后写入S3但问题未解决,需要确认该现象是否符合预期,以及可行的解决方案。
原因说明
该现象完全符合Spark写入Parquet的底层逻辑,不属于异常,核心诱因有三点:
- 小文件过多导致元数据占比暴涨
distinct是宽依赖操作,会触发全量Shuffle,你配置的spark.sql.shuffle.partitions=20000会让Shuffle后默认生成20000个输出分区,每个分区写入时会生成独立的Parquet文件。2000亿条记录分摊到2万个分区,单文件平均仅100万条记录,Parquet文件自带的列统计信息、字典表、文件Footer等元数据的占比会大幅提升,远高于大文件场景下的元数据占比。而你之前直接写入查询结果时,广播Join不会触发全量Shuffle,输出分区数和input_table_1的读入分区数一致(13GB的表读入分区通常仅数百个),单文件大小合理,元数据占比极低。 - 小文件场景下Parquet压缩效率骤降
Parquet的字典编码、RLE编码以及底层通用压缩算法的压缩率,高度依赖单文件内同列数据的重复度。小文件场景下样本量不足,字典覆盖度低,编码和压缩都无法达到最优效率,实际存储膨胀比远高于大文件场景。你之前用coalesce没生效,是因为coalesce是窄依赖操作,只能合并同一节点上的分区,无法跨节点拉取数据合并,最终还是会存在大量远小于最优大小的文件,甚至可能因为强制合并导致单节点内存溢出。 - 去重后数据分布变化进一步拉低压缩率
原始带重复数据的结果集中,相同记录连续存储的概率更高,Parquet的页级编码可以高效压缩连续重复值;去重后相同记录被打散到不同分区、不同文件,单文件内重复值占比下降,也会导致整体压缩率降低。
解决方案
按落地优先级排序:
- 替换coalesce为repartition,控制单文件大小在128MB~256MB最优区间
不要直接使用Shuffle默认的20000分区写入,在distinct之后调用repartition触发全量Shuffle,将数据均匀打散到合理数量的分区。按去重后理想压缩后大小300GiB500GiB估算,单文件按128MB计算,分区数设置为20003000即可,示例代码:dataframe.distinct.repartition(2500).write.parquet("s3path") - 分区内排序提升压缩率
对去重后的数据按低基数字段(如account、role、start_date)做分区内排序后再写入,相同取值的记录会集中在同一个文件、同一个数据页内,字典编码和压缩算法的效率会大幅提升,通常能将存储体积降到比原始未去重结果更小的水平,示例:dataframe.distinct .repartition(2500, col("account")) .sortWithinPartitions("account", "role", "start_date") .write.parquet("s3path") - 开启AQE自动优化分区与压缩配置
写入时开启Spark自适应查询执行,自动合并Shuffle后过小的分区,同时更换压缩比更高的ZSTD压缩算法,不需要手动计算分区数:spark.conf.set("spark.sql.parquet.compression.codec", "zstd") spark.conf.set("spark.sql.adaptive.enabled", "true") spark.conf.set("spark.sql.adaptive.shuffle.targetPostShuffleInputSize", "134217728") // 单分区目标大小128MB - 从源头消除重复,减少全量去重的开销
重复数据本质是Join逻辑产生的:如果input_table_2本身存在重复的membership_id记录,或者input_table_1中同一条account的membership数组包含多个命中input_table_2的membership_id,就会生成重复记录。可以在Join前先对input_table_2按membership_id去重,同时提前炸开input_table_1的membership数组做去重后再Join,能大幅减少Shuffle阶段处理的数据量,避免后续全量Distinct的性能开销。
内容的提问来源于stack exchange,提问作者PSD
相关产品推荐
相关产品推荐

