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
相关产品推荐
相关产品推荐

