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

AWS Glue连接两张表时Job Bookmark增量加载失效问题排查

AWS Glue Job Bookmark增量加载失效问题排查与修复

问题概述

AWS Glue可视化ETL生成的脚本可正常运行,但Job Bookmark增量加载功能失效,每次运行都会全量上传数据至S3。已执行以下操作:

  • 启用Job Bookmark
  • 配置多列作为Bookmark键保证行唯一性
  • 手动修改可视化生成的脚本
  • 代码包含job.init()和job.commit()

问题根源分析

从提供的代码来看,主要存在以下几个导致增量失效的问题:

  1. Job Bookmark参数大小写错误
    代码中t_df的additional_options里使用了"JobBookmarkKeys"(大写J),而AWS Glue的正确参数名应为小写开头的"jobBookmarkKeys",参数名大小写不匹配会导致Glue无法识别Bookmark配置,直接失效。

  2. 关联表未配置Job Bookmark
    脚本中对mywarehouse_m表(m_df)未配置任何Bookmark参数,执行SQL join操作时,若其中一个表为全量读取,即使另一个表配置了增量,也会导致最终输出为全量数据——Glue无法追踪关联表的变化,只能重新处理所有数据。

  3. Spark SQL节点的Bookmark上下文传递风险
    使用Spark SQL将两个DynamicFrame join后,转换回DynamicFrame时需确保transformation_ctx唯一且正确,否则可能丢失Bookmark追踪信息。

修复方案

1. 修正Bookmark参数名大小写

修改t_df的additional_options中参数名,将"JobBookmarkKeys"改为"jobBookmarkKeys":

t_df = glueContext.create_dynamic_frame.from_options(
    connection_type = "oracle",
    connection_options = {
        "useConnectionProperties": "true",
        "dbtable": "mywarehouse_t",
        "connectionName": "warehouse",
        "hashfield":"id"
    },
    transformation_ctx = "t_node",
    additional_options = {
        "jobBookmarkKeys": ["id", "key2"],  # 修正参数名大小写
        "jobBookmarkKeysSortOrder": "asc"
    }   
)

2. 为关联表配置Job Bookmark(按需)

如果mywarehouse_m是需要增量加载的表(如维度表有更新),需为其添加Bookmark配置,选择合适的增量键(如唯一ID、更新时间戳):

m_df = glueContext.create_dynamic_frame.from_options(
    connection_type = "oracle",
    connection_options = {
        "useConnectionProperties": "true",
        "dbtable": "mywarehouse_m",
        "connectionName": "base2",
        "hashfield":"key2"
    },
    transformation_ctx = "m_node",
    additional_options = {
        "jobBookmarkKeys": ["key2", "update_time"],  # 替换为实际增量标识列
        "jobBookmarkKeysSortOrder": "asc"
    }
)

如果mywarehouse_m是静态维度表,无需增量加载,可添加缓存减少重复全量读取:

m_df = m_df.cache()

3. 确保Job Bookmark状态正常

  • 登录AWS Glue控制台,进入对应Job配置页面,确认Job Bookmark已启用
  • 如果之前的Job运行产生了错误的Bookmark记录,点击Reset Job Bookmark重置状态后再重新运行测试

4. 验证Spark SQL节点上下文

当前代码中SQL查询转换回DynamicFrame的transformation_ctx(sql_ctx)配置无误,无需修改。

修改后的完整脚本示例

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 awsgluedq.transforms import EvaluateDataQuality
from awsglue import DynamicFrame
import gs_now

def sparkSqlQuery(glueContext, query, mapping, transformation_ctx) -> DynamicFrame:
    for alias, frame in mapping.items():
        frame.toDF().createOrReplaceTempView(alias)
    result = spark.sql(query)
    return DynamicFrame.fromDF(result, glueContext, transformation_ctx)

args = getResolvedOptions(sys.argv, ['JOB_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# Default ruleset used by all target nodes with data quality enabled
DEFAULT_DATA_QUALITY_RULESET = """
    Rules = [
        ColumnCount >= max(last(1)),
        RowCount > avg(last(10))*0.6
    ]
"""

t_df = glueContext.create_dynamic_frame.from_options(
    connection_type = "oracle",
    connection_options = {
        "useConnectionProperties": "true",
        "dbtable": "mywarehouse_t",
        "connectionName": "warehouse",
        "hashfield":"id"
    },
    transformation_ctx = "t_node",
    additional_options = {
        "jobBookmarkKeys": ["id", "key2"],
        "jobBookmarkKeysSortOrder": "asc"
    }   
)

m_df = glueContext.create_dynamic_frame.from_options(
    connection_type = "oracle",
    connection_options = {
        "useConnectionProperties": "true",
        "dbtable": "mywarehouse_m",
        "connectionName": "base2",
        "hashfield":"key2"
    },
    transformation_ctx = "m_node",
    # 按需添加Bookmark配置
    additional_options = {
        "jobBookmarkKeys": ["key2", "update_time"],
        "jobBookmarkKeysSortOrder": "asc"
    }
)
# 静态表可添加缓存
# m_df = m_df.cache()

# Script generated for node SQL Query
sqlQuery = '''
select t.*, m.tag from mywarehouse_t t
inner join mywarehouse_m m on t.key2 = m.key2
'''

sql_node = sparkSqlQuery(glueContext, query = sqlQuery, mapping = {"mywarehouse_m":m_df , "mywarehouse_t":t_df}, 
transformation_ctx = "sql_ctx")

s3_node = glueContext.write_dynamic_frame.from_options(frame=sql_node, connection_type="s3", format="glueparquet", 
    connection_options={"path": "s3://my_s3", "partitionKeys": []}, 
    format_options={"compression": "snappy"}, 
    transformation_ctx="s3_node")

job.commit()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:29:51