Airflow Dag Patterns

作者 wshobson46891e7e60da无许可证收录于 2026年10月8日更新于 2026年10月8日

Build production Apache Airflow DAGs with best practices for operators, sensors, testing, and deployment. Use when creating data pipelines, orchestrating workflows, or scheduling batch jobs.

AI 生成的概览

指导构建生产级 Apache Airflow DAG,涵盖算子、传感器、测试与部署模式。

功能
提供面向生产环境的 Apache Airflow DAG 设计指导,涵盖 DAG 设计原则、任务依赖、算子、传感器、测试与部署策略。其中包含快速入门示例 DAG 以及最佳实践的正反建议,更详细的模式文档放在参考文件中。它产出的是说明与代码示例,本身不执行任何操作。
适用场景
适用于创建数据管道编排、设计 DAG 结构与依赖、实现自定义算子或传感器、在本地测试 DAG、在生产环境搭建 Airflow,或调试失败的 DAG 运行。
运行要求
不附带脚本,仅为说明文档。参照示例需要 Apache Airflow 环境和 Python,但阅读这些指导本身不需要任何依赖。

Apache Airflow DAG Patterns

Production-ready patterns for Apache Airflow including DAG design, operators, sensors, testing, and deployment strategies.

When to Use This Skill

  • Creating data pipeline orchestration with Airflow
  • Designing DAG structures and dependencies
  • Implementing custom operators and sensors
  • Testing Airflow DAGs locally
  • Setting up Airflow in production
  • Debugging failed DAG runs

Core Concepts

1. DAG Design Principles

PrincipleDescription
IdempotentRunning twice produces same result
AtomicTasks succeed or fail completely
IncrementalProcess only new/changed data
ObservableLogs, metrics, alerts at every step

2. Task Dependencies

python
# Lineartask1 >> task2 >> task3
# Fan-outtask1 >> [task2, task3, task4]
# Fan-in[task1, task2, task3] >> task4
# Complextask1 >> task2 >> task4task1 >> task3 >> task4

Quick Start

python
# dags/example_dag.pyfrom datetime import datetime, timedeltafrom airflow import DAGfrom airflow.operators.python import PythonOperatorfrom airflow.operators.empty import EmptyOperator
default_args = {    'owner': 'data-team',    'depends_on_past': False,    'email_on_failure': True,    'email_on_retry': False,    'retries': 3,    'retry_delay': timedelta(minutes=5),    'retry_exponential_backoff': True,    'max_retry_delay': timedelta(hours=1),}
with DAG(    dag_id='example_etl',    default_args=default_args,    description='Example ETL pipeline',    schedule='0 6 * * *',  # Daily at 6 AM    start_date=datetime(2024, 1, 1),    catchup=False,    tags=['etl', 'example'],    max_active_runs=1,) as dag:
    start = EmptyOperator(task_id='start')
    def extract_data(**context):        execution_date = context['ds']        # Extract logic here        return {'records': 1000}
    extract = PythonOperator(        task_id='extract',        python_callable=extract_data,    )
    end = EmptyOperator(task_id='end')
    start >> extract >> end

Detailed patterns and worked examples

Detailed pattern documentation lives in references/details.md. Read that file when the navigation tier above is insufficient.

Best Practices

Do's

  • Use TaskFlow API - Cleaner code, automatic XCom
  • Set timeouts - Prevent zombie tasks
  • Use mode='reschedule' - For sensors, free up workers
  • Test DAGs - Unit tests and integration tests
  • Idempotent tasks - Safe to retry

Don'ts

  • Don't use depends_on_past=True - Creates bottlenecks
  • Don't hardcode dates - Use {{ ds }} macros
  • Don't use global state - Tasks should be stateless
  • Don't skip catchup blindly - Understand implications
  • Don't put heavy logic in DAG file - Import from modules

来源与署名

来源:wshobson/agents位于plugins/data-engineering/skills/airflow-dag-patterns提交46891e7

许可证: 无许可证

内容归原作者所有。SourceWeft 从公开仓库中收录这些内容。

举报或申请下架