Airflow 核心概念

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

认识您的讲师

  Volker Janz 个人照片

$$

Volker Janz

$$

  • Astronomer 高级开发者布道师
  • 在游戏领域担任数据工程师 14+ 年
  • 1.x 版本 起使用 Airflow
  • Data Engineer Things 的演讲者、导师与通讯负责人
使用 Airflow 构建数据流水线

您将构建什么

 

  • 使用 TaskFlow API 编写 Dags
  • 通过任务映射与基于资产的调度构建动态工作流
  • 重试与回调处理失败
  • 通过 Airflow 运行 SQL 工作负载

各章节内容可视化

使用 Airflow 构建数据流水线

开始之前

 

$$

  • 熟悉 Dag、任务与算子
  • 理解 调度基础

Airflow 入门——课程页面

使用 Airflow 构建数据流水线

快速回顾

from airflow.sdk import dag, task

@dag
def star_wars_dag():


@task def get_star_wars_person(): import requests return requests.get("https://swapi.dev/api/people/1/").json()
@task.bash def print_name(person): return f"echo '{person['name']}'"
person = get_star_wars_person() print_name(person) star_wars_dag()
  • Dag:带依赖的任务集合
  • 任务(Tasks):独立的工作单元
  • 算子 / 装饰器:定义每个任务做什么
  • 依赖:定义执行顺序

$$

简单 Dag

使用 Airflow 构建数据流水线

Airflow 架构

$$

Airflow 3 架构

 

$$

  • 编排:Scheduler,Dag Processor
  • 执行:Worker,Triggerer
  • 接口与存储:API Server,Metadata DB
使用 Airflow 构建数据流水线

调度方式

# No automatic runs: trigger manually (default)
@dag(schedule=None)
def my_pipeline(): ...

# Time-based: runs every day at 6 AM @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Data-aware: runs when an Asset updates @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
使用 Airflow 构建数据流水线

两种编写 Dag 的方式

经典算子

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • 当无可用装饰器时使用

TaskFlow API

@task
def extract():
    return {"users": 150}

data = extract()
transform(data)
  • 简洁的 Python 装饰器
  • 更少样板代码
  • 可与经典算子组合
使用 Airflow 构建数据流水线

本课程的练习

$$

IDE 练习截图

 

  • IDE 练习:编辑真实的 .py 文件
  • 点击 "Run this file" 或使用 python3 filename.py
  • dag.test() 在单进程中运行整个 Dag
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...