Các khái niệm cốt lõi của Airflow

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Gặp gỡ giảng viên của bạn

  Ảnh hồ sơ Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ năm làm data engineer trong gaming
  • Làm việc với Airflow từ phiên bản 1.x
  • Diễn giả, mentor, và phụ trách newsletter tại Data Engineer Things
Xây dựng Data Pipeline với Airflow

Bạn sẽ xây dựng gì

 

  • Viết Dags với TaskFlow API
  • Xây workflow động bằng task mapping và lịch dựa trên asset
  • Xử lý lỗi bằng retry và callback
  • Chạy tác vụ SQL qua Airflow

Minh họa nội dung từng chương

Xây dựng Data Pipeline với Airflow

Trước khi bắt đầu

 

$$

  • Thành thạo Dags, task và operator
  • Nắm các nguyên tắc lập lịch

Trang khóa học Introduction to Airflow

Xây dựng Data Pipeline với Airflow

Ôn nhanh

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 là tập hợp các task có quan hệ phụ thuộc
  • Task là đơn vị công việc riêng lẻ
  • Operator/decorator xác định mỗi task làm gì
  • Phụ thuộc quyết định thứ tự thực thi

$$

Dag đơn giản

Xây dựng Data Pipeline với Airflow

Kiến trúc Airflow

$$

Kiến trúc Airflow 3

 

$$

  • Điều phối: Scheduler, Dag Processor
  • Thực thi: Worker, Triggerer
  • Giao diện & Lưu trữ: API Server, Metadata DB
Xây dựng Data Pipeline với Airflow

Các cách lập lịch

# Không tự động chạy: kích hoạt thủ công (mặc định)
@dag(schedule=None)
def my_pipeline(): ...

# Theo thời gian: chạy mỗi ngày lúc 6h sáng @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Theo dữ liệu: chạy khi một Asset được cập nhật @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Xây dựng Data Pipeline với Airflow

Hai cách viết Dags

Classic operators

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Chọn khi không có decorator tương ứng

TaskFlow API

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

data = extract()
transform(data)
  • Python decorator đơn giản
  • Ít mã lặp khuôn
  • Có thể kết hợp với classic operators
Xây dựng Data Pipeline với Airflow

Bài tập trong khóa học này

$$

Ảnh chụp bài tập trong IDE

 

  • Bài tập IDE: chỉnh sửa file .py thật
  • Bấm "Run this file" hoặc dùng python3 filename.py
  • dag.test() chạy toàn bộ Dag trong một tiến trình
Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...