Airflow コードのテスト

Airflow によるデータパイプラインの構築

Volker Janz

Senior Developer Advocate at Astronomer

テストの3つのレベル

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: プロバイダーパッケージが不足しているか、インポートパスが誤っている
  • 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 パイプライン

  • 整合性+単体: コミットのたびに実行(高速)
  • 統合: プルリクエストまたは夜間実行(低速)
  • 3つすべてに合格しない限り、DAG は本番環境に到達しない
Airflow によるデータパイプラインの構築

練習しましょう!

Airflow によるデータパイプラインの構築

Preparing Video For Download...