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

Glue上PySpark写入Apache Iceberg生成过多小文件的解决方法

问题描述

在AWS Glue上运行PySpark任务,处理数据后保存为Apache Iceberg表时,每个分区内生成多个小文件,期望每个分区仅保留一个合并后的文件。当前代码片段如下:

import pyspark.sql.functions as f
from pyspark.conf import SparkConf
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from pyspark.sql import DataFrame

conf = (
    SparkConf()
    .setAppName(APP_NAME)
    .set(
        "spark.sql.extensions",
        "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
    )
    .set("spark.sql.catalog.glue_catalog", "org.apache.iceberg.spark.SparkCatalog")
    .set("spark.sql.catalog.glue_catalog.warehouse", BRONZE_PATH)
    .set(
        "spark.sql.catalog.glue_catalog.catalog-impl",
        "org.apache.iceberg.aws.glue.GlueCatalog",
    )
    .set("spark.sql.catalog.glue_catalog.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
    .set("spark.sql.shuffle.partitions", "100")
)
sc = SparkContext(conf=conf)
glueContext = GlueContext(sc)
glue_db = glueContext.create_dynamic_frame.from_catalog(database=DATABASE_NAME, table_name=LANDING_TABLE_NAME)
df = glue_db.toDF()
df.createOrReplaceTempView(APP_NAME)

# do all processing here...

df = df.sortWithinPartitions("issueid")
(
    df.writeTo(f"glue_catalog.bronze_{ENVIRONMENT}.{BRONZE_TABLE_NAME}")
    .using("iceberg")
    .tableProperty("format-version", "2")
    .tableProperty("location", BRONZE_PATH + BRONZE_TABLE_OUTPUT)
    .tableProperty("write.distribution.mode", "hash")
    .tableProperty("write.target-file-size-bytes", "536870912")
    .partitionedBy("issueid")
    .createOrReplace()
)
解决方案

要实现每个Iceberg分区仅一个文件,需从Spark分区控制和Iceberg写入配置两方面调整,具体步骤如下:

1. 按分区键重分区,对齐Spark与Iceberg分区

当前spark.sql.shuffle.partitions=100会将数据分散到100个Spark分区,即使设置Iceberg目标文件大小,也会因Spark分区过多生成小文件。处理完成后,直接按Iceberg的分区键issueid重分区,确保每个issueid对应一个Spark分区:

# 处理完成后添加重分区逻辑
df = df.repartition("issueid")
df = df.sortWithinPartitions("issueid")

2. 优化Iceberg写入配置

  • 移除write.distribution.mode=hash:该配置会额外打乱已按分区键整理的数据,导致多文件生成,直接删除或设为none。
  • 启用write.merge.enabled=true:让Iceberg自动合并写入过程中产生的小文件,确保最终每个分区仅保留符合目标大小的文件。

3. 完整调整后的代码

import pyspark.sql.functions as f
from pyspark.conf import SparkConf
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from pyspark.sql import DataFrame

conf = (
    SparkConf()
    .setAppName(APP_NAME)
    .set(
        "spark.sql.extensions",
        "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
    )
    .set("spark.sql.catalog.glue_catalog", "org.apache.iceberg.spark.SparkCatalog")
    .set("spark.sql.catalog.glue_catalog.warehouse", BRONZE_PATH)
    .set(
        "spark.sql.catalog.glue_catalog.catalog-impl",
        "org.apache.iceberg.aws.glue.GlueCatalog",
    )
    .set("spark.sql.catalog.glue_catalog.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
    .set("spark.sql.shuffle.partitions", "100")
)
sc = SparkContext(conf=conf)
glueContext = GlueContext(sc)
glue_db = glueContext.create_dynamic_frame.from_catalog(database=DATABASE_NAME, table_name=LANDING_TABLE_NAME)
df = glue_db.toDF()
df.createOrReplaceTempView(APP_NAME)

# do all processing here...

# 关键调整:按分区键重分区,确保每个issueid对应一个Spark分区
df = df.repartition("issueid")
df = df.sortWithinPartitions("issueid")

(
    df.writeTo(f"glue_catalog.bronze_{ENVIRONMENT}.{BRONZE_TABLE_NAME}")
    .using("iceberg")
    .tableProperty("format-version", "2")
    .tableProperty("location", BRONZE_PATH + BRONZE_TABLE_OUTPUT)
    .tableProperty("write.target-file-size-bytes", "536870912")
    .tableProperty("write.merge.enabled", "true")
    .partitionedBy("issueid")
    .createOrReplace()
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 00:50:08