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

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


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

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)

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는 각 쌍을 키워드 인수로 언패킹합니다
| 패턴 | 사용 사례 |
|---|---|
expand() |
단일 리스트 매핑 |
partial() |
인스턴스 간 공유 파라미터 고정 |
zip() |
여러 리스트를 위치 기준으로 쌍 구성 |
Airflow로 데이터 파이프라인 구축하기