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

Cloud Composer(Airflow) test_dag循环创建任务时调度被跳过问题求助

问题场景

在Cloud Composer(Airflow)中运行test_dag时,通过循环批量创建SQL执行任务会触发调度日志报错:DAG test_dag scheduling was skipped, probably because the DAG record was locked,一段时间后DAG直接失败;但手动按顺序逐个创建任务时一切正常。当前使用单个调度器Pod,生产环境无法直接重启调度器。

引发问题的代码片段

...

directory_path = '/home/airflow/gcs/data/test_dag/sql'
file_contents = {}
for filename in filenames:
    try:
        print(f"Processing file: {filename}")
        with open(os.path.join(directory_path, filename), 'r') as file:
            template = Template(file.read())
            query = template.render(**env)
            file_contents[filename] = query
    except Exception as e:
        print(f"Error processing file {filename}: {e}")


taskA >> taskB
previous_task = bq_streaming_buffer_wait
for filename, content in file_contents.items():
    try:
        execute_query_task = PythonOperator(
            task_id=f'execute_query_{filename.replace(".sql", "")}',
            python_callable=execute_query,
            op_args=[filename, content],
            provide_context=True,
            trigger_rule=TriggerRule.NONE_SKIPPED
        )
        taskB >> taskC
    except Exception as e:
        print(f"Error creating task for file {filename}: {e}")

代码层面问题定位

  1. 任务依赖逻辑错误:循环内重复执行taskB >> taskC,而非将新创建的execute_query_task加入依赖链。这会导致所有动态任务没有被正确挂载到DAG的依赖关系中,Airflow调度器在解析DAG时会陷入异常状态,引发锁冲突。
  2. 未使用previous_task构建链式依赖:代码中定义了previous_task = bq_streaming_buffer_wait但未在循环中使用,导致动态生成的任务没有形成正确的执行序列,调度器无法识别合法的任务流,进而触发DAG锁死。

代码修复方案

修改循环内的依赖构建逻辑,将动态生成的任务正确加入依赖链:

...

# 先构建基础依赖链
taskA >> taskB >> bq_streaming_buffer_wait
previous_task = bq_streaming_buffer_wait
for filename, content in file_contents.items():
    try:
        execute_query_task = PythonOperator(
            task_id=f'execute_query_{filename.replace(".sql", "")}',
            python_callable=execute_query,
            op_args=[filename, content],
            provide_context=True,
            trigger_rule=TriggerRule.NONE_SKIPPED
        )
        # 将当前任务挂载到前序任务之后,并更新前序任务指针
        previous_task >> execute_query_task
        previous_task = execute_query_task
    except Exception as e:
        print(f"Error creating task for file {filename}: {e}")

# 所有动态SQL任务执行完成后,再执行taskC
previous_task >> taskC

调度器配置调优(无需重启)

如果代码修复后仍存在锁相关问题,可以调整以下Airflow配置(在Cloud Composer的环境配置中修改):

  • scheduler_zombie_task_threshold:延长任务被标记为僵尸任务的阈值,默认300秒,可调整为600秒,避免调度器误判任务状态引发锁竞争。
  • scheduler_min_file_process_interval:增加DAG文件的扫描间隔,默认30秒,可调整为60秒,减少调度器频繁解析DAG带来的锁冲突。
  • dag_file_processor_timeout:延长DAG文件处理器的超时时间,默认50秒,可调整为120秒,给复杂DAG解析足够时间,避免中途超时导致锁未释放。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 14:21:13