태스크플로우·워크플로우 (TaskFlow)

태스크플로우·워크플로우 (TaskFlow)

Airflow 2.0부터 추가됨 (Added in version 2.0)

대부분의 Dag을 Operator 대신 순수 파이썬 코드로 작성한다면, TaskFlow API를 쓰면 불필요한 보일러플레이트 없이도 훨씬 깔끔한 Dag을 만들 수 있어요. 핵심은 @task 데코레이터 하나예요.

TaskFlow는 태스크 사이의 입력과 출력을 XCom으로 옮기는 일, 그리고 의존성을 자동으로 계산하는 일까지 대신 처리해 줍니다. Dag 파일에서 TaskFlow 함수를 호출할 때 실제로 함수를 실행하는 게 아니라, 그 결과의 XCom을 나타내는 객체(XComArg)를 돌려받는다는 점이 포인트예요. 그 객체를 다운스트림 태스크나 Operator의 입력으로 쓰면 됩니다. 예시를 볼게요.

from airflow.sdk import task
from airflow.providers.smtp.operators.smtp import EmailOperator

@task
def get_ip():
    return my_ip_service.get_main_ip()

@task(multiple_outputs=True)
def compose_email(external_ip):
    return {
        'subject': f'Server connected from {external_ip}',
        'body': f'Your server executing Airflow is connected from the external IP {external_ip}<br>'
    }

email_info = compose_email(get_ip())

EmailOperator(
    task_id='send_email_notification',
    to='[email protected]',
    subject=email_info['subject'],
    html_content=email_info['body']
)

여기에는 get_ip, compose_email, send_email_notification 세 개의 태스크가 있어요.

앞의 두 개는 TaskFlow로 선언됐고, get_ip의 반환값을 자동으로 compose_email에 넘겨줍니다. XCom을 이어줄 뿐만 아니라, compose_emailget_ip다운스트림이라는 것도 자동으로 선언해 주는 거예요.

send_email_notification은 조금 더 전통적인 Operator인데, 이 역시 compose_email의 반환값을 파라미터로 쓸 수 있고, 마찬가지로 compose_email의 다운스트림이어야 한다는 것도 자동으로 계산해 줍니다.

평범한 값이나 변수를 넘겨서 TaskFlow 함수를 호출할 수도 있어요. 바로 아래 예시처럼요. 물론 이렇게 해서 기대한 대로 동작하긴 하는데, 코드 안의 로직이 실제로 도는 건 Dag이 실행될 때까지가 아니라는 점만 기억해 두세요. 그 전까지는 name 값이 태스크 파라미터로 보관될 뿐입니다.

@task
def hello_name(name: str):
    print(f'Hello {name}!')

hello_name('Airflow users')

TaskFlow를 더 자세히 배우고 싶다면 TaskFlow 튜토리얼을 확인해 보세요.

컨텍스트 (Context)

아래 예시처럼 키워드 인자로 추가하면 Airflow 컨텍스트 변수에 접근할 수 있어요.

from airflow.sdk import TaskInstance
from airflow.sdk.types import DagRunProtocol

@task
def print_ti_info(task_instance: TaskInstance, dag_run: DagRunProtocol):
    print(f"Run ID: {task_instance.run_id}")  # Run ID: scheduled__2023-08-09T00:00:00+00:00
    print(f"Task start date: {task_instance.start_date}")  # 2023-08-10 00:00:01+00:00
    print(f"Dag Run logical date: {dag_run.logical_date}")  # 2023-08-09 00:00:00+00:00

그 대신 태스크 시그니처에 **kwargs를 추가할 수도 있는데, 이렇게 하면 모든 Airflow 컨텍스트 변수를 kwargs 딕셔너리에서 꺼내 쓸 수 있어요.

from airflow.sdk import TaskInstance
from airflow.sdk.types import DagRunProtocol

@task
def print_ti_info(**kwargs):
    ti: TaskInstance = kwargs["task_instance"]
    print(f"Run ID: {ti.run_id}")  # Run ID: scheduled__2023-08-09T00:00:00+00:00
    print(f"Task start date: {ti.start_date}")  # 2023-08-10 00:00:01+00:00
    dr: DagRunProtocol = kwargs["dag_run"]
    print(f"Dag Run logical date: {dr.logical_date}")  # 2023-08-09 00:00:00+00:00

컨텍스트 변수의 전체 목록은 컨텍스트 변수 문서를 보세요.

로깅 (Logging)

태스크 함수 안에서 로깅을 쓰고 싶다면, 파이썬의 logging 시스템을 import 해서 쓰면 됩니다.

logger = logging.getLogger("airflow.task")

이렇게 만든 모든 로깅 라인은 태스크 로그에 기록돼요.

임의의 객체를 인자로 넘기기 (Passing Arbitrary Objects As Arguments)

Airflow 2.5.0부터 추가됨

앞에서 말했듯 TaskFlow는 변수를 태스크에 넘길 때 XCom을 사용해요. 그러다 보니 인자로 쓰는 변수는 직렬화(serialize)가 가능해야 한다는 조건이 붙습니다. Airflow는 기본적으로 내장 타입(int나 str 같은)을 모두 지원하고, @dataclass@attr.define으로 데코레이트된 객체도 지원해요. 아래 예시는 @attr.define으로 꾸며진 Asset을 TaskFlow와 함께 쓰는 모습을 보여줍니다.

Note

Asset을 쓰면 얻는 추가 이점이 있어요. 입력 인자로 쓰면 자동으로 inlet으로 등록되고, 태스크의 반환값이 Asset이거나 list[Asset]이라면 자동으로 outlet으로도 등록됩니다.

import json
import pendulum
import requests
from airflow import Asset
from airflow.sdk import dag, task

SRC = Asset("https://www.ncei.noaa.gov/access/monitoring/climate-at-a-glance/global/time-series/globe/land_ocean/ytd/12/1880-2022.json")

now = pendulum.now()

@dag(start_date=now, schedule="@daily", catchup=False)
def etl():
    @task()
    def retrieve(src: Asset) -> dict:
        resp = requests.get(url=src.uri)
        data = resp.json()
        return data["data"]

    @task()
    def to_fahrenheit(temps: dict[int, dict[str, float]]) -> dict[int, float]:
        ret: dict[int, float] = {}
        for year, info in temps.items():
            ret[year] = float(info["anomaly"]) * 1.8 + 32
        return ret

    @task()
    def load(fahrenheit: dict[int, float]) -> Asset:
        filename = "/tmp/fahrenheit.json"
        s = json.dumps(fahrenheit)
        f = open(filename, "w")
        f.write(s)
        f.close()
        return Asset(f"file:///{filename}")

    data = retrieve(SRC)
    fahrenheit = to_fahrenheit(data)
    load(fahrenheit)

etl()

커스텀 객체 (Custom Objects)

직렬화를 직접 제어하고 싶은 커스텀 객체가 있을 수 있어요. 보통은 클래스에 @dataclass@attr.define을 붙이면 Airflow가 알아서 처리해 줍니다. 다만 때로는 직접 직렬화를 제어하고 싶을 때가 있죠. 그럴 땐 클래스에 serialize() 메서드와 deserialize(data:dict, version:int) 정적 메서드를 추가하면 됩니다.

from typing import ClassVar

class MyCustom:
    __version__: ClassVar[int] = 1

    def __init__(self, x):
        self.x = x

    def serialize(self) -> dict:
        return dict({"x": self.x})

    @staticmethod
    def deserialize(data: dict, version: int):
        if version > 1:
            raise TypeError(f"version > {MyCustom.version}")
        return MyCustom(data["x"])

객체 버전 관리 (Object Versioning)

직렬화에 쓰이는 객체는 버전을 관리해 두는 게 좋은 습관이에요. 클래스에 __version__: ClassVar[int] = <x>를 추가하면 됩니다. Airflow는 클래스가 하위 호환(backwards compatible)된다고 가정해요. 그래서 버전 2가 버전 1을 역직렬화(deserialize)할 수 있다고 봅니다. 역직렬화에 커스텀 로직이 필요하다면 deserialize(data:dict, version:int)를 명시해 주세요.

Note

__version__의 타입 표기는 필수이고 ClassVar[int]여야 해요.

Sensor와 TaskFlow API (Sensors and the TaskFlow API)

Airflow 2.5.0부터 추가됨

TaskFlow API로 Sensor를 작성하는 예시는 TaskFlow API를 Sensor Operator와 함께 쓰기를 참고하세요.

역사 (History)

TaskFlow API는 Airflow 2.0에서 새로 등장했어요. 그래서 이전 버전의 Airflow로 작성된 Dag에는 비슷한 목적을 코드가 훨씬 더 많은 PythonOperator로 해결한 경우를 자주 만나게 됩니다.

TaskFlow API가 추가된 배경과 설계에 대한 더 자세한 내용은 Airflow 개선 제안서 AIP-31: “TaskFlow API” for clearer/simpler Dag definition에서 확인할 수 있어요.