이름 있는 매핑과 Classic Operator 매핑 (map_index_template, expand_kwargs)
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/dynamic-task-mapping.rst - Named mapping, Mapping with non-TaskFlow operators, Mapping over result of classic operators, Mixing TaskFlow and classic operators, Assigning multiple parameters to a non-TaskFlow operator (약 183-345줄)
이 모듈을 다 읽으면
- map_index_template으로 매핑된 태스크 인스턴스에 의미 있는 이름을 붙이는 방법을 설명할 수 있다
- Classic Operator에 partial/expand를 적용할 때의 제약을 설명할 수 있다
- expand_kwargs와 expand의 차이를 설명할 수 있다
기본적으로 매핑된 태스크는 정수 인덱스로 구분되지만 map_index_template으로 Jinja 템플릿 기반 이름을 부여할 수 있고, Classic(비-TaskFlow) Operator에도 partial/expand를 적용할 수 있으며 여러 인자를 함께 매핑하려면 expand_kwargs를 쓴다.
map_index_template으로 이름 있는 매핑
기본적으로 매핑된 태스크에는 정수 인덱스가 부여된다. Airflow UI에서 각 매핑 태스크의 이 정수 인덱스를, 태스크의 입력값 기반 이름으로 오버라이드할 수 있는데, 태스크에 `map_index_template`으로 Jinja 템플릿을 주면 된다. `.expand(<property>=...)` 형태의 확장이라면 보통 `map_index_template="{{ task.<property> }}"` 형태를 쓴다. 이 템플릿은 각 확장된 태스크가 실행을 마친 뒤 태스크 컨텍스트를 이용해 렌더링된다.
템플릿이 메인 실행 블록 이후에 렌더링된다는 점을 이용해, 렌더링 컨텍스트에 동적으로 값을 주입할 수도 있다. Jinja 템플릿 문법만으로 원하는 이름을 표현하기 어려울 때, 특히 TaskFlow 함수 안에서 유용하다 - `get_current_context()`로 컨텍스트를 가져와 `context["my_variable"] = my_value * 3`처럼 값을 직접 넣고 `@task(map_index_template="{{ my_variable }}")`로 그 값을 참조하면 된다.
핵심 포인트
- map_index_template은 확장된 각 태스크 인스턴스가 실행을 마친 뒤 태스크 컨텍스트를 이용해 렌더링되며, UI에 정수 인덱스 대신 의미 있는 이름을 보여준다
- 템플릿 로직을 Jinja만으로 표현하기 어려울 때는 get_current_context()로 컨텍스트에 값을 직접 주입해(context[...] = ...) 렌더링에 활용할 수 있다
Classic Operator에 partial/expand 적용하기
TaskFlow 스타일이 아닌 Classic 오퍼레이터에도 `partial`과 `expand`를 적용할 수 있다. 다만 `task_id`, `queue`, `pool`을 비롯한 `BaseOperator`의 대부분의 공통 인자는 매핑할 수 없어, 이런 값들은 `partial()`에 넘겨야 한다. `partial()`에도 키워드 인자만 넘길 수 있다는 점은 `expand()`와 같다.
Classic 오퍼레이터의 실행 결과로 매핑하고 싶다면, 오퍼레이터 자체가 아니라 명시적으로 그 *출력*(`.output`)을 참조해야 한다. 예를 들어 `extract`가 만든 데이터를 변환하는 `transform`, 그 결과를 적재하는 `load`를 이어 붙일 때 `TransformOperator.partial(task_id="transform").expand(input=extract.output)`처럼 `.output`을 명시적으로 참조한다.
핵심 포인트
- Classic Operator를 매핑할 때 task_id, queue, pool 등 대부분의 BaseOperator 공통 인자는 매핑할 수 없어 partial()에 넘겨야 한다
- Classic Operator의 실행 결과로 매핑하려면 오퍼레이터 자체가 아니라 명시적으로 .output을 참조해야 한다 (예: TransformOperator.partial(...).expand(input=extract.output))
TaskFlow와 Classic 혼합, expand_kwargs로 여러 인자 함께 매핑
TaskFlow 태스크와 Classic 오퍼레이터를 섞어 쓰는 것도 가능하다 - 예를 들어 S3 버킷에 정기적으로 도착하는 파일들을, 몇 개가 오든 모두 같은 처리를 하고 싶을 때, Classic `S3ListOperator`로 파일 목록을 얻고 그 `.output`을 TaskFlow 태스크의 `.partial().expand()`로 넘기면 된다.
한 다운스트림 오퍼레이터에 업스트림이 여러 인자를 함께 지정해야 하는 경우에는 `expand_kwargs` 함수를 쓴다 - 이 함수는 서로 매핑할 딕셔너리들의 시퀀스를 받는다. 예를 들어 `BashOperator.partial(task_id="bash").expand_kwargs([{...}, {...}])`처럼 완전한 kwargs 딕셔너리 두 개를 넘기면 실행 시점에 각각 다른 조합의 인자로 실행되는 두 개의 태스크 인스턴스가 생긴다. `expand_kwargs`는 `PythonOperator`의 `op_kwargs`처럼 대부분의 오퍼레이터 인자와도 함께 섞어 쓸 수 있다.
`expand`와 마찬가지로, 딕셔너리 리스트를 반환하는 XCom이나 딕셔너리를 반환하는 XCom들의 리스트에 대해서도 매핑할 수 있다. 이를 이용하면, 예를 들어 파일 확장자에 따라 서로 다른 버킷으로 복사하는 것처럼 아이템별 분기(branching) 로직을 만들 수 있다 - 업스트림 태스크가 파일마다 목적지 버킷을 담은 kwargs 딕셔너리를 XCom 으로 반환하면, 그 딕셔너리 리스트를 그대로 `expand_kwargs`로 넘겨 처리한다.
핵심 포인트
- expand_kwargs는 매핑 리스트(각 원소가 완전한 kwargs 딕셔너리)를 받아 원소 개수만큼 태스크 인스턴스를 만든다는 점에서, 여러 파라미터의 모든 조합을 만드는 expand의 cross product와 다르다
- 업스트림 태스크가 XCom으로 딕셔너리 리스트를 반환하게 하면(create_copy_kwargs 패턴), 파일 확장자별로 다른 버킷에 복사하는 식의 아이템별 분기 로직을 expand_kwargs로 구현할 수 있다