为什么需要工作流编排?

上一章我们介绍了 Airbyte 和 dbt,它们各自负责 ETL 的不同环节。但一个完整的生产环境数据管道包含多个步骤:Airbyte 同步、dbt 运行、数据质量检查、异常告警、失败重试。这些步骤需要按特定顺序执行,并且要每天自动运行。这就是工作流编排工具的职责。

没有编排工具时,你可能用 cron 脚本把一堆命令串起来:

# cron 方案,很脆弱
0 2 * * * /usr/bin/airbyte sync-orders && /usr/bin/dbt run && /usr/bin/python check_quality.py
[/code]

这种方案的问题很明显:没有任务依赖管理,失败无法自动重试,没有可视化的执行状态,难以扩展。当你有几十个数据管道时,cron 脚本很快就会变成一团乱麻。

Apache Airflow 就是为了解决这些问题而生的。

## Apache Airflow 核心概念

### DAG(有向无环图)

DAG 是 Airflow 中最核心的概念。一个 DAG 定义了一个完整的工作流,包含所有任务以及它们之间的依赖关系。"有向无环"意味着任务之间有方向性的依赖,且不能形成循环依赖。

一个最简单的 DAG

from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime

def extract(): print(“抽取数据…”)

def transform(): print(“转换数据…”)

def load(): print(“加载数据…”)

with DAG( dag_id=“simple_etl”, start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False ) as dag: extract_task = PythonOperator( task_id=“extract”, python_callable=extract ) transform_task = PythonOperator( task_id=“transform”, python_callable=transform ) load_task = PythonOperator( task_id=“load”, python_callable=load )

extract_task >> transform_task >> load_task

[/code]

这个 DAG 定义了一个三步的 ETL 流程:extract -> transform -> load。» 符号表示依赖关系,extract 完成后才能执行 transform,transform 完成后才能执行 load。

Operator

Operator 是 Airflow 中执行具体工作的单元。Airflow 提供了丰富的内置 Operator:

Operator 用途 适用场景
PythonOperator 执行 Python 函数 自定义逻辑、API 调用
PostgresOperator 执行 SQL 语句 数据库操作
BashOperator 执行 Shell 命令 文件操作、系统命令
S3Operator 操作 S3 文件 云存储操作
SimpleHttpOperator 发送 HTTP 请求 API 调用
EmailOperator 发送邮件 告警通知
SlackWebhookOperator 发送 Slack 消息 团队通知

Sensor

Sensor 是一种特殊的 Operator,它等待某个条件满足后再继续执行。这在数据管道中非常有用。

from airflow.sensors.sql import SQLSensor

# 等待 staging 表中有数据再执行下一步
wait_for_data = SQLSensor(
    task_id="wait_for_staging_data",
    conn_id="postgres_default",
    sql="SELECT count(*) FROM staging.orders WHERE load_date = CURRENT_DATE",
    success=lambda res: res[0][0] > 0,
    timeout=600,  # 最多等 10 分钟
    poke_interval=30  # 每 30 秒检查一次
)
[/code]

## 安装与配置

使用 pip 安装

pip install apache-airflow

安装 PostgreSQL 和 Redis 相关依赖(生产环境推荐)

pip install apache-airflow[postgres,redis,celery]

初始化数据库

airflow db init

创建管理员用户

airflow users create
–username admin
–password admin
–firstname Admin
–lastname User
–role Admin
–email admin@example.com

启动 Web 服务

airflow webserver –port 8080

启动 Scheduler

airflow scheduler [/code]

生产环境推荐使用 Docker Compose 部署:

# docker-compose.yaml
version: '3.8'
services:
  postgres:
    image: postgres:15
    environment:
      POSTGRES_USER: airflow
      POSTGRES_PASSWORD: airflow
      POSTGRES_DB: airflow
    volumes:
      - postgres_data:/var/lib/postgresql/data

  redis:
    image: redis:7

  airflow-webserver:
    image: apache/airflow:2.10.0
    command: webserver
    ports:
      - "8080:8080"
    environment:
      AIRFLOW__CORE__EXECUTOR: CeleryExecutor
      AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
      AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
      AIRFLOW__CELERY__BROKER_URL: redis://redis:6379/0
    volumes:
      - ./dags:/opt/airflow/dags
[/code]

## 编写 ETL DAG

下面是一个完整的 ETL DAG,结合了 PythonOperator、PostgresOperator 和 BashOperator:

from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.postgres.operators.postgres import PostgresOperator from airflow.operators.bash import BashOperator from airflow.operators.email import EmailOperator from datetime import datetime, timedelta import pandas as pd import requests

default_args = { “owner”: “data_team”, “depends_on_past”: False, “email_on_failure”: True, “email”: [“alert@company.com”], “retries”: 2, “retry_delay”: timedelta(minutes=5), “execution_timeout”: timedelta(hours=2), }

def extract_orders_from_api(**context): “““从 API 抽取订单数据””” api_url = “https://api.example.com/orders" params = { “start_date”: context[“ds”], # execution date “end_date”: context[“next_ds”] } response = requests.get(api_url, params=params, timeout=30) response.raise_for_status() data = response.json()

df = pd.DataFrame(data)
output_path = f"/tmp/orders_{context['ds']}.parquet"
df.to_parquet(output_path, index=False)

# 通过 XCom 传递文件路径
context["task_instance"].xcom_push(key="output_path", value=output_path)
return len(df)

def validate_data(**context): “““验证数据质量””” ti = context[“task_instance”] output_path = ti.xcom_pull( task_ids=“extract_orders”, key=“output_path” ) df = pd.read_parquet(output_path)

# 数据质量检查
checks = {
    "总行数": len(df),
    "缺失订单ID": df["order_id"].isna().sum(),
    "负金额": (df["amount"] < 0).sum(),
    "重复订单": df["order_id"].duplicated().sum()
}

for check, value in checks.items():
    print(f"{check}: {value}")

# 如果缺失值超过阈值,直接报错
if checks["缺失订单ID"] > 0:
    raise ValueError("发现缺失订单ID的数据")

return checks

with DAG( dag_id=“etl_daily_orders”, default_args=default_args, start_date=datetime(2026, 1, 1), schedule=“0 3 * * *”, # 每天凌晨 3 点 catchup=True, tags=[“etl”, “orders”], description=“每日订单 ETL 管道” ) as dag:

# 任务 1:从 API 抽取数据
extract_orders = PythonOperator(
    task_id="extract_orders",
    python_callable=extract_orders_from_api,
)

# 任务 2:数据验证
validate = PythonOperator(
    task_id="validate_data",
    python_callable=validate_data,
)

# 任务 3:创建 staging 表
create_staging = PostgresOperator(
    task_id="create_staging_table",
    postgres_conn_id="warehouse",
    sql="""
    CREATE TABLE IF NOT EXISTS staging.orders (
        order_id VARCHAR(50) PRIMARY KEY,
        customer_id VARCHAR(50),
        amount NUMERIC(12,2),
        status VARCHAR(20),
        order_date DATE,
        load_date DATE DEFAULT CURRENT_DATE
    );
    """,
)

# 任务 4:加载数据到 staging
load_to_staging = BashOperator(
    task_id="load_to_staging",
    bash_command="""
    psql  \
      -c "\\copy staging.orders FROM '/tmp/orders_{{ ds }}.parquet'"
    """,
)

# 任务 5:清洗和转换
transform = PostgresOperator(
    task_id="transform_data",
    postgres_conn_id="warehouse",
    sql="""
    INSERT INTO dw.daily_order_summary
    SELECT
        order_date,
        COUNT(DISTINCT customer_id) AS customer_count,
        COUNT(*) AS order_count,
        SUM(amount) AS total_revenue
    FROM staging.orders
    WHERE load_date = CURRENT_DATE
    GROUP BY order_date
    ON CONFLICT (order_date) DO UPDATE
    SET order_count = EXCLUDED.order_count,
        total_revenue = EXCLUDED.total_revenue;
    """,
)

# 任务 6:发送成功通知
send_notification = EmailOperator(
    task_id="send_success_email",
    to=["data_team@company.com"],
    subject="ETL 完成: {{ ds }}",
    html_content="<h3>每日订单 ETL 已成功完成</h3>",
)

# 设置依赖关系
extract_orders >> validate >> create_staging >> load_to_staging >> transform >> send_notification

[/code]

任务依赖与分支

实际场景中,任务之间的依赖关系往往更加复杂。Airflow 支持多种依赖模式:

# 并行执行
[extract_a, extract_b, extract_c] >> transform

# 条件分支
from airflow.operators.python import BranchPythonOperator

def decide_branch(**context):
    record_count = context["ti"].xcom_pull(task_ids="extract", key="count")
    if record_count > 10000:
        return "full_transform"
    return "quick_transform"

branch = BranchPythonOperator(
    task_id="branch_decider",
    python_callable=decide_branch,
)
branch >> [full_transform, quick_transform]

# 汇聚
[full_transform, quick_transform] >> load
[/code]

## 调度策略

Airflow 提供了灵活的调度方式:

from airflow import DAG from datetime import timedelta

Cron 表达式

dag = DAG(dag_id=“cron_dag”, schedule=“0 2 * * *”)

预设间隔

dag = DAG(dag_id=“preset_dag”, schedule="@daily”) # 每天 dag = DAG(dag_id=“hourly_dag”, schedule="@hourly") # 每小时 dag = DAG(dag_id=“weekly_dag”, schedule="@weekly") # 每周

timedelta

dag = DAG(dag_id=“interval_dag”, schedule=timedelta(hours=6))

数据驱动调度(等待上游数据就绪)

from airflow.sensors.external_task import ExternalTaskSensor

wait_for_dbt = ExternalTaskSensor( task_id=“wait_for_dbt_run”, external_dag_id=“dbt_transformations”, external_task_id=“dbt_run_complete”, timeout=3600, ) [/code]

失败处理与告警

生产环境中的数据管道难免会失败。Airflow 提供了多层级的失败处理机制:

重试机制

from airflow import DAG
from datetime import timedelta

default_args = {
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,  # 指数退避
    "max_retry_delay": timedelta(hours=1),
}
[/code]

### SLA 监控

from airflow import DAG from datetime import timedelta

dag = DAG( dag_id=“sla_monitored_dag”, sla=timedelta(hours=4), # 超过 4 小时没完成就告警 default_args={ “sla_miss_callback”: send_sla_alert, } ) [/code]

自定义告警

from airflow.operators.slack import SlackWebhookOperator

def task_failure_alert(context):
    """任务失败时发送 Slack 告警"""
    slack_msg = f"""
        :x: 任务失败
        DAG: {context['dag'].dag_id}
        Task: {context['task'].task_id}
        Execution: {context['ds']}
        Log: {context['task_instance'].log_url}
    """
    SlackWebhookOperator(
        task_id="slack_alert",
        webhook_token="T00000000/B00000000/xxxxxxxxxx",
        message=slack_msg,
    ).execute(context=context)

default_args = {
    "on_failure_callback": task_failure_alert,
}
[/code]

## XCom:任务间数据传递

Airflow 的任务是隔离执行的,任务间不能直接共享变量。XCom(Cross-Communication)是 Airflow 提供的任务间数据传递机制。

将数据推送到 XCom

def push_function(**context): context[“ti”].xcom_push(key=“order_count”, value=1000)

从 XCom 拉取数据

def pull_function(**context): count = context[“ti”].xcom_pull( task_ids=“push_task”, key=“order_count” ) print(f"获取到订单数: {count}") [/code]

XCom 适合传递小数据(配置参数、文件路径、记录数等)。大数据应该写入文件或数据库,然后在 XCom 中传递路径。

生产环境最佳实践

1. 使用 Celery Executor

# airflow.cfg
[core]
executor = CeleryExecutor

[celery]
worker_concurrency = 8
worker_prefetch_multiplier = 1
worker_umask = 0007
[/code]

### 2. 连接管理

敏感信息通过 Airflow 的 Connection 管理,不要写在代码里:

在代码中使用 Connection ID

from airflow.providers.postgres.operators.postgres import PostgresOperator

task = PostgresOperator( task_id=“query”, postgres_conn_id=“warehouse_prod”, # 引用 Connection sql=“SELECT 1” ) [/code]

Connection 的配置通过 Web UI 或环境变量设置,不会出现在代码仓库中。

3.变量管理

from airflow.models import Variable

# 设置变量(通过 Web UI 或 CLI)
config = Variable.get("etl_config", deserialize_json=True)

# 在代码中使用
batch_size = config.get("batch_size", 10000)
retry_times = config.get("retry_times", 3)
[/code]

### 4. DAG 设计原则

- **每个 DAG 只做一个业务域的事情**,不要把所有逻辑塞进一个 DAG
- **任务粒度要适中**,太粗难以重试,太细增加调度开销
- **幂等设计**,同一个 DAG 在相同时间点运行多次应该得到相同结果
- **参数化**,使用 Airflow 的宏和变量,不要硬编码日期和路径

## 小结

Airflow 是现代数据工程领域最核心的工作流编排工具。本文介绍了 DAG、Operator、Sensor 等核心概念,展示了如何编写完整的 ETL DAG,并讨论了任务依赖、调度策略、失败处理和最佳实践。有了 Airflow,你的数据管道就能从手工操作升级为自动化、可监控的生产级系统。下一章我们将用前面学到的所有知识,构建一个完整的实战项目。

*Summary:* Airflow DAG 概念、ETL 编排实战与生产环境配置。