← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 61번째

Airflow 모듈 61/151 airflow-learn-61

매핑 데이터 필터링·변환·결합 (filter, map(), zip, concat)

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/dynamic-task-mapping.rst - Filtering items from a mapped task, Transforming expanding data, Combining upstream data (aka zipping), Concatenating multiple upstreams (약 479-601줄)

이 모듈을 다 읽으면

  • 매핑 대상에서 특정 항목을 제외하는 두 가지 방식(None 반환 vs AirflowSkipException)의 차이를 설명할 수 있다
  • .map()이 일반 태스크와 어떻게 다르게 실행되는지 설명할 수 있다
  • .zip()과 .concat()의 각각의 용도와 길이 처리 방식을 설명할 수 있다

매핑 태스크는 None을 반환해 특정 항목을 다운스트림에서 걸러낼 수 있고, .map()으로 XCom 반환값을 태스크로 만들지 않고도 사전 변환할 수 있으며, .zip()과 .concat()으로 여러 업스트림 iterable을 각각 튜플 결합 또는 순차 연결해 하나의 매핑 대상으로 합칠 수 있다.

None 반환으로 매핑 항목 걸러내기

실제 태스크로 실행되는 매핑 태스크는 특정 입력에 대해 `None`을 반환함으로써 그 항목을 다운스트림으로 넘기지 않고 걸러낼 수 있다. 예를 들어 S3 버킷에서 특정 확장자를 가진 파일만 다른 버킷으로 복사하고 싶다면, `create_copy_kwargs` 태스크가 `.json`이나 `.yml`로 끝나지 않는 파일에 대해 `None`을 반환하도록 구현하면 된다. 이렇게 하면 `copy_files`는 `.json`/`.yml` 파일에 대해서만 확장되고 나머지는 무시된다.

핵심 포인트

  • 매핑 태스크(실제 태스크로 실행되는 것)가 특정 입력에 대해 None을 반환하면, 그 항목은 다운스트림 매핑 대상에서 제외된다

map()으로 사전 변환하기

매핑을 위한 출력 데이터 형식을 변환하고 싶은 경우가 흔한데, 특히 출력 형식이 미리 정해져 있어 쉽게 바꿀 수 없는 비-TaskFlow 오퍼레이터에서 그렇다. 이럴 때 `.map()` 함수로 이런 변환을 간단히 수행할 수 있다.

`.map()`에 넘기는 콜러블(위 예의 `create_copy_kwargs`)은 태스크가 아니라 순수한 파이썬 함수여야 한다 - 이 변환은 다운스트림 태스크(예: `copy_files`)의 "전처리"의 일부로 취급되며, Dag 상의 독립된 태스크가 아니다. 이 콜러블은 항상 정확히 하나의 위치 인자를 받으며, 파이썬 내장 `map()`처럼 태스크 매핑에 쓰이는 iterable의 각 항목마다 한 번씩 호출된다.

이 콜러블은 다운스트림 태스크의 일부로 실행되므로 기존에 태스크 함수를 작성하던 모든 기법을 그대로 쓸 수 있다. 다만 항목을 스킵 처리하고 싶다면 `AirflowSkipException`을 raise해야 한다는 점에 유의해야 한다 - 매핑 태스크 자체의 필터링과 달리, `.map()` 안에서는 `None`을 반환해도 스킵 처리가 되지 않는다.

핵심 포인트

  • .map()에 넘기는 콜러블은 태스크가 아니라 순수 파이썬 함수여야 하며, Dag 상의 독립된 태스크가 아니라 다운스트림 태스크의 전처리 일부로 실행된다
  • .map() 안에서 항목을 건너뛰려면 None을 반환하는 것은 동작하지 않으며 반드시 AirflowSkipException을 raise해야 한다 - 이는 매핑 태스크 자체의 None 필터링과 다른 규칙이다

zip()과 concat()으로 업스트림 결합하기

여러 입력 소스를 하나의 태스크 매핑 iterable로 결합하고 싶은 경우도 흔하다. 파이썬 내장 `zip()`처럼 여러 iterable을 튜플로 묶는 것을 흔히 "zipping"이라 부르며, 이 역시 다운스트림 태스크의 전처리로 수행된다. 이는 태스크 매핑에서 조건부 로직을 쓸 때 특히 유용한데, 예를 들어 S3에서 파일을 내려받으면서 이름을 새로 짓고 싶다면 파일 목록과 새 이름 목록을 `.zip()`으로 묶어 각 다운로드 태스크가 `(원본파일명, 새이름)` 튜플을 받게 만들 수 있다.

파이썬 내장 `zip`처럼, 임의 개수의 iterable을 함께 zip해 위치 인자 개수만큼의 튜플들의 iterable을 얻을 수 있다. 기본적으로 zip된 iterable의 길이는 zip 대상 중 가장 짧은 것과 같아지고, 남는 항목은 버려진다. 선택적 키워드 인자 `default`를 넘기면 파이썬의 `itertools.zip_longest`처럼 동작이 바뀌어, zip된 iterable의 길이가 가장 *긴* 것과 같아지고 부족한 자리는 `default`로 채워진다.

또 다른 흔한 결합 패턴은 여러 iterable에 대해 같은 태스크를 실행하는 것이다 - 물론 각 iterable마다 별도로 같은 코드를 실행해도 되지만(overrider된 `task_id`로 각각 expand), `concat`을 쓰면 Dag가 더 확장성 있고 살펴보기도 쉬워진다. `concat`은 여러 iterable을 순서대로 이어 붙여 하나의 매핑 대상으로 만들며, 임의 개수의 iterable을 함께 concat할 수 있고(`foo.concat(bar, rex)`), 반환값도 XCom 참조이므로 체이닝도 가능하다 (`foo.concat(bar).concat(rex)`) - 파이썬의 `itertools.chain`과 비슷하게 순서를 유지하며 모두 이어 붙인 하나의 iterable이 된다.

핵심 포인트

  • .zip()은 여러 업스트림 iterable을 튜플로 묶어 매핑 대상을 만들며, 기본적으로 가장 짧은 iterable 길이에 맞추고 초과 항목은 버려지지만 default 키워드를 주면 itertools.zip_longest처럼 가장 긴 길이에 맞춰 부족한 자리를 채운다
  • .concat()은 여러 iterable을 순서대로 이어 붙여 하나의 매핑 대상으로 만들어, 업스트림마다 태스크를 따로 expand하지 않고 하나의 태스크로 통합할 수 있게 한다