测试 Airflow 代码

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

三层测试

DAG 测试金字塔

  • 完整性:DAG 能否无错误加载?
  • 单元:业务逻辑是否产出正确结果?
  • 集成:完整 DAG 能否端到端运行?
使用 Airflow 构建数据流水线

使用 DagBag 做完整性测试

from airflow.models import DagBag

dag_bag = DagBag(include_examples=False)
def test_no_import_errors(): assert len(dag_bag.import_errors) == 0
def test_dag_loaded(): assert "daily_etl" in dag_bag.dags
使用 Airflow 构建数据流水线

为何导入时 DAG 会出错

ModuleNotFoundError

  • ModuleNotFoundError:缺少 provider 包或导入路径错误
  • NameError:变量或函数被重命名,或拼写错误
  • ImportError:DAG 文件间循环导入
  • 完整性测试可避免以上问题
使用 Airflow 构建数据流水线

为任务函数做单元测试

在 DAG 文件中:

def clean_record(record):
    return {
        "name": record["name"].strip(),
        "email": record["email"].lower(),
    }

@task
def transform(records):
    return [clean_record(r)
            for r in records]
  • 抽取业务逻辑

在测试文件中:

from dags.data_cleaning import (
    clean_record,
)

def test_strips_whitespace():
    result = clean_record(
      {"name": "  Alice  ",
       "email": "[email protected]"}
    )
    assert result["name"] == "Alice"
  • 将单元测试聚焦于业务逻辑
使用 Airflow 构建数据流水线

使用 dag.test() 做集成测试

import pytest
from airflow.models import DagBag
from pendulum import datetime

dag_bag = DagBag(include_examples=False)

def test_etl_pipeline(): dag = dag_bag.get_dag("etl_output") assert dag is not None dag.test(logical_date=datetime(2026, 1, 15)) output = Path("/tmp/etl_results.json") assert output.exists() results = json.loads(output.read_text()) assert len(results) == 2
  • 在可控输入下运行真实 DAG,并校验输出
使用 Airflow 构建数据流水线

在 CI 中测试

CI pipeline

  • 完整性 + 单元:每次提交(速度快)
  • 集成:PR 或每晚运行(较慢)
  • 未通过三类测试的 DAG 不得上线
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...