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

本地启动Apache Airflow报错:无法获取可用DAG

Apache Airflow 2.1.4无法加载DAG:No viable dags retrieved问题排查

环境与问题概述

  • Airflow版本:2.1.4
  • DAG定义文件:dags/client.py
  • 配置文件:项目根目录下config/dag_configs.json
  • 已确认airflow.cfg中dags_folder配置路径正确,DAG定义结构符合最佳实践
  • 执行airflow db init后启动airflow webserver --port 8080,日志提示No viable dags retrieved,无法获取并执行DAG

代码片段

from airflow import DAG
from airflow.operators.dummy_operator import DummyOperator
from datetime import datetime, timedelta
from airflow.utils.task_group import TaskGroup
import json
from airflow.utils.dates import days_ago
import os

config_dir = os.path.join(os.path.dirname(__file__))
class TaskID:
    RunClientQuery = "run_client_query"
    CreateTempTable = "create_temp_table"
    GetTransferConfigJson = "get_transfer_config_json"
    CreateTransfer = "create_transfer"
    StartTransfer = "start_transfer"
    ExtractData = "extract_data"
    ExportToClient = "export_to_client"


DEFAULT_ARGS = {
    "owner": "delivery-squad",
    "depends_on_past": False,
    "wait_for_downstream": False,
    "email": [""],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 2,
    "retry_delay": timedelta(minutes=2),
    "catchup": True,
    'start_date': days_ago(2),
}

def create_dag(dag_id, schedule_interval, query_directory):
    dag = DAG(
        dag_id=f"{dag_id}_client_delivery",
        default_args=DEFAULT_ARGS,
        schedule_interval=schedule_interval,
        catchup=False,
    )

    run_client_query = DummyOperator(
        task_id=TaskID.RunClientQuery,
        dag=dag
    )
    create_temp_data = DummyOperator(
        task_id=TaskID.CreateTempTable,
        dag=dag
    )
    transfer_data =  DummyOperator(
        task_id=TaskID.StartTransfer,
        dag=dag)
    extract_data = DummyOperator(
        task_id=TaskID.ExtractData,
        dag=dag
    )
    export_to_client = DummyOperator(
        task_id=TaskID.ExportToClient,
        dag=dag
    )

    run_client_query >> create_temp_data >> transfer_data >> extract_data >> export_to_client

    return dag

config_file_path = os.path.join(config_dir, '../config/dag_configs.json')

with open(config_file_path) as f:
    config_data = json.load(f)

query_directory = "queries"  # Replace this with the actual path

for config in config_data:
    dag_id = config["dag_id"]
    schedule_interval = config["schedule_interval"]
    query_directory_placeholder = config["query_directory"]

    query_directory_actual = query_directory_placeholder.replace("{{QUERY_DIRECTORY}}", query_directory)

    globals()[dag_id] = create_dag(dag_id, schedule_interval, query_directory_actual)

错误日志

[2023-08-08 10:40:12,582] {processor.py:162} INFO - Started process (PID=27475) to work on ...
[2023-08-08 10:40:12,582] {processor.py:613} INFO - Processing file ... client_delivery_dag.py for tasks to queue
[2023-08-08 10:40:12,583] {dagbag.py:496} INFO - Filling up the DagBag from ... client_delivery_dag.py
[2023-08-08 10:40:12,585] {processor.py:625} WARNING - No viable dags retrieved from ... client_delivery_dag.py

排查与解决方案

1. 修复配置文件路径问题

当前代码中路径拼接逻辑可能因Airflow运行工作目录不同导致文件找不到,替换为绝对路径拼接:

# 从当前文件向上两级到项目根目录,再定位config文件
project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
config_file_path = os.path.join(project_root, 'config', 'dag_configs.json')

添加验证代码,确认文件可访问:

print(f"Config file path: {config_file_path}")
print(f"File exists: {os.path.exists(config_file_path)}")

执行airflow dags list查看输出,确认文件路径是否正确。

2. 校验dag_configs.json格式

确保配置文件是有效的JSON数组,且每个配置项包含必填字段:
示例正确格式:

[
  {
    "dag_id": "client_dag_1",
    "schedule_interval": "@daily",
    "query_directory": "{{QUERY_DIRECTORY}}/client1"
  }
]

检查是否存在语法错误、字段缺失或拼写错误。

3. 确认DAG注册逻辑

在循环中添加打印,验证DAG是否被正确创建并注册:

for config in config_data:
    dag_id = config["dag_id"]
    schedule_interval = config["schedule_interval"]
    query_directory_placeholder = config["query_directory"]

    query_directory_actual = query_directory_placeholder.replace("{{QUERY_DIRECTORY}}", query_directory)
    dag = create_dag(dag_id, schedule_interval, query_directory_actual)
    print(f"Created DAG: {dag.dag_id}")
    globals()[dag_id] = dag

执行airflow dags list,若能看到打印的DAG ID,说明注册成功;否则检查循环逻辑是否正常执行。

4. 排查代码语法与运行错误

直接执行DAG文件,检查是否有语法或运行时错误:

python dags/client.py

若出现JSON解析错误、导入错误等,修复对应问题。

5. 检查文件权限

确保Airflow运行用户拥有client.py和dag_configs.json的可读权限:

chmod +r dags/client.py
chmod +r config/dag_configs.json

6. 查看详细日志

开启DEBUG日志级别,获取更多加载细节:

airflow webserver --port 8080 --log-level DEBUG

根据日志中的错误信息定位问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:44:59