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

在AWS Glue作业中读取Iceberg表时遇作业书签格式不支持错误求助

问题分析与解决

错误原因

AWS Glue Job Bookmarks目前不支持Iceberg格式的DynamicFrame。你通过glueContext.create_dynamic_frame.from_catalog读取Iceberg表时,开启Job Bookmarks会触发该错误——因为Glue的书签机制无法识别Iceberg的表格式,无法追踪增量数据状态。

解决方案

方案1:关闭Job Bookmarks(无需增量同步场景)

如果任务不需要增量同步,直接关闭Job Bookmarks即可:

  • 在Glue控制台的Job配置页,找到"Job bookmarks"选项,设置为Disable
  • 或在脚本开头初始化Job时显式关闭:
from awsglue.context import GlueContext
from pyspark.context import SparkContext

sc = SparkContext()
glueContext = GlueContext(sc)
job = glueContext.job
job.init(args['JOB_NAME'], bookmarks_enabled=False)

方案2:改用Spark SQL读取Iceberg并手动实现增量同步(需要增量场景)

若需增量处理,放弃Glue原生Job Bookmarks,改用Spark SQL直接读取Iceberg表,借助Iceberg自身的增量能力实现同步:

  1. 确保Spark配置正确(注意你之前的配置需统一添加--conf前缀):
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
--conf spark.sql.catalog.glue_catalog=org.apache.iceberg.spark.SparkCatalog
--conf spark.sql.catalog.glue_catalog.warehouse=s3://veeva-od-delta-sync/cn_lastest/hco_data/
--conf spark.sql.catalog.glue_catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog
--conf spark.sql.catalog.glue_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO
--conf spark.sql.catalog.glue_catalog.lock-impl=org.apache.iceberg.aws.glue.DynamoLockManager
--conf spark.sql.catalog.glue_catalog.lock.table=hco_data
  1. 修改脚本,用Spark SQL读取Iceberg表并手动维护增量状态:
from awsglue.context import GlueContext
from pyspark.context import SparkContext
import boto3

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session

# 从S3读取上次处理的快照ID(示例用S3存储状态,也可改用Glue参数存储)
s3 = boto3.client('s3')
last_snapshot_id = None
try:
    response = s3.get_object(Bucket='veeva-od-delta-sync', Key='cn_lastest/hco_data/last_snapshot.txt')
    last_snapshot_id = response['Body'].read().decode('utf-8')
except s3.exceptions.NoSuchKey:
    pass

# 读取增量/全量数据
if last_snapshot_id:
    df = spark.sql(f"SELECT * FROM glue_catalog.cn_lastest.hco_data WHERE snapshot_id > {last_snapshot_id}")
else:
    df = spark.sql("SELECT * FROM glue_catalog.cn_lastest.hco_data")

# 转换为DynamicFrame并执行业务处理(替换为你的逻辑)
dynamic_frame = glueContext.create_dynamic_frame.from_df(df, glueContext, "hco_dynamic_dataFrame")

# 保存最新快照ID,供下次增量使用
latest_snapshot = spark.sql("SELECT MAX(snapshot_id) AS latest_id FROM glue_catalog.cn_lastest.hco_data.history").collect()[0]['latest_id']
s3.put_object(Bucket='veeva-od-delta-sync', Key='cn_lastest/hco_data/last_snapshot.txt', Body=str(latest_snapshot))

job.commit()

关键说明

  • Glue Job Bookmarks仅支持Glue原生兼容的数据源格式(如Parquet、CSV等),Iceberg作为外部表格式不在支持范围内,因此无法通过create_dynamic_frame.from_catalog结合Bookmarks使用。
  • Iceberg本身提供了快照、时间旅行、增量扫描等完善的增量能力,可通过Spark SQL直接利用这些特性,自行维护增量状态替代Glue的Job Bookmarks。

内容的提问来源于stack exchange,提问作者Jing Wang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 07:42:37