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

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的底层逻辑,不属于异常,核心诱因有三点:

  1. 小文件过多导致元数据占比暴涨
    distinct是宽依赖操作,会触发全量Shuffle,你配置的spark.sql.shuffle.partitions=20000会让Shuffle后默认生成20000个输出分区,每个分区写入时会生成独立的Parquet文件。2000亿条记录分摊到2万个分区,单文件平均仅100万条记录,Parquet文件自带的列统计信息、字典表、文件Footer等元数据的占比会大幅提升,远高于大文件场景下的元数据占比。而你之前直接写入查询结果时,广播Join不会触发全量Shuffle,输出分区数和input_table_1的读入分区数一致(13GB的表读入分区通常仅数百个),单文件大小合理,元数据占比极低。
  2. 小文件场景下Parquet压缩效率骤降
    Parquet的字典编码、RLE编码以及底层通用压缩算法的压缩率,高度依赖单文件内同列数据的重复度。小文件场景下样本量不足,字典覆盖度低,编码和压缩都无法达到最优效率,实际存储膨胀比远高于大文件场景。你之前用coalesce没生效,是因为coalesce是窄依赖操作,只能合并同一节点上的分区,无法跨节点拉取数据合并,最终还是会存在大量远小于最优大小的文件,甚至可能因为强制合并导致单节点内存溢出。
  3. 去重后数据分布变化进一步拉低压缩率
    原始带重复数据的结果集中,相同记录连续存储的概率更高,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 18:09:16