如何用Airflow动态任务触发Glue Job并按顺序执行
Airflow动态任务映射相关问题
问题描述
尝试使用Airflow动态任务映射功能时遇到以下问题:
- 用
@task装饰run_glue_ingestion函数,生成的3个映射任务显示成功但未触发Glue Job; - 改用
@task_group装饰后可以触发Glue Job(因脚本问题失败),但多了一层任务层级,请问该分组是否必要? - 当前映射任务为并发执行,如何改为顺序执行(等待一个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
相关产品推荐
相关产品推荐

