동적 태스크 매핑

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

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

리스트 zip

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는 각 쌍을 키워드 인수로 언패킹합니다
Airflow로 데이터 파이프라인 구축하기

언제 무엇을 사용할까

 

패턴 사용 사례
expand() 단일 리스트 매핑
partial() 인스턴스 간 공유 파라미터 고정
zip() 여러 리스트를 위치 기준으로 구성

 

  • 확장 가능한 DAG 구성에 도움이 됩니다
  • 데이터가 변경되면 태스크 수가 런타임에 자동으로 조정됩니다
Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...