Apache Airflow:从零开始理解数据流水线

早八人AI炼丹师 专家 6小时前 217 浏览 5 点赞 约 1 分钟

很多人把 Airflow 想得太复杂,其实把它看作一个“带依赖关系的智能待办清单”就足够了。它最核心的能力就是让你用 Python 代码定义工作流,然后由它来负责定时触发、按顺序执行,以及在某个环节崩溃时自动重试。

在 Airflow 里,最关键的几个概念得理清楚,否则看文档会非常痛苦:

  • DAG (有向无环图): 简单说就是你的整个工作流地图。它规定了任务 A 必须在任务 B 之前完成,且绝对不能形成环路(否则程序会死循环)。
  • Operator (算子): 相当于任务模板。你不需要从零写怎么调用 Python 或 SQL,直接用现成的 PythonOperatorBashOperator 就像搭积木一样把任务拼起来。
  • 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_datecatchup 的配置。如果你把 start_date 设在一年以前且 catchup=True,启动的一瞬间 Airflow 会尝试把过去一年所有漏掉的任务全部补回来,直接把服务器 CPU 跑满。建议新手在测试阶段统一把 catchup 设为 False

AI编程AI编程实战Tutorialwebdevprogramming

全部回复 (3)

远程办公技术宅 中级 10小时前
那如果任务依赖太多,DAG界面卡死怎么优化?
0 回复
咖啡续命折腾党 中级 10小时前
建议把retries设高点,很多时候网络抖一下就报错,自动重试能救命。
0 回复
创业者阿杰 中级 10小时前
确实,刚上手时被概念绕晕了,其实就是个定时跑脚本的。
0 回复

发表回复

支持 Markdown 格式