TriggerDagRunOperator触发的DAG任务卡在队列的问题咨询
控制器与目标DAG代码及任务卡队列排查思路
控制器DAG代码(example_trigger_controller_dag.py)
import pendulum from airflow import DAG from airflow.operators.trigger_dagrun import TriggerDagRunOperator with DAG( dag_id="example_trigger_controller_dag", start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, schedule="@once", tags=["example"], ) as dag: trigger = TriggerDagRunOperator( task_id="test_trigger_dagrun", trigger_dag_id="example_trigger_target_dag", # 需与待触发的DAG的dag_id完全一致 conf={"message": "Hello World"}, )
目标DAG代码(example_trigger_target_dag.py)
import pendulum from airflow import DAG from airflow.decorators import task from airflow.operators.bash import BashOperator @task(task_id="run_this") def run_this_func(dag_run=None): """ 打印传递到DagRun conf属性中的"message"参数 :param dag_run: DagRun对象 """ print(f"远程接收到的message值为: {dag_run.conf.get('message')}") with DAG( dag_id="example_trigger_target_dag", start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, schedule=None, tags=["example"], ) as dag: run_this = run_this_func() bash_task = BashOperator( task_id="bash_task", bash_command='echo "Here is the message: $message"', env={"message": '{{ dag_run.conf.get("message") }}'}, )
任务卡队列排查思路
- 检查Airflow Worker状态:确认Worker节点是否正常运行,有无进程崩溃或资源耗尽情况。可通过Airflow UI的Admin > Workers查看,或执行对应命令:
- Celery executor:
airflow celery worker status - Sequential executor:
airflow jobs list(注意该模式同一时间仅能执行一个任务,若有其他任务占用则会阻塞)
- Celery executor:
- 确认目标DAG启用状态:在Airflow UI中检查
example_trigger_target_dag的开关是否开启,未启用的DAG不会被调度执行任务。 - 核对队列配置:若使用Celery executor,检查目标任务的
queue参数是否指向Worker正在监听的队列,避免任务被分配到不存在的队列导致无法被拾取。 - 排查资源瓶颈:查看Worker节点的CPU、内存使用率,若资源被占满,任务会卡在队列中等待资源释放。可通过
top、htop等工具监控系统资源。 - 查看日志信息:到Airflow UI中打开目标DAG的任务详情页,查看任务日志;同时检查Scheduler和Worker的日志文件,定位调度或执行环节的具体错误。
- 验证DAG依赖:确认目标DAG所需的Python包、数据库连接、外部服务等依赖是否配置正确且可用,缺失依赖可能导致任务无法启动。
- 检查权限配置:确保执行任务的角色拥有足够权限,包括读取DAG配置、执行任务、访问相关资源的权限。
内容的提问来源于stack exchange,提问作者Arda Burak Mamur
相关产品推荐
相关产品推荐

