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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 15:15:32