Object Storage를 활용한 클라우드 네이티브 워크플로우

Object Storage를 활용한 클라우드 네이티브 워크플로우 (Cloud-Native Workflows with Object Storage)

이 튜토리얼은 Object Storage API를 소개해요. Amazon S3, GCS, Azure Blob Storage 같은 클라우드 스토리지에 provider 전용 SDK나 낮은 수준의 자격 증명 관리 없이 읽고 쓸 수 있게 해주는 API랍니다.

출처: 문서

본문

버전 2.8에 추가됨.

우리 Airflow 시리즈의 마지막 튜토리얼에 오신 것을 환영해요! 이제 Python과 TaskFlow API로 DAG을 만들고, XCom으로 데이터를 전달하고, Task들을 명확하고 재사용 가능한 워크플로우로 연결해 봤을 거예요.

이 튜토리얼에서는 한 단계 더 나아가 Object Storage API를 소개할게요. 이 API는 Amazon S3, Google Cloud Storage (GCS), Azure Blob Storage 같은 클라우드 스토리지에 provider 전용 SDK나 낮은 수준의 자격 증명 관리에 신경 쓰지 않고 읽고 쓸 수 있게 해줘요.

실제 사용 사례를 함께 살펴볼 거예요:

  1. 공개 API에서 데이터 가져오기
  2. 그 데이터를 Parquet 형식으로 object storage에 저장하기
  3. DuckDB로 SQL을 사용해 분석하기

그 과정에서 새로운 ObjectStoragePath 추상화를 강조하고, Airflow가 connection을 통해 클라우드 자격 증명을 어떻게 처리하는지 설명하며, 이것이 이식 가능하고 클라우드에 구애받지 않는 파이프라인을 어떻게 가능하게 하는지 보여줄게요.

왜 이런 게 중요한가요?

많은 데이터 워크플로우는 파일에 의존해요 — 원시 CSV든, 중간 Parquet 파일이든, 모델 아티팩트든 말이죠. 전통적으로는 이에 대해 S3 전용이나 GCS 전용 코드를 작성해야 했어요. 이제 ObjectStoragePath를 사용하면 올바른 Airflow connection만 구성했다면 provider와 무관하게 동작하는 일반 코드를 작성할 수 있어요.

시작해 볼까요!

사전 요구사항 (Prerequisites)

시작하기 전에 다음이 필요해요:

  • DuckDB, 인프로세스 SQL 데이터베이스: pip install duckdb로 설치
  • Amazon S3 접근s3fs가 포함된 Amazon Provider: pip install apache-airflow-providers-amazon[s3fs] (스토리지 URL 프로토콜을 바꾸고 관련 provider를 설치해 선호하는 provider로 대체할 수 있어요.)
  • Pandas 테이블 데이터 작업용: pip install pandas

ObjectStoragePath 만들기

이 튜토리얼의 핵심은 클라우드 object store의 경로를 처리하기 위한 새 추상화인 ObjectStoragePath예요. pathlib.Path와 비슷하지만 파일시스템 대신 버킷용이라고 생각하면 돼요.

airflow/example_dags/tutorial_objectstorage.py[source]

base = ObjectStoragePath("s3://aws_default@airflow-tutorial-data/")

URL 구문은 간단해요: protocol://bucket/path/to/file

  • protocol(s3, gs, azure 등)이 백엔드를 결정해요
  • URL의 "username" 부분은 conn_id가 될 수 있으며, Airflow에 인증 방법을 알려줘요
  • conn_id를 생략하면 Airflow는 해당 백엔드의 기본 connection으로 폴백해요

명확성을 위해 conn_id를 키워드 인자로 제공할 수도 있어요:

ObjectStoragePath("s3://airflow-tutorial-data/", conn_id="aws_default")

이것은 다른 곳(예: Asset)에서 정의된 경로를 재사용할 때나, connection이 URL에 포함되어 있지 않을 때 특히 유용해요. 키워드 인자가 항상 우선해요.

Tip

ObjectStoragePath를 전역 DAG 스코프에서 안전하게 만들 수 있어요. Connection은 경로가 생성될 때가 아니라 사용될 때만 해석돼요.

Object Storage에 데이터 저장하기

데이터를 가져와 클라우드에 저장해 볼게요.

airflow/example_dags/tutorial_objectstorage.py[source]

    @task
    def get_air_quality_data(**kwargs) -> ObjectStoragePath:
        """
        #### Get Air Quality Data
        This task gets air quality data from the Finnish Meteorological Institute's
        open data API. The data is saved as parquet.
        """
        import pandas as pd

        logical_date = kwargs["logical_date"]

        latitude = 28.6139
        longitude = 77.2090

        params: Mapping[str, str | float] = {
            "latitude": latitude,
            "longitude": longitude,
            "hourly": ",".join(aq_fields.keys()),
            "timezone": "UTC",
        }

        response = requests.get(API, params=params)
        response.raise_for_status()

        data = response.json()
        hourly_data = data.get("hourly", {})

        df = pd.DataFrame(hourly_data)

        df["time"] = pd.to_datetime(df["time"])

        # ensure the bucket exists
        base.mkdir(exist_ok=True)

        formatted_date = logical_date.format("YYYYMMDD")
        path = base / f"air_quality_{formatted_date}.parquet"

        with path.open("wb") as file:
            df.to_parquet(file)
        return path

여기서 일어나는 일이에요:

  • 핀란드 기상청의 헬싱키 대기질 데이터 공개 API를 호출해요
  • JSON 응답을 pandas DataFrame으로 파싱해요
  • Task의 logical date를 기반으로 파일명을 생성해요
  • ObjectStoragePath를 사용해 데이터를 Parquet으로 클라우드 스토리지에 직접 기록해요

이것은 전형적인 TaskFlow 패턴이에요. object key가 매일 바뀌므로 매일 실행하며 시간이 지남에 따라 데이터셋을 만들 수 있어요. 다운스트림 Task에서 사용할 최종 object 경로를 반환해요.

이것이 멋진 이유: boto3도, GCS 클라이언트 설정도, 자격 증명 뒤섞기도 필요 없어요. 스토리지 백엔드를 넘나드는 단순한 파일 의미론만 있으면 돼요.

DuckDB로 데이터 분석하기

이제 SQL을 사용해 DuckDB로 그 데이터를 분석해 볼게요.

airflow/example_dags/tutorial_objectstorage.py[source]

    @task
    def analyze(
        path: ObjectStoragePath,
    ):
        """
        #### Analyze
        This task analyzes the air quality data, prints the results
        """
        import duckdb

        conn = duckdb.connect(database=":memory:")
        conn.register_filesystem(path.fs)
        s3_path = path.path
        conn.execute(
            f"CREATE OR REPLACE TABLE airquality_urban AS SELECT * FROM read_parquet('{path.protocol}://{s3_path}')"
        )

        df2 = conn.execute("SELECT * FROM airquality_urban").fetchdf()

        print(df2.head())

주목할 몇 가지 핵심 사항:

  • DuckDB는 Parquet 읽기를 기본 지원해요
  • DuckDB와 ObjectStoragePath 모두 fsspec에 의존하므로 object storage 백엔드를 쉽게 등록할 수 있어요
  • path.fs를 사용해 올바른 파일시스템 객체를 가져와 DuckDB에 등록해요
  • 마지막으로 SQL을 사용해 Parquet 파일을 질의하고 pandas DataFrame을 반환해요

함수가 경로를 수동으로 다시 만들지 않고, Xcom을 사용해 업스트림 Task에서 전체 경로를 가져온다는 점을 눈여겨보세요. 이렇게 하면 Task가 이식 가능하고 이전 로직에서 분리돼요.

모두 함께 묶기

모든 것을 묶는 전체 DAG은 다음과 같아요:

airflow/example_dags/tutorial_objectstorage.py[source]

import pendulum
import requests

from airflow.sdk import ObjectStoragePath, dag, task

API = "https://air-quality-api.open-meteo.com/v1/air-quality"

aq_fields = {
    "pm10": "float64",
    "pm2_5": "float64",
    "carbon_monoxide": "float64",
    "nitrogen_dioxide": "float64",
    "sulphur_dioxide": "float64",
    "ozone": "float64",
    "european_aqi": "float64",
    "us_aqi": "float64",
}
base = ObjectStoragePath("s3://aws_default@airflow-tutorial-data/")


@dag(
    schedule=None,
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    catchup=False,
    tags=["example"],
)
def tutorial_objectstorage():
    """
    ### Object Storage Tutorial Documentation
    This is a tutorial DAG to showcase the usage of the Object Storage API.
    Documentation that goes along with the Airflow Object Storage tutorial is
    located
    [here](https://airflow.apache.org/docs/apache-airflow/stable/tutorial/objectstorage.html)
    """
    @task
    def get_air_quality_data(**kwargs) -> ObjectStoragePath:
        """
        #### Get Air Quality Data
        This task gets air quality data from the Finnish Meteorological Institute's
        open data API. The data is saved as parquet.
        """
        import pandas as pd

        logical_date = kwargs["logical_date"]

        latitude = 28.6139
        longitude = 77.2090

        params: Mapping[str, str | float] = {
            "latitude": latitude,
            "longitude": longitude,
            "hourly": ",".join(aq_fields.keys()),
            "timezone": "UTC",
        }

        response = requests.get(API, params=params)
        response.raise_for_status()

        data = response.json()
        hourly_data = data.get("hourly", {})

        df = pd.DataFrame(hourly_data)

        df["time"] = pd.to_datetime(df["time"])

        # ensure the bucket exists
        base.mkdir(exist_ok=True)

        formatted_date = logical_date.format("YYYYMMDD")
        path = base / f"air_quality_{formatted_date}.parquet"

        with path.open("wb") as file:
            df.to_parquet(file)
        return path
    @task
    def analyze(
        path: ObjectStoragePath,
    ):
        """
        #### Analyze
        This task analyzes the air quality data, prints the results
        """
        import duckdb

        conn = duckdb.connect(database=":memory:")
        conn.register_filesystem(path.fs)
        s3_path = path.path
        conn.execute(
            f"CREATE OR REPLACE TABLE airquality_urban AS SELECT * FROM read_parquet('{path.protocol}://{s3_path}')"
        )

        df2 = conn.execute("SELECT * FROM airquality_urban").fetchdf()

        print(df2.head())
    obj_path = get_air_quality_data()
    analyze(obj_path)
tutorial_objectstorage()

이 DAG을 트리거하고 Airflow UI의 Graph View에서 볼 수 있어요. 각 Task는 입력과 출력을 명확하게 기록하며, Xcom 탭에서 반환된 경로를 검사할 수 있어요.

다음에 살펴볼 내용

더 나아갈 수 있는 몇 가지 방법:

  • object sensors(S3KeySensor 같은)를 사용해 외부 시스템이 업로드한 파일을 기다리기
  • S3-to-GCS 전송이나 리전 간 데이터 동기화 오케스트레이션
  • 누락되거나 손상된 파일을 처리하는 분기 로직 추가
  • CSV나 JSON 같은 다른 형식 실험

참고 (See Also)

더 알아보기 (Learn more)