使用AWS Glue DataFrame API配置Iceberg表分桶分区报错,求指导
解决AWS Glue DataFrame API写入Iceberg表时的Bucket分区错误
问题原因
你遇到的AnalysisException: Invalid partition transformation: bucket(3, age__pk3)错误,是因为Glue DataFrame API的partitionedBy方法无法直接通过expr()识别Iceberg的bucket分区转换语法,需要使用Iceberg Spark集成提供的专用转换函数。
解决方案
修改代码,使用Iceberg官方提供的SparkExpressions.bucket函数来定义分桶分区,具体步骤如下:
- 导入必要的依赖类和函数
- 在
partitionedBy中替换expr()为SparkExpressions.bucket()调用
修正后的完整代码
# ===================================================== # 🧊 Step 4. Write Data to Iceberg Table (Glue Catalog) # ===================================================== from pyspark.sql.functions import col from org.apache.iceberg.spark.expressions import SparkExpressions table_name = "glue_catalog.cmt_test_db.iceberg_table" ( df.writeTo(table_name) .using("iceberg") .tableProperty("format-version", "2") .partitionedBy( col("city__pk1"), SparkExpressions.bucket(3, col("age__pk3")), SparkExpressions.bucket(3, col("id__pk2")), ) .createOrReplace() ) print(f"✅ Iceberg table '{table_name}' created/loaded successfully!")
替代方案(字符串形式定义分区转换)
如果上述方式遇到类导入问题,也可以直接传递字符串格式的Iceberg分区转换规则给partitionedBy,代码如下:
# ===================================================== # 🧊 Step 4. Write Data to Iceberg Table (Glue Catalog) # ===================================================== table_name = "glue_catalog.cmt_test_db.iceberg_table" ( df.writeTo(table_name) .using("iceberg") .tableProperty("format-version", "2") .partitionedBy( "city__pk1", "bucket(3, age__pk3)", "bucket(3, id__pk2)", ) .createOrReplace() ) print(f"✅ Iceberg table '{table_name}' created/loaded successfully!")
验证说明
两种方式都能正确映射Iceberg的分桶分区规则,和你在Athena SQL中使用的逻辑一致。执行后可以通过Athena查询表的分区信息,确认分桶配置生效。
内容的提问来源于stack exchange,提问作者Mohammed Suhail
相关产品推荐
相关产品推荐

