Dynamic Task Mapping

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

Volker Janz

Senior Developer Advocate at Astronomer

ปัญหาของการ hardcode task

@task
def process_file_1():
    transform("/data/file_1.csv")
@task
def process_file_2():
    transform("/data/file_2.csv")
@task
def process_file_3():
    transform("/data/file_3.csv")

ตัวอย่าง Dynamic Task Mapping

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

.expand()

@task
def fetch_files():
    return ["/data/file_1.csv", "/data/file_2.csv", "/data/file_3.csv"]

@task(max_active_tis_per_dagrun=2) # control max parallelism def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

ปรับขนาดขณะ runtime

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

Mapped task ใน UI

 

  • Mapped instance จะแสดงเป็น indexed instance ภายใต้ task เดียว
  • คลิกแต่ละรายการเพื่อดู log และสถานะแยกกัน
  • DAG ยัง อ่านง่าย แม้มี instance หลายร้อยรายการ

Mapped task ใน UI

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

.partial()

@task
def process_file(path: str, output_dir: str):
    transform(path, output_dir)

files = get_files() process_file.partial(output_dir="/out").expand(path=files)

Partial และ expand

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

ปัญหา cross-product

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Cross product: 3 x 3 = 9 instances (not what we want!) process.expand(path=files, dest=destinations)

Zip list

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

.zip()

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# Paired: 3 instances (a->a, b->b, c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() จับคู่รายการ ทีละตำแหน่ง
  • expand_kwargs แตก pair แต่ละคู่เป็น keyword argument
การสร้าง Data Pipeline ด้วย Airflow

เลือกใช้แบบไหนดี

 

Pattern กรณีใช้งาน
expand() วนซ้ำ list เดียว
partial() กำหนด parameter ร่วม ให้ทุก instance
zip() จับคู่ หลาย list ตามตำแหน่ง

 

  • ช่วยสร้าง DAG ที่ขยายได้
  • เมื่อข้อมูลเปลี่ยน จำนวน task จะปรับ อัตโนมัติขณะ runtime
การสร้าง Data Pipeline ด้วย Airflow

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

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

Preparing Video For Download...