🎉 75% of content is free forever — Unlock Premium from $10/mo →
CW
NEWSLIVESearch All Content
💼 Servicesℹ️ About✉️ ContactView Pricing Plansfrom $10

Testing DAGs and Operators in Apache Airflow

🟢 Free Lesson

Advertisement

Testing DAGs and Operators

Testing PyramidUnit Tests (80%)Task logic, XCom, branchingIntegration Tests (15%)Hooks, connections, external systemsE2E Tests (5%)Full pipeline execution, scheduler, executor

Architecture Diagram

Formal Definitions

Detailed Explanation

DAG Validation Tests

Validate DAG structure and imports before deployment.


Key Validations:

TestDescription
Import SuccessAll DAGs load without import errors
DAG CountVerify expected number of DAGs
Required AttributesCheck dag_id, default_args, start_date
AcyclicityNo circular dependencies
Operator TypesUse only supported operators

Unit Testing Tasks

Test task functions in isolation with mocked dependencies.


Testing Strategy:

Test TypeDescriptionMock
Success PathNormal operationNone
Empty InputEdge case handlingData source
Connection ErrorExternal system failureHook/Connection
Invalid DataData validationData source

Mocking External Services

Isolate tests from external dependencies.


Common Mocks:

ServiceMock Target
DatabasePostgresHook, SqlAlchemy
S3/GCSS3Hook, GCSHook
APIsHttpHook, requests
Airflow ServicesVariable, Connection

def test_transform_data(self): """Test data transformation logic.""" from dags.my_dag import transform_data

input_data = [ {"name": "alice", "score": 85}, {"name": "bob", "score": 92}, ]

result = transform_data(input_data)

assert len(result) == 2 assert result[0]["name"] == "ALICE" # Uppercase assert result[0]["grade"] == "B" # Score-based grade

def test_transform_empty_data(self): """Test transformation with empty input.""" from dags.my_dag import transform_data

result = transform_data([]) assert result == []

def test_transform_invalid_data(self): """Test transformation handles invalid data gracefully.""" from dags.my_dag import transform_data

input_data = [{"name": None, "score": "invalid"}]

with pytest.raises(ValueError): transform_data(input_data)

class TestLoadTask: """Unit tests for load task function."""

@patch('dags.my_dag.get_db_connection') def test_load_success(self, mock_db): """Test successful data loading.""" from dags.my_dag import load_data

mock_conn = MagicMock() mock_db.return_value = mock_conn

data = [{"id": 1, "name": "test"}] result = load_data(data)

assert result['loaded'] == 1 mock_conn.execute.assert_called_once()

@patch('dags.my_dag.get_db_connection') def test_load_partial_failure(self, mock_db): """Test loading with partial failures.""" from dags.my_dag import load_data

mock_conn = MagicMock() mock_conn.execute.side_effect = [ None, # First insert succeeds Exception("Duplicate key"), # Second fails None, # Third succeeds ] mock_db.return_value = mock_conn

data = [{"id": 1}, {"id": 2}, {"id": 3}] result = load_data(data)

assert result['loaded'] == 2 assert result['failed'] == 1

Architecture Diagram

### Mocking External Systems

```python
# tests/test_with_mocks.py
import pytest
from unittest.mock import MagicMock, patch, AsyncMock
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.providers.postgres.hooks.postgres import PostgresHook

class TestWithMockedExternalSystems:
    """Test tasks with mocked external dependencies."""
    
    @patch('airflow.providers.amazon.aws.hooks.s3.S3Hook')
    def test_s3_upload(self, MockS3Hook):
        """Test S3 upload with mocked hook."""
        from dags.my_dag import upload_to_s3
        
        mock_hook = MagicMock()
        MockS3Hook.return_value = mock_hook
        mock_hook.load_file.return_value = True
        
        result = upload_to_s3(
            local_path='/tmp/file.csv',
            bucket='my-bucket',
            key='data/file.csv',
        )
        
        assert result is True
        mock_hook.load_file.assert_called_once_with(
            filename='/tmp/file.csv',
            key='data/file.csv',
            bucket_name='my-bucket',
            replace=True,
        )
    
    @patch('airflow.providers.postgres.hooks.postgres.PostgresHook')
    def test_database_operations(self, MockPostgresHook):
        """Test database operations with mocked hook."""
        from dags.my_dag import process_database
        
        mock_hook = MagicMock()
        MockPostgresHook.return_value = mock_hook
        mock_hook.get_pandas_df.return_value = MagicMock()
        
        result = process_database(query="SELECT * FROM users")
        
        assert result is not None
        mock_hook.get_pandas_df.assert_called_once()
    
    @patch('requests.get')
    def test_api_call(self, mock_get):
        """Test API call with mocked requests."""
        from dags.my_dag import fetch_from_api
        
        mock_response = MagicMock()
        mock_response.status_code = 200
        mock_response.json.return_value = {"data": [1, 2, 3]}
        mock_get.return_value = mock_response
        
        result = fetch_from_api(url="http://api.example.com/data")
        
        assert result == {"data": [1, 2, 3]}
        mock_get.assert_called_once()

Key Concepts Table

Test TypeScopeSpeedDependenciesCoverage
Unit TestSingle functionFast (~ms)Mocked80%
Integration TestMultiple componentsMedium (~s)Real services15%
E2E TestFull pipelineSlow (~min)Full stack5%
DAG ValidationStructureFast (~ms)NoneAll DAGs
Performance TestExecution timeSlow (~min)Load envCritical paths

Code Examples

Comprehensive Test Suite

CI/CD Test Configuration

# .github/workflows/test.yml
name: Airflow Tests

on:
  push:
    branches: [main, develop]
  pull_request:
    branches: [main]

jobs:
  lint:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v3
      - name: Set up Python
        uses: actions/setup-python@v4
        with:
          python-version: '3.10'
      - name: Install dependencies
        run: |
          pip install flake8 black isort mypy
      - name: Lint with flake8
        run: flake8 dags/ tests/
      - name: Check formatting with black
        run: black --check dags/ tests/
      - name: Sort imports
        run: isort --check-only dags/ tests/

  test:
    runs-on: ubuntu-latest
    needs: lint
    services:
      postgres:
        image: postgres:15
        env:
          POSTGRES_USER: airflow
          POSTGRES_PASSWORD: airflow
          POSTGRES_DB: airflow
        ports:
          - 5432:5432
    steps:
      - uses: actions/checkout@v3
      - name: Set up Python
        uses: actions/setup-python@v4
        with:
          python-version: '3.10'
      - name: Install Airflow
        run: |
          pip install apache-airflow==2.8.0
          pip install pytest pytest-cov
      - name: Run DAG validation
        run: pytest tests/test_dag_validation.py -v
      - name: Run unit tests
        run: pytest tests/test_tasks.py -v --cov=dags --cov-report=xml
      - name: Upload coverage
        uses: codecov/codecov-action@v3

  integration:
    runs-on: ubuntu-latest
    needs: test
    steps:
      - uses: actions/checkout@v3
      - name: Run integration tests
        run: |
          docker-compose -f docker-compose-test.yml up -d
          pytest tests/integration/ -v
          docker-compose -f docker-compose-test.yml down

Performance Metrics

Test Suite Performance

MetricTargetWarningCritical
Unit Test Duration< 5min5-10min> 10min
Integration Test Duration< 15min15-30min> 30min
Total CI Duration< 30min30-60min> 60min
Test Coverage> 80%70-80%< 70%
Flaky Test Rate< 1%1-5%> 5%

Test Distribution

Test CategoryCountDurationPriority
DAG Validation10-2010sP0
Unit Tests50-1002-5minP0
Integration Tests10-2010-15minP1
E2E Tests5-1020-30minP2

See Also

☆☆☆☆☆
0 ratings

Rate & Feedback

Need Expert Airflow Help?

Get personalized tutoring, project support, or professional consulting.

Advertisement