Airflow 해부 (5) Operator, Connection, Sensor
PythonOperator와 BashOperator 등 주요 Operator의 위치, 접속 정보를 코드 밖으로 빼는 Connection과 Variable, 조건을 기다리는 Sensor와 deferrable의 개념을 정리합니다.
Airflow 해부 시리즈의 5편입니다. 전체 목차는 0편에 있습니다.
Operator: task를 만드는 부품
3편의 @task는 PythonOperator의 데코레이터 표현입니다. Python 함수 외의 작업에는 용도별 Operator가 있습니다.
| Operator | 용도 |
|---|---|
@task (PythonOperator) | Python 함수 실행. 기본 선택지 |
| BashOperator | 셸 명령 실행. 기존 스크립트를 그대로 태울 때 |
| 각종 SQL Operator | DB에 쿼리 실행 (provider 패키지로 제공) |
| DockerOperator, KubernetesPodOperator | task를 별도 컨테이너에서 실행. 의존성 격리가 필요할 때 |
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 선택, 멱등성 원칙을 다룹니다.