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

Spark 3.3如何利用AQE动态合并减少输出文件数量?

问题描述

以下Spark函数用于对比两个版本的DataFrame,仅保留修改或删除的行:

def calculate_delta(old_df: DataFrame,
                    new_df: DataFrame,
                    id_keys: List[str],
                    deletion_field_name: str,
                    excluded_fields: List[str] = None) -> DataFrame:
    # make sure both the new_df and old_df have the same columns order
    if excluded_fields:
        old_df = old_df.drop(*excluded_fields)
    shared_columns = set(old_df.columns).intersection(new_df.columns)
    new_df = new_df.select(*id_keys, f.struct(*shared_columns).alias("new_struct"))
    old_df = old_df.select(*id_keys, f.struct(*shared_columns).alias("old_struct"))

    return new_df.join(old_df, on=id_keys, how="outer") \
        .filter("old_struct IS NULL OR new_struct IS NULL OR new_struct != old_struct") \
        .withColumn(deletion_field_name, f.expr("new_struct is null"))\
        .withColumn("output_struct", f.when(f.col("new_struct").isNull(), f.col("old_struct")).otherwise(f.col("new_struct")))\
        .selectExpr(f"(new_struct IS NULL) as {deletion_field_name}",
                    "output_struct.*")

当前场景:关联后数据量约150M行,过滤后的行数动态波动(从50行到50M行不等)。已设置spark.sql.shuffle.partitions=720,写入数据时会生成大量文件(最多720个),希望利用Spark 3.3的AQE动态合并特性减少分区及输出文件数,且避免冗余shuffle操作。

可行优化建议
  • 开启AQE核心合并参数
    Spark 3.3的AQE动态合并小分区功能需确保以下参数配置生效(部分参数默认开启,建议显式声明):

    • spark.sql.adaptive.enabled=true:启用自适应执行框架
    • spark.sql.adaptive.coalescePartitions.enabled=true:开启小分区自动合并
    • spark.sql.adaptive.coalescePartitions.minPartitionNum=1:允许合并至最少1个分区
    • spark.sql.adaptive.advisoryPartitionSizeInBytes=128m:设置目标分区大小(可根据集群存储调整为256m等)
      这些参数会让AQE在shuffle阶段(join触发)完成后,根据实际数据量自动合并小分区,无需额外shuffle操作。
  • 保留初始shuffle分区数,交给AQE动态调整
    无需修改spark.sql.shuffle.partitions=720的配置,这个值是shuffle阶段的初始分区数,AQE会根据过滤后的数据量自动合并冗余分区。不要手动调用coalesce或repartition,避免触发额外shuffle。

  • 保持计算链连续性,让AQE感知数据变化
    确保过滤后的后续操作(如selectExpr)不会触发新的shuffle,保持数据处理链的连续性,让AQE能完整感知过滤后的数据量变化,从而精准执行分区合并。

  • 配合写入参数控制文件大小
    增加以下写入相关配置,进一步优化输出文件数量:

    • spark.sql.files.maxRecordsPerFile=1000000:限制单个文件的最大记录数(根据单条记录大小调整)
    • spark.sql.files.openCostInBytes=134217728:与目标分区大小一致,帮助Spark判断是否需要合并文件
  • 验证AQE生效状态
    提交任务后,可在Spark UI的SQL页面查看Adaptive Execution模块,确认是否触发了小分区合并,以及最终生成的分区数是否符合预期。

内容的提问来源于stack exchange,提问作者Infinity

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 03:53:11