직렬화

직렬화 (Serialization)

이 페이지는 Airflow가 Task 사이의 인자 같은 데이터 교환을 위해 직렬화/역직렬화를 수행하는 방식을 다뤄요. 직렬화는 또한 webserver와 scheduler(Dag processor와 대조적으로)가 Dag 파일을 읽을 필요 없게 하기 위해서도 일어나며, 이는 보안과 효율성을 위한 거예요. 커스텀 직렬화 해석 순서와 Airflow 객체 예제를 설명해요.

출처: 문서

본문

Task 사이의 인자 같은 데이터 교환을 지원하기 위해, Airflow는 교환할 데이터를 직렬화하고 다운스트림 Task에서 필요할 때 다시 역직렬화해야 해요. 직렬화는 또한 webserver와 scheduler(Dag processor와 대조적으로)가 Dag 파일을 읽을 필요가 없도록 하기 위해서도 일어나요. 이는 보안 목적과 효율성을 위한 것이에요.

직렬화는 놀랍도록 어려운 작업이에요. Python은 기본적으로 strint 같은 프리미티브와 iterable의 순환만 직렬화를 지원해요. 더 복잡해지면 커스텀 직렬화가 필요해요.

직렬화 해석 순서

Airflow는 커스텀 직렬화를 다음 순서로 해석해요:

  1. 프리미티브 값(str이나 int 같은)과 프리미티브의 iterable은 추가 인코딩 없이 그대로 반환돼요.
  2. 객체가 프리미티브가 아니면, Airflow는 airflow.sdk.serde.serializers 네임스페이스에서 등록된 직렬화기와 역직렬화기를 찾아요.
  3. 등록된 직렬화기가 없으면, Airflow는 객체가 serialize() 메서드(그리고 역직렬화를 위해 대응하는 deserialize(data, version: int) 메서드)를 정의하는지 확인해요.
  4. 마지막으로, 객체가 @dataclass@attr.define으로 데코레이트되어 있으면, Airflow는 그 데코레이터가 제공하는 공개 필드를 사용해 객체를 직렬화해요.

Airflow를 새 직렬화기로 확장하려 한다면, 언제 어떤 직렬화 방식을 선택해야 하는지 아는 것이 좋아요. Airflow의 제어 아래 있는 객체, 즉 airflow.model.dag.DAG 같은 airflow.* 네임스페이스 아래 있거나 개발자 제어 아래 있는 객체(예: my.company.Foo)는 먼저 @attr.define이나 @dataclass로 장식할 수 있는지 살펴봐야 해요. 그것이 불가능하면 serializedeserialize 메서드를 구현해야 해요. serialize 메서드는 프리미티브나 dict를 반환해야 해요. dict 안의 값을 직렬화할 필요는 없어요(그것은 알아서 처리되므로), 하지만 키는 프리미티브 형태여야 해요.

Airflow의 제어 아래 있지 않은 객체, 예를 들어 numpy.int16은 등록된 직렬화기와 역직렬화기가 필요해요. 버전 관리가 필수예요. bytes를 제외한 프리미티브는 dict로 반환될 수 있어요. 다시 말해 dict 값은 직렬화될 필요가 없지만, 키는 프리미티브 형태여야 해요. 등록된 직렬화기를 구현할 때는 순환 import가 없도록 특별히 주의해요. 보통 직렬화기 목록을 채울 때 str을 사용해 이를 피할 수 있어요. 이렇게: serializers = ["my.company.Foo"] 대신 serializers = [Foo].

Note

직렬화·역직렬화는 속도에 의존적이에요. 가능한 한 dict 같은 내장 함수를 많이 사용하고 클래스나 다른 복잡한 구조를 피하세요.

Airflow 객체

from typing import Any, ClassVar

class Foo:
    __version__: ClassVar[int] = 1

    def __init__(self, a, v) -> None:
        self.a = a
        self.b = {"x": v}

    def serialize(self) -> dict[str, Any]:
        return {
            "a": self.a,
            "b": self.b,
        }

    @staticmethod
    def deserialize(data: dict[str, Any], version: int):
        f = Foo(a=data["a"], v=data["b"])
        return f

등록됨

from __future__ import annotations

from typing import TYPE_CHECKING

from airflow.sdk.module_loading import qualname

if TYPE_CHECKING:
    import decimal

    from airflow.sdk.serde import U

serializers = [
    "decimal.Decimal"
]  # this can be a type or a fully qualified str. Str can be used to prevent circular imports
deserializers = serializers  # in some cases you might not have a deserializer (e.g. k8s pod)

__version__ = 1  # required

# the serializer expects output, classname, version, is_serialized?
def serialize(o: object) -> tuple[U, str, int, bool]:
    from decimal import Decimal

    if not isinstance(o, Decimal):
        return "", "", 0, False
    name = qualname(o)
    _, _, exponent = o.as_tuple()
    if isinstance(exponent, int) and exponent >= 0:  # No digits after the decimal point.
        return int(o), name, __version__, True
    # Technically lossy due to floating point errors, but the best we
    # can do without implementing a custom encode function.
    return float(o), name, __version__, True

# the deserializer sanitizes the data for you, so you do not need to deserialize values yourself
def deserialize(cls: type, version: int, data: object) -> Decimal:
    from decimal import Decimal

    # always check version compatibility
    if version > __version__:
        raise TypeError(f"serialized {version} of {qualname(cls)} > {__version__}")

    if cls is not Decimal:
        raise TypeError(f"do not know how to deserialize {qualname(cls)}")

    return Decimal(str(data))

더 알아보기 (Learn more)