如何为AWS Airflow调度的Glue作业关联Redshift连接?
问题解决:Airflow触发Glue作业时保留Redshift连接
核心原因
你遇到的问题是因为GlueJobOperator默认会在每次运行时,根据create_job_kwargs的配置创建或更新Glue作业,这会覆盖你手动添加的Redshift连接配置。要让连接持久显示在作业详情页,需要在Airflow的作业配置中直接指定该连接。
解决方案1:在Airflow配置中指定Redshift连接
在GlueJobOperator的create_job_kwargs参数中添加Connections字段,直接引用你在Glue控制台已创建好的Redshift连接名称。这样Airflow每次创建/更新作业时,都会自动带上这个连接,不会再被覆盖。
修改后的代码片段:
submit_glue_job = GlueJobOperator( task_id="submit_glue_job", job_name=f"{project}-{job_name}-{env}-{region_name}", iam_role_name="<role>", s3_bucket=f"s3://{project}-g-a-{env}-{region_name}", script_location=f"s3://{project}-g-a-{env}-{region_name}/scripts/{glue_script}", create_job_kwargs={ "GlueVersion": "3.0", "NumberOfWorkers": 2, "WorkerType": "G.1X", # 添加Redshift连接配置 "Connections": { "Connections": ["你的Redshift连接名称"] # 替换为实际的Glue连接名 } }, region_name=region_name, script_args={'--AWS_REGION': region_name, '--REDSHIFT_SECRET': '<secret>', '--REDSHIFT_DATABASE': '<db>', '--RS_DB_TABLE': '<table>', '--DYNAMODB_TABLE': '<dynamotable>'}, dag=dag )
注意:这里的连接名称必须是你已经在Glue控制台创建完成的Redshift连接,Airflow只是引用该已存在的资源。
解决方案2:预先创建Glue作业,Airflow仅触发运行
如果你希望完全由Glue控制台管理作业配置,可以预先创建好包含Redshift连接的作业,然后让Airflow只触发作业运行,不修改配置:
- 在Glue控制台创建好带Redshift连接的作业
- 在Airflow的
GlueJobOperator中设置update_config=False(需确保apache-airflow-providers-amazon版本≥3.0.0),禁止Airflow更新作业配置
修改后的代码片段:
submit_glue_job = GlueJobOperator( task_id="submit_glue_job", job_name=f"{project}-{job_name}-{env}-{region_name}", iam_role_name="<role>", region_name=region_name, update_config=False, # 关键:禁止更新现有作业配置 script_args={'--AWS_REGION': region_name, '--REDSHIFT_SECRET': '<secret>', '--REDSHIFT_DATABASE': '<db>', '--RS_DB_TABLE': '<table>', '--DYNAMODB_TABLE': '<dynamotable>'}, dag=dag )
避免第三方库的实践
在Glue脚本中使用Glue原生的DynamicFrame读写Redshift,完全不需要依赖psycopg2等第三方库:
- 读取Redshift示例:
from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext() glueContext = GlueContext(sc) # 通过Glue连接读取Redshift数据 redshift_dyf = glueContext.create_dynamic_frame.from_options( connection_type="redshift", connection_options={ "connectionName": "你的Redshift连接名称", "dbtable": "<table>", "database": "<db>" } )
- 写入Redshift示例:
glueContext.write_dynamic_frame.from_options( frame=redshift_dyf, connection_type="redshift", connection_options={ "connectionName": "你的Redshift连接名称", "dbtable": "<table>", "database": "<db>", "preactions": "TRUNCATE TABLE <table>" } )
内容的提问来源于stack exchange,提问作者mad4red
相关产品推荐
相关产品推荐

