Apache Airflow:从零开始理解数据流水线
很多人把 Airflow 想得太复杂,其实把它看作一个“带依赖关系的智能待办清单”就足够了。它最核心的能力就是让你用 Python 代码定义工作流,然后由它来负责定时触发、按顺序执行,以及在某个环节崩溃时自动重试。
给大家看一个最典型的 ETL(提取-转换-加载)实操逻辑,用代码定义一个简单的 DAG 结构大概是这样的:
下一篇
用 Rust 写 Polars 插件 →
在 Airflow 里,最关键的几个概念得理清楚,否则看文档会非常痛苦:
- DAG (有向无环图): 简单说就是你的整个工作流地图。它规定了任务 A 必须在任务 B 之前完成,且绝对不能形成环路(否则程序会死循环)。
- Operator (算子): 相当于任务模板。你不需要从零写怎么调用 Python 或 SQL,直接用现成的
PythonOperator或BashOperator就像搭积木一样把任务拼起来。 - Scheduler (调度器): 那个盯着时钟的“监工”,到了预定时间就踢启动开关。
- XCom: 任务之间的“传话筒”,用来传递少量的中间数据。
给大家看一个最典型的 ETL(提取-转换-加载)实操逻辑,用代码定义一个简单的 DAG 结构大概是这样的:
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
# 定义一个简单的 DAG
with DAG(
dag_id='my_first_etl_pipeline',
start_date=datetime(2023, 1, 1),
schedule_interval='@daily',
catchup=False
) as dag:
def extract():
print("正在从 API 抓取原始数据...")
return "raw_data"
def transform(ti):
# 使用 XCom 获取上一个任务的返回值
data = ti.xcom_pull(task_ids='extract_task')
print(f"正在清洗数据: {data}")
def load():
print("数据已成功写入仓库")
# 定义任务
extract_task = PythonOperator(task_id='extract_task', python_callable=extract)
transform_task = PythonOperator(task_id='transform_task', python_callable=transform)
load_task = PythonOperator(task_id='load_task', python_callable=load)
# 设置依赖关系:提取 -> 转换 -> 加载
extract_task >> transform_task >> load_task在实际部署和实操中,最容易踩的坑就是 start_date 和 catchup 的配置。如果你把 start_date 设在一年以前且 catchup=True,启动的一瞬间 Airflow 会尝试把过去一年所有漏掉的任务全部补回来,直接把服务器 CPU 跑满。建议新手在测试阶段统一把 catchup 设为 False。