Airflow 1.10中持续检查Glue Job状态并按序执行任务的方案咨询
解决Airflow 1.10中Python Operator等待Glue Job完成的问题
核心思路
在Python Operator的执行函数里,用boto3提交Glue Job后,循环轮询Job的运行状态,直到Job进入终端状态。只有当Job成功完成时,函数才正常返回(Airflow标记任务成功);如果Job失败、超时或被终止,直接抛出异常让Airflow标记任务失败,中断后续依赖任务的执行。
具体实现步骤
- 初始化boto3的Glue客户端,配置目标区域
- 编写通用的Glue Job执行+状态轮询函数,包含提交Job、循环查状态、终端状态判断逻辑
- 将函数绑定到Python Operator,并设置任务间的依赖顺序
代码示例
import boto3 import time from airflow.models import DAG from airflow.operators.python_operator import PythonOperator from datetime import datetime # 初始化Glue客户端,替换为你的AWS区域 glue_client = boto3.client('glue', region_name='us-east-1') def run_glue_job_wait_completion(job_name): # 提交Glue Job submit_response = glue_client.start_job_run(JobName=job_name) job_run_id = submit_response['JobRunId'] # 轮询配置:可根据Job实际运行时长调整 poll_interval = 30 # 每30秒查一次状态 max_wait_minutes = 30 # 最长等待30分钟 max_retries = int(max_wait_minutes * 60 / poll_interval) for _ in range(max_retries): time.sleep(poll_interval) # 查询当前Job运行状态 run_details = glue_client.get_job_run(JobName=job_name, RunId=job_run_id) current_status = run_details['JobRun']['JobRunState'] # 处理终端状态 if current_status == 'SUCCEEDED': print(f"Glue Job [{job_name}] 执行成功,Run ID: {job_run_id}") return elif current_status in ['FAILED', 'TIMEOUT', 'STOPPED']: error_msg = run_details['JobRun'].get('ErrorMessage', '无详细错误信息') raise Exception(f"Glue Job [{job_name}] 执行失败,状态: {current_status},错误信息: {error_msg}") # 非终端状态(如RUNNING、STARTING)继续轮询 # 超过最长等待时间仍未完成 raise Exception(f"Glue Job [{job_name}] 执行超时,已等待{max_wait_minutes}分钟") # 定义DAG default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'retries': 1 } with DAG( 'sequential_glue_jobs', default_args=default_args, schedule_interval='@daily', catchup=False ) as dag: # 第一个Glue任务 glue_task_1 = PythonOperator( task_id='run_glue_job_1', python_callable=run_glue_job_wait_completion, op_kwargs={'job_name': 'your-first-glue-job'} ) # 第二个Glue任务 glue_task_2 = PythonOperator( task_id='run_glue_job_2', python_callable=run_glue_job_wait_completion, op_kwargs={'job_name': 'your-second-glue-job'} ) # 设置任务执行顺序 glue_task_1 >> glue_task_2
关键注意事项
- 权限配置:确保Airflow运行的服务角色拥有Glue的
StartJobRun和GetJobRun权限,以及Glue Job依赖的S3、IAM等相关权限。 - 轮询参数调整:根据你的Glue Job平均运行时长,修改
poll_interval和max_wait_minutes,避免过早超时或不必要的频繁查询。 - 异常处理:抛出异常后,Airflow会将当前任务标记为失败,后续依赖任务不会触发,符合DAG的错误流转逻辑。
- 代码复用:通过
op_kwargs传递不同的Job名称,复用同一个轮询函数,减少重复代码。
内容的提问来源于stack exchange,提问作者cknolhoff
相关产品推荐
相关产品推荐

