Airflow动态任务expand行为异常:输入类型与执行次数不符问题
Airflow动态任务expand行为差异问题
运行recon_rule_setup任务时,每次会从上游read_recon_config任务接收Dictionary类型的recon_conf作为输入;但运行recon_rule_exec任务时,却从上游任务接收List类型的recon_rule作为输入。
预期是recon_rule_setup运行2次,recon_rule_exec运行4次(取决于返回值),为何expand的表现会有这种差异?
from datetime import datetime from airflow.models import DAG, XCom from airflow.utils.dates import days_ago from airflow.decorators import task from airflow.utils.db import provide_session from airflow.operators.bash import BashOperator from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup @provide_session def clear_xcom_data(session=None, **kwargs): dag_instance = kwargs["dag"] dag_id = dag_instance._dag_id session.query(XCom).filter(XCom.dag_id == dag_id).delete() @task(task_id="read_recon_config") def read_recon_config(dag_run=None): parent_dict = dag_run.conf d1 = {"name": "Santanu"} d2 = {"name": "Ghosh"} return [d1, d2] @task(task_id="recon_rule_setup") def recon_rule_setup(recon_conf): print(f"type of recon_conf_dict: {type(recon_conf)}") print(f"recon_conf_dict: {recon_conf}") return [recon_conf, {"name": "Kolkata"}] @task(task_id="recon_rule_exec") def recon_rule_exec(recon_rule, master_key): print(f"master_key type: {type(master_key)}") print(f"master_key: {master_key}") print(f"recon_rule type: {type(recon_rule)}") print(f"recon_rule: {recon_rule}") default_args = { 'owner': 'Airflow', 'start_date': days_ago(1), 'depends_on_past': False, 'email_on_failure': False, 'email_on_retry': False, 'retries': 1 } dag_name_id = "dynamic_demo" cur_datetime = datetime.utcnow().strftime("%Y%m%d%H%M%S%f")[:-3] master_dag_key = f"{dag_name_id}_{cur_datetime}" with DAG( dag_id=dag_name_id, default_args=default_args, schedule_interval=None, catchup=False ) as dag: start = BashOperator(task_id="start", bash_command='echo "starting reconciliation"', do_xcom_push=False) stop = BashOperator(task_id="stop", bash_command='echo "stopping reconciliation"', do_xcom_push=False) delete_xcom = PythonOperator( task_id="delete_xcom", python_callable=clear_xcom_data ) with TaskGroup(group_id="reconciliation_process") as tg1: recon_config_list = read_recon_config() recon_rule_list = recon_rule_setup.expand(recon_conf=recon_config_list) recon_rule_exec.partial(master_key=master_dag_key).expand(recon_rule=recon_rule_list) start >> tg1 >> delete_xcom >> stop
内容的提问来源于stack exchange,提问作者Santanu Ghosh
相关产品推荐
相关产品推荐

