在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自身的增量能力实现同步:
- 确保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
- 修改脚本,用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
相关产品推荐
相关产品推荐

