在AWS Glue中Spark处理数据写入Iceberg表时会话启动失败
问题排查与解决方案
核心问题分析
你遇到的Session failed to reach READY instead reaching terminal state FAILED错误,结合Parquet写入正常、Iceberg写入失败的现象,大概率是Iceberg配置冲突、权限缺失或代码逻辑遗漏导致的会话初始化失败。
针对性修复步骤
1. 移除重复的Spark配置
代码中同时用%%configure和spark.conf.set重复设置Iceberg配置,会导致会话初始化时的配置冲突。保留%%configure的配置,删除代码中的spark.conf.set块:
需要删除的代码段:
# Set Spark configurations for Iceberg spark.conf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") spark.conf.set("spark.sql.catalog.glue_catalog", "org.apache.iceberg.spark.SparkCatalog") spark.conf.set("spark.sql.catalog.glue_catalog.warehouse", "s3://s3-bucket-fabric/") spark.conf.set("spark.sql.catalog.glue_catalog.catalog-impl", "org.apache.iceberg.aws.glue.GlueCatalog") spark.conf.set("spark.sql.catalog.glue_catalog.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
2. 修复Job初始化逻辑
代码直接调用job.commit()但未初始化Job对象,会导致作业提交失败进而中断会话。需添加Job初始化代码:
# 初始化Job(添加在spark = glueContext.spark_session之后) args = getResolvedOptions(sys.argv, ['JOB_NAME']) job = Job(glueContext) job.init(args['JOB_NAME'], args)
3. 验证Iceberg表前置条件
- 确保目标Iceberg表
transactions-db.transactions-iceberg已存在,或添加createTable选项自动创建:df.write.format("iceberg")\ .mode("overwrite")\ .option("createTable", "true")\ .save("glue_catalog.transactions-db.transactions-iceberg") - 检查S3仓库路径
s3://s3-bucket-fabric/权限:Glue角色需具备该路径的s3:PutObject、s3:GetObject、s3:ListBucket权限,同时具备Glue Data Catalog的CreateTable、UpdateTable权限。
4. 调整Glue会话的Iceberg版本兼容性
Glue 4.0+对Iceberg支持更稳定,若使用旧版本建议切换到Glue 4.0,并在%%configure中补充配置:
{ "conf": { "spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions", "spark.sql.catalog.glue_catalog": "org.apache.iceberg.spark.SparkCatalog", "spark.sql.catalog.glue_catalog.warehouse": "s3://s3-bucket-fabric/", "spark.sql.catalog.glue_catalog.catalog-impl": "org.apache.iceberg.aws.glue.GlueCatalog", "spark.sql.catalog.glue_catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO", "spark.sql.catalog.glue_catalog.glue.skip-name-validation": "true" }, "args": [ "--datalake-format", "iceberg", "--enable-spark-glue-datasource", "true" ], "glueVersion": "4.0", "workerType": "G.1X", "numberOfWorkers": 2 }
5. 排查会话失败的详细日志
在Glue控制台查看会话日志定位具体异常:
- 进入Glue Studio → 选择对应Notebook → 点击"查看日志"
- 重点搜索
ClassNotFoundException(依赖缺失)、AccessDeniedException(权限问题)这类具体报错。
完整修复后的代码示例
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 import DataFrame from pyspark.sql.functions import col, regexp_replace, to_timestamp, to_date, when # Initialize GlueContext glueContext = GlueContext(SparkContext.getOrCreate()) spark = glueContext.spark_session # 初始化Job args = getResolvedOptions(sys.argv, ['JOB_NAME']) job = Job(glueContext) job.init(args['JOB_NAME'], args) # Load the data from the Glue Data Catalog (Lake Formation database and table) datasource = glueContext.create_dynamic_frame.from_catalog(database="transactions-db", table_name="raw", transformation_ctx="datasource") # Define a function to convert DynamicFrame to DataFrame def dynamic_frame_to_df(dynamic_frame): return dynamic_frame.toDF() # Convert to Spark DataFrame df = dynamic_frame_to_df(datasource) # Define a list of columns to clean columns_to_clean = df.columns # Example Cleaning and Formatting Transformations df = df.withColumn("Amount", regexp_replace(col("Amount"), "[^0-9.]", "").cast("float")) df = df.withColumn("Time", to_timestamp(col("Time"), "H:m")) df = df.withColumn("Date", to_date(col("Year") + "-" + col("Month") + "-" + col("Day"))) df = df.withColumn("MerchantName", regexp_replace(col("Merchant Name"), "[^a-zA-Z0-9 ]", "")) df = df.withColumn("TransactionType", when(col("Use Chip") == 'Swipe Transaction', 'Physical').otherwise('Online')) # Remove extra quotation marks from the "errors?" column values df = df.withColumn("errors?", regexp_replace(col("errors?"), '"', '')) # Remove unnecessary special characters and fill empty rows with "N/A" for column in columns_to_clean: df = df.withColumn(column, regexp_replace(col(column), '[?!"#$%&\'()*+,./:;<=>@\[\\]^_`{|}~]', '')) df = df.withColumn(column, when(col(column).isNull() | (col(column) == ''), 'N/A').otherwise(col(column))) # Write the transformed data to the Iceberg table (自动创建表) df.write.format("iceberg")\ .mode("overwrite")\ .option("createTable", "true")\ .save("glue_catalog.transactions-db.transactions-iceberg") # Commit the job job.commit()
内容的提问来源于stack exchange,提问作者user31081998
相关产品推荐
相关产品推荐

