← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 60번째

Airflow 모듈 60/151 airflow-learn-60

태스크 그룹 매핑과 depth-first 실행

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/dynamic-task-mapping.rst - Mapping over a task group, Value references in a task group function, Depth-first execution, Depending on a mapped task group's output, Branching on a mapped task group's output (약 346-478줄)

이 모듈을 다 읽으면

  • @task_group에 expand/expand_kwargs를 적용했을 때의 동작을 설명할 수 있다
  • depth-first 실행이 breadth-first(그룹 없는 개별 expand)와 어떻게 다른지 설명할 수 있다
  • 태스크 그룹 함수 안에서 매핑 값이 왜 즉시 분기 조건으로 쓰일 수 없는지 설명할 수 있다

@task_group 함수에도 expand/expand_kwargs를 적용해 매핑된 태스크 그룹을 만들 수 있으며, 이때 그룹 내 모든 태스크가 함께 확장되어 같은 확장 인스턴스 안의 태스크끼리만 의존하는 depth-first 실행이 이뤄지지만, 그룹 함수 코드 자체는 워커 없이 실행되므로 매핑 인자의 실제 값을 즉시 참조할 수 없다.

태스크 그룹 매핑과 값 참조 제약

일반 TaskFlow 태스크와 마찬가지로, `@task_group`으로 데코레이트된 함수에도 `expand`나 `expand_kwargs`를 호출해 매핑된 태스크 그룹을 만들 수 있다. 예를 들어 `file_transforms.expand(filename=["data1.json", "data2.json"])`처럼 쓰면, 그룹 안의 `convert_to_yaml` 태스크가 두 개의 인스턴스로 확장되며 각각 "data1.json"과 "data2.json"을 입력으로 받는다.

TaskFlow 태스크 함수(`@task`)와 태스크 *그룹* 함수(`@task_group`)의 중요한 차이는, 태스크 그룹에는 대응하는 워커가 없다는 점이다. 그래서 태스크 그룹 함수 안의 코드는 넘겨받은 매핑 인자를 실제 값으로 해석(resolve)할 수 없고, 그 인자는 실제 태스크에 전달되어야만 비로소 해석되는 참조일 뿐이다. 예를 들어 태스크 그룹 함수 본문에서 `if not value:` 같은 조건 분기를 쓰면 기대한 대로 동작하지 않는다 - `value`가 아직 참조일 뿐이기 때문이다. 따라서 매핑된 값에 어떤 로직(조건 분기나 반복)을 적용하려면, 반드시 그 로직을 실제 태스크(조건이면 `@task.branch`/`BranchPythonOperator`, 반복이면 태스크 매핑 자체)로 옮겨서 값이 해석되게 해야 한다.

핵심 포인트

  • @task_group 함수에도 expand/expand_kwargs를 적용해 매핑된 태스크 그룹을 만들 수 있다
  • 태스크 그룹 함수에는 대응하는 워커가 없어서, 함수 본문 안의 매핑 인자는 실제 값이 아니라 참조일 뿐이다 - if not value: 같은 조건 분기는 기대한 대로 동작하지 않으며, 반드시 실제 태스크(예: @task.branch)에 값을 넘겨야 값이 해석(resolve)된다

Depth-first 실행과 breadth-first의 차이

매핑된 태스크 그룹 안에 여러 태스크가 있으면, 그룹 안의 모든 태스크가 같은 입력에 대해 "함께" 확장된다. 예를 들어 그룹 안에 `convert_to_yaml`과 `replace_defaults` 두 태스크가 있고 그룹이 두 개로 확장되면, 두 태스크 모두 각각 두 개의 인스턴스가 된다.

비슷한 효과는 그룹 없이 두 태스크를 각각 개별적으로 `expand`해도 낼 수 있는데(이를 그룹 없는 방식, breadth-first라 부른다), 태스크 그룹의 depth-first 실행과의 차이는 그룹 안의 각 태스크가 오직 자신의 "관련 있는 입력"에만 의존한다는 점이다. 위 예에서 첫 번째 `replace_defaults`는 같은 그룹 인스턴스의 `convert_to_yaml`에만 의존할 뿐, 다른 그룹 인스턴스(예: 다른 파일)의 `convert_to_yaml`에는 의존하지 않는다 - 즉 두 번째 `convert_to_yaml("data2.json")`이 끝나기 전에도 첫 번째 `replace_defaults`가 실행될 수 있고, 그 성공 여부를 신경 쓸 필요가 없다. 이 depth-first 전략은 더 논리적인 태스크 분리, 세밀한 의존성 규칙, 정확한 자원 배분을 가능하게 한다.

다만 매핑된 태스크 그룹 안에 중첩된 태스크 매핑은 현재 허용되지 않는다 - 기술적으로 크게 어렵지는 않지만, UI 복잡도를 상당히 높이고 일반적인 사용 사례에 꼭 필요하지 않다고 판단해 의도적으로 제외한 기능이며, 사용자 피드백에 따라 향후 재검토될 수 있다.

핵심 포인트

  • 매핑된 태스크 그룹 안에 여러 태스크가 있으면 그룹의 모든 태스크가 같은 입력에 대해 함께 확장되며, 이를 depth-first 실행이라 부른다 - 그룹 없이 각 태스크를 개별적으로 expand하는 breadth-first 방식과 대비된다
  • depth-first 실행에서는 각 확장 인스턴스의 태스크가 같은 그룹 인스턴스 안의 선행 태스크에만 의존해, 다른 그룹 인스턴스(예: 다른 파일)의 처리 성공 여부를 신경 쓸 필요 없이 더 세밀한 의존성과 자원 배분이 가능하다
  • 매핑된 태스크 그룹 안에 중첩된 태스크 매핑은 UI 복잡도 문제로 현재 의도적으로 지원하지 않는다

매핑된 태스크 그룹의 출력에 의존/분기하기

매핑된 태스크 그룹의 출력에 의존하는 것도, 매핑된 일반 태스크의 출력에 의존하는 것과 비슷하게 결과가 자동 집계된다. 예를 들어 그룹 안에서 `add_one`을 거쳐 `double`을 적용하는 `add_to`를 `add_to.expand(value=[1, 2, 3])`로 확장하면, 그 결과에 의존하는 `consumer`는 `[4, 6, 8]`을 받는다 - 일반 매핑 태스크의 결과에 대해 할 수 있는 모든 연산을 여기서도 똑같이 할 수 있다.

다만 매핑된 태스크 그룹의 *출력*을 기준으로 `@task.branch` 같은 분기 로직을 구현하는 것은 불가능하다. 대신 태스크 그룹의 *입력*을 기준으로 분기하는 것은 가능하다 - 그룹 함수 안에서 입력 값에 따라 `"my_task_group.a"`, `"my_task_group.b"`, `"my_task_group.c"` 중 하나를 반환하는 `@task.branch` 태스크를 두어, 매핑된 입력 각각에 대해 세 태스크 중 하나를 실행하도록 만들 수 있다.

핵심 포인트

  • 매핑된 태스크 그룹의 출력에 의존하면 그룹 내부 처리(예: add_one 후 double)를 거친 결과가 자동으로 집계되어 컨슈머에 전달된다
  • 매핑된 태스크 그룹의 출력을 기준으로 분기(@task.branch)하는 것은 불가능하지만, 그룹의 입력을 기준으로 그룹 내부에서 분기하는 것은 가능하다