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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 18:28:10