แนวคิดหลักของ Airflow

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

แนะนำผู้สอน

  รูปโปรไฟล์ Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • ประสบการณ์ด้าน data engineer ในวงการเกม กว่า 14 ปี
  • ใช้งาน Airflow ตั้งแต่เวอร์ชัน 1.x
  • วิทยากร, พี่เลี้ยง, และบรรณาธิการจดหมายข่าวที่ Data Engineer Things
การสร้าง Data Pipeline ด้วย Airflow

สิ่งที่จะได้สร้าง

 

  • เขียน Dag ด้วย TaskFlow API
  • สร้าง dynamic workflow ด้วย task mapping และ asset-based scheduling
  • จัดการความล้มเหลวด้วย retries และ callbacks
  • รัน SQL workload ผ่าน Airflow

ภาพรวมเนื้อหาแต่ละบท

การสร้าง Data Pipeline ด้วย Airflow

ก่อนเริ่มต้น

 

$$

  • คุ้นเคยกับ Dag, task และ operator
  • เข้าใจพื้นฐาน การกำหนดตารางเวลา (scheduling)

หน้าคอร์ส Introduction to Airflow

การสร้าง Data Pipeline ด้วย 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 คือกลุ่มของ task ที่มี dependency
  • Task คือหน่วยงานย่อยแต่ละชิ้น
  • Operator / decorator กำหนดว่าแต่ละ task ทำอะไร
  • Dependency กำหนดลำดับการทำงาน

$$

Dag อย่างง่าย

การสร้าง Data Pipeline ด้วย Airflow

สถาปัตยกรรมของ Airflow

$$

สถาปัตยกรรม Airflow 3

 

$$

  • การประสานงาน (Orchestration): Scheduler, Dag Processor
  • การประมวลผล (Execution): Worker, Triggerer
  • อินเทอร์เฟซและการจัดเก็บ (Interface & Storage): API Server, Metadata DB
การสร้าง Data Pipeline ด้วย 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(): ...
การสร้าง Data Pipeline ด้วย Airflow

2 วิธีในการเขียน Dag

Classic operator

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • เลือกใช้เมื่อไม่มี decorator ให้ใช้งาน

TaskFlow API

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

data = extract()
transform(data)
  • ใช้ Python decorator ที่เรียบง่าย
  • โค้ดกระชับขึ้น
  • ใช้ร่วมกับ classic operator ได้
การสร้าง Data Pipeline ด้วย Airflow

แบบฝึกหัดในคอร์สนี้

$$

ภาพหน้าจอแบบฝึกหัด IDE

 

  • แบบฝึกหัด IDE: แก้ไขไฟล์ .py จริง
  • คลิก "Run this file" หรือใช้ python3 filename.py
  • dag.test() รัน Dag ทั้งหมดในกระบวนการเดียว
การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...