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}")
代码层面问题定位
- 任务依赖逻辑错误:循环内重复执行
taskB >> taskC,而非将新创建的execute_query_task加入依赖链。这会导致所有动态任务没有被正确挂载到DAG的依赖关系中,Airflow调度器在解析DAG时会陷入异常状态,引发锁冲突。 - 未使用
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
相关产品推荐
相关产品推荐

