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

如何用Airflow动态任务触发Glue Job并按顺序执行

Airflow动态任务映射相关问题

问题描述

尝试使用Airflow动态任务映射功能时遇到以下问题:

  1. 用@task装饰run_glue_ingestion函数,生成的3个映射任务显示成功但未触发Glue Job;
  2. 改用@task_group装饰后可以触发Glue Job(因脚本问题失败),但多了一层任务层级,请问该分组是否必要?
  3. 当前映射任务为并发执行,如何改为顺序执行(等待一个GlueJobOperator完成后再启动下一个)?

现有代码

import json
import os
from datetime import datetime, timedelta

import boto3
import urllib3
from airflow.decorators import task, task_group
from airflow.models import DAG
import hvac
from airflow.operators.empty import EmptyOperator
from airflow.providers.amazon.aws.operators.glue import GlueJobOperator

default_args = {

}

DAG_ID = "dag-id"

with DAG(dag_id=DAG_ID,
         schedule=None,
         description="testing",
         default_args=default_args,
         start_date=datetime(2023, 8, 21),
         catchup=False) as dag:

    def get_credentials():
        # code to get aws creds

        return region_name, access_key, secret_key, session_token

    @task
    def processing():
        # this just gets some s3_keys by calling a lambda
        region_name, access_key, secret_key, session_token = get_credentials()
        session = boto3.Session(aws_access_key_id=access_key, aws_secret_access_key=secret_key,
                                aws_session_token=session_token, region_name=region_name)

        lambda_client = session.client('lambda')

        response = lambda_client.invoke(
            FunctionName='lambda_name',
            InvocationType='RequestResponse',
        )

        response_payload = json.loads(response['Payload'].read().decode('utf-8'))

        body = response_payload['body']

        print(type(body))
        return body.strip('][').split(', ')

    @task
    def run_glue_ingestion(s3_key):
        GlueJobOperator(
                task_id=f"test-job-{DAG_ID}",
                job_name="glue-job",
                job_desc="Glue test",
                script_location="s3_path",
                retries=0,
                region_name="us-east-1",
                iam_role_name="role_name",
                run_job_kwargs={
                    "SecurityConfiguration": "sec_config"
                },
                verbose=True,
                script_args={
                    "--s3_path": f"s3_path/{s3_key}",
                    "--environment": "dev",
                },
                num_of_dpus=2,
                aws_conn_id='aws-personal-conn'
            )

    values = processing()
    glue_ingestion_task = run_glue_ingestion.expand(s3_key=values)

    values >> glue_ingestion_task

问题解答

1. @task_group是否必要?

不必要。完全可以在不使用@task_group的前提下实现Glue任务的动态映射,只需修正@task装饰的函数写法即可。

2. 用@task装饰时未触发Glue Job的原因

在@task装饰的函数中,你仅仅实例化了GlueJobOperator对象,但没有调用它的execute方法。Airflow Operator的核心逻辑都在execute方法里,仅创建对象不会触发任何实际操作,所以任务会直接标记为成功,但Glue Job根本没被启动。

正确的写法有两种:

  • 方法一:在@task函数内调用Operator的execute方法(需传入上下文参数):
@task
def run_glue_ingestion(s3_key, context):
    glue_op = GlueJobOperator(
            task_id=f"test-job-{s3_key}",  # task_id需保证唯一,避免冲突
            job_name="glue-job",
            job_desc="Glue test",
            script_location="s3_path",
            retries=0,
            region_name="us-east-1",
            iam_role_name="role_name",
            run_job_kwargs={
                "SecurityConfiguration": "sec_config"
            },
            verbose=True,
            script_args={
                "--s3_path": f"s3_path/{s3_key}",
                "--environment": "dev",
            },
            num_of_dpus=2,
            aws_conn_id='aws-personal-conn'
        )
    glue_op.execute(context=context)
  • 方法二:直接使用GlueJobOperator的partial+expand组合,无需@task包装:
glue_ingestion_task = GlueJobOperator.partial(
        task_id="test-job",
        job_name="glue-job",
        job_desc="Glue test",
        script_location="s3_path",
        retries=0,
        region_name="us-east-1",
        iam_role_name="role_name",
        run_job_kwargs={
            "SecurityConfiguration": "sec_config"
        },
        verbose=True,
        script_args={
            "--environment": "dev",
        },
        num_of_dpus=2,
        aws_conn_id='aws-personal-conn'
    ).expand(script_args={"--s3_path": [f"s3_path/{key}" for key in values]})

注:partial用于定义固定参数,expand用于传入动态参数,task_id会自动添加索引后缀保证唯一。

3. 如何将映射任务改为顺序执行

Airflow动态映射默认并发执行,要改为顺序执行,只需限制该任务的同时运行实例数为1即可,有两种方式:

  • 方式一:在任务定义时设置max_active_tis_per_dagrun=1:
    如果用@task装饰:
@task(max_active_tis_per_dagrun=1)
def run_glue_ingestion(s3_key, context):
    # 逻辑同上

如果用GlueJobOperator.partial:

glue_ingestion_task = GlueJobOperator.partial(
        # 其他固定参数
        max_active_tis_per_dagrun=1
    ).expand(...)
  • 方式二:在DAG的default_args中设置max_active_tasks=1(会限制整个DAG的并发任务数,仅适用于整个DAG都需要顺序执行的场景):
default_args = {
    "max_active_tasks": 1
}

设置后,Airflow会等待当前映射任务完成后再启动下一个,实现顺序执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 06:34:55