如何触发AWS Glue Iceberg表的Compaction操作
解决AWS Glue Iceberg表Compaction未触发的问题
核心问题分析
你之前通过Spark直接将parquet文件写入S3路径的方式,不会被Iceberg表的元数据系统识别。Iceberg依赖自身的元数据(manifest文件、快照等)追踪表的数据文件,直接写入的文件不在Iceberg的元数据中,所以Glue的Compaction机制无法检测到这些小文件,自然不会触发操作。
解决方案步骤
1. 修改Glue Job,通过Iceberg API写入数据
更新你的Glue Job代码,使用Iceberg格式写入数据,确保元数据被正确维护。修改后的代码如下:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType from faker import Faker import random # Initialize Glue job context args = getResolvedOptions(sys.argv, ['JOB_NAME', 'database_name', 'table_name']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # Parameters database_name = args['database_name'] table_name = args['table_name'] # Define schema schema = StructType([ StructField("customer_id", IntegerType(), False), StructField("name", StringType(), True), StructField("age", IntegerType(), True), StructField("dob", DateType(), True), StructField("email", StringType(), True), StructField("address", StringType(), True), StructField("testcol", StringType(), True) ]) # Initialize Faker fake = Faker() def generate_data(num_records): """Generates sample data for testing.""" return [(random.randint(1, 10000), fake.name(), random.randint(18, 70), fake.date_of_birth(minimum_age=18), fake.email(), fake.address(), "test") for _ in range(num_records)] # 创建并写入多个小文件到Iceberg表 num_files = 150 # 生成超过50个文件,满足Compaction触发条件 records_per_file = 100 # 每个文件100条记录,确保单文件远小于128MB for _ in range(num_files): data = generate_data(records_per_file) df = spark.createDataFrame(data, schema=schema) # 使用Iceberg格式写入Glue表,自动维护元数据 df.write.format("iceberg") \ .mode("append") \ .option("database", database_name) \ .option("table", table_name) \ .save() job.commit()
注意:运行此Job前,确保你的Glue Job角色拥有访问Iceberg表的权限,以及S3存储桶的读写权限。
2. 检查Glue表的Compaction配置
进入Glue控制台,找到你的test_compaction表,进入优化标签页,确认以下配置:
- 已启用表优化(你已经完成这一步)
- 检查Compaction的触发阈值:默认是当某个分区下的小文件数超过50个且单文件小于128MB时触发。如果你的表没有分区,会检查整个表的文件情况。
- 确认优化角色的权限包含
glue:StartJobRun、s3:GetObject、s3:PutObject等必要权限
3. 手动触发Compaction(可选)
如果自动Compaction仍未触发,你可以手动触发测试:
在Glue控制台的表优化标签页,点击运行优化,选择Compaction类型,然后启动任务。这会立即触发一次Compaction操作,你可以在运行历史中查看进度。
4. 验证Compaction效果
Compaction完成后,你可以:
- 查看S3路径下的文件:会生成新的大文件(接近128MB),同时原来的小文件会被标记为已删除(Iceberg的快照机制不会立即删除旧文件,需要手动清理或启用过期快照)
- 在Athena中执行
SELECT * FROM msk_ingestion__qa_apps_us.test_compaction LIMIT 10,确认数据正常可访问
其他注意事项
- 确保你的Iceberg表是通过Glue或Athena正确创建的(你的CREATE TABLE语句是正确的)
- 避免直接修改S3路径下的文件,所有数据写入都要通过Iceberg的API(Spark、Athena、Glue Job),否则元数据会不一致
内容的提问来源于stack exchange,提问作者ErnieAndBert
相关产品推荐
相关产品推荐

