포스트

Airflow 해부 (5) Operator, Connection, Sensor

PythonOperator와 BashOperator 등 주요 Operator의 위치, 접속 정보를 코드 밖으로 빼는 Connection과 Variable, 조건을 기다리는 Sensor와 deferrable의 개념을 정리합니다.

Airflow 해부 (5) Operator, Connection, Sensor

Airflow 해부 시리즈의 5편입니다. 전체 목차는 0편에 있습니다.

Operator: task를 만드는 부품

3편의 @task는 PythonOperator의 데코레이터 표현입니다. Python 함수 외의 작업에는 용도별 Operator가 있습니다.

Operator용도
@task (PythonOperator)Python 함수 실행. 기본 선택지
BashOperator셸 명령 실행. 기존 스크립트를 그대로 태울 때
각종 SQL OperatorDB에 쿼리 실행 (provider 패키지로 제공)
DockerOperator, KubernetesPodOperatortask를 별도 컨테이너에서 실행. 의존성 격리가 필요할 때
1
2
3
4
5
6
from airflow.providers.standard.operators.bash import BashOperator

dbt_run = BashOperator(
    task_id="dbt_run",
    bash_command="dbt run --profiles-dir /opt/dbt",
)

데코레이터 방식과 Operator 객체 방식은 한 DAG 안에서 섞을 수 있고, 의존성은 >> 연산자로도 연결합니다.

1
extract() >> dbt_run >> train()

DB, 클라우드, Slack 등 외부 시스템용 Operator는 본체가 아니라 provider 패키지로 분리되어 있습니다. PostgreSQL을 쓰려면 apache-airflow-providers-postgres를 별도 설치하는 식입니다. 설치할 때는 본체와 마찬가지로 constraint를 붙입니다.

Connection: 접속 정보를 코드 밖으로

DB 비밀번호를 DAG 코드에 넣지 않기 위한 장치가 Connection입니다. UI(Admin, Connections) 또는 CLI로 등록하고, 코드에서는 id로만 부릅니다.

1
2
3
4
5
airflow connections add warehouse_db \
  --conn-type postgres \
  --conn-host localhost --conn-port 5432 \
  --conn-login analyst --conn-password '...' \
  --conn-schema warehouse
1
2
3
4
5
6
from airflow.providers.postgres.hooks.postgres import PostgresHook

@task
def load_to_db(path: str) -> None:
    hook = PostgresHook(postgres_conn_id="warehouse_db")
    hook.run("COPY features FROM %s", parameters=[path])

Hook은 Connection을 읽어 실제 접속을 만들어주는 클라이언트입니다. Operator가 “task 껍데기”라면 Hook은 “접속 도구”라서, @task 함수 안에서 Hook만 꺼내 쓰는 조합이 실무에서 가장 흔합니다.

비슷한 장치로 Variable이 있습니다. 환경별 설정값(버킷 이름, 임계값)을 UI에서 관리하고 코드에서 Variable.get("bucket_name")으로 읽습니다. 접속 정보는 Connection, 일반 설정은 Variable로 나눠 담습니다.

Sensor: 조건을 기다리는 task

“파일이 도착하면 시작”처럼 외부 조건을 기다리는 task가 Sensor입니다.

1
2
3
4
5
6
7
8
9
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor

wait_for_file = S3KeySensor(
    task_id="wait_for_raw",
    bucket_name="raw-data",
    bucket_key="trips/{{ ds }}.parquet",
    poke_interval=300,
    timeout=60 * 60 * 6,
)
  • poke_interval: 조건 확인 주기입니다. 5분마다 S3를 확인합니다
  • timeout: 이 시간까지 조건이 안 되면 실패 처리합니다. timeout 없는 Sensor는 영원히 기다리며 슬롯을 점유합니다

기본 모드의 Sensor는 기다리는 동안에도 worker 슬롯을 하나 차지합니다. Sensor가 많아지면 기다리는 task가 슬롯을 다 먹어 정작 일할 task가 못 도는 상황이 생깁니다. 이 문제의 해법이 deferrable 모드로, 대기 중에 슬롯을 반납하고 triggerer라는 별도 컴포넌트가 조건을 대신 감시합니다. Sensor를 몇 개 이상 쓰게 되면 deferrable 지원 여부를 확인하는 것이 좋습니다.

실패 알림 붙이기

3편에서 미룬 알림입니다. Slack provider를 설치하고 Connection을 등록하면 콜백 한 줄로 연결됩니다.

1
2
3
4
5
6
7
8
9
10
from airflow.providers.slack.notifications.slack import send_slack_notification

@dag(
    ...,
    on_failure_callback=send_slack_notification(
        slack_conn_id="slack_alerts",
        text="DAG {{ dag.dag_id }} 실패: {{ ds }}",
        channel="#data-alerts",
    ),
)

다음 편은 운영입니다. 컨테이너 구동과 executor 선택, 멱등성 원칙을 다룹니다.

다음 글: Airflow 해부 (6) 운영: 컨테이너 구동, executor, 멱등성

이 기사는 저작권자의 CC BY 4.0 라이센스를 따릅니다.