AWS Glue连接两张表时Job Bookmark增量加载失效问题排查
问题概述
AWS Glue可视化ETL生成的脚本可正常运行,但Job Bookmark增量加载功能失效,每次运行都会全量上传数据至S3。已执行以下操作:
- 启用Job Bookmark
- 配置多列作为Bookmark键保证行唯一性
- 手动修改可视化生成的脚本
- 代码包含
job.init()和job.commit()
问题根源分析
从提供的代码来看,主要存在以下几个导致增量失效的问题:
Job Bookmark参数大小写错误
代码中t_df的additional_options里使用了"JobBookmarkKeys"(大写J),而AWS Glue的正确参数名应为小写开头的"jobBookmarkKeys",参数名大小写不匹配会导致Glue无法识别Bookmark配置,直接失效。关联表未配置Job Bookmark
脚本中对mywarehouse_m表(m_df)未配置任何Bookmark参数,执行SQL join操作时,若其中一个表为全量读取,即使另一个表配置了增量,也会导致最终输出为全量数据——Glue无法追踪关联表的变化,只能重新处理所有数据。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

