動態任務對映

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

硬編任務的問題

@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")

動態任務對映範例

使用 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)

執行期擴充

使用 Airflow 建置資料管線

UI 中的對映任務

 

  • 對映出的執行個體以索引個體呈現在同一個任務下
  • 點進每個個體可查看個別記錄與狀態
  • 即使有數百個個體,Dag 仍易讀

UI 中的對映任務

使用 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

使用 Airflow 建置資料管線

交叉乘積的問題

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

# 笛卡兒積:3 x 3 = 9 個執行個體(不是我們要的!) process.expand(path=files, dest=destinations)

對齊配對清單

使用 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"]

# 成對:3 個執行個體(a->a、b->b、c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() 依位置一對一配對
  • expand_kwargs 將每對拆成關鍵字參數
使用 Airflow 建置資料管線

何時用哪個

 

模式 使用情境
expand() 對映單一清單
partial() 固定跨個體的共用參數
zip() 依位置配對多個清單

 

  • 有助打造可擴充的 Dags
  • 資料改變時,任務數會在執行期自動調整
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...