상위: Airflow
요약
Operator는 Airflow에서 Task를 실행하는 객체입니다. BashOperator, PythonOperator 등 다양한 종류가 있으며, 각각 특정 작업에 최적화되어 있습니다. Action, Sensor, Transfer, Database, Big Data, Container 등 7가지 카테고리로 분류되며, 사용자 정의 Operator도 만들 수 있습니다.
Operator란?
- Airflow에서 Task를 실행하는 역할을 수행하는 객체
- DAG 내에서 개별 Task로 사용됨
- 다양한 실행 방식 제공
- Python 기반으로 확장 가능
- 기본 제공(내장) 오퍼레이터
- 사용자 정의 오퍼레이터
Operator 종류
| 오퍼레이터 종류 | 설명 | 예제 |
|---|---|---|
| Action Operators | 특정 동작을 수행하는 오퍼레이터 | PythonOperator, BashOperator, EmailOperator |
| Sensor Operators | 특정 이벤트를 감지할 때까지 대기 | FileSensor, HttpSensor, S3KeySensor |
| Transfer Operators | 한 위치에서 다른 위치로 데이터 이동 | S3ToGCSOperator, MySQLToGCSOperator |
| Database Operators | DB에서 SQL을 실행하는 오퍼레이터 | PostgresOperator, MySqlOperator, SnowflakeOperator |
| Big Data & ML Operators | Spark, Hive, Dataproc, ML 관련 오퍼레이터 | SparkSubmitOperator, DataflowOperator |
| Docker & Kubernetes Operators | 컨테이너 환경에서 실행 | DockerOperator, KubernetesPodOperator |
| Dummy Operators | Task 흐름을 설정하는 데 사용 | DummyOperator (또는 EmptyOperator) |
BashOperator
Bash 명령어를 실행하는 Task를 정의할 때 사용
기본 사용법
from airflow.operators.bash import BashOperator
bash_t1 = BashOperator(
task_id="bash_t1",
bash_command="echo 'Hello, Airflow!'"
)
bash_t2 = BashOperator(
task_id="bash_t2",
bash_command="ls -al"
)
bash_t1 >> bash_t2 # 실행 순서 지정
외부 Shell Script 실행
Shell script 파일 생성:
# select_fruit.sh
FRUIT=$1
if [ $FRUIT == APPLE ]; then
echo "You selected Apple!"
elif [ $FRUIT == ORANGE ]; then
echo "You selected Orange!"
else
echo "You selected other Fruit!"
fi
실행 권한 부여:
chmod +x select_fruit.sh
docker-compose.yaml 내 plugins 경로 바인딩:
volumes:
- ./plugins:/opt/airflow/plugins
DAG에서 실행:
bash_task = BashOperator(
task_id='run_script',
bash_command="/opt/airflow/plugins/shell/select_fruit.sh ORANGE"
)
BashOperator 주요 파라미터
| 파라미터 | 설명 |
|---|---|
bash_command | 실행할 Bash 명령어 또는 스크립트 |
env | 환경 변수 딕셔너리 |
cwd | 작업 디렉토리 |
append_env | 기존 환경변수에 추가 여부 |
PythonOperator
Python 함수를 실행하는 Task를 정의할 때 사용
기본 사용법
from airflow.operators.python import PythonOperator
def my_function():
print("Hello, Airflow!")
python_t1 = PythonOperator(
task_id='python_t1',
python_callable=my_function
)
인자 전달
def greet(name, age):
print(f"Hello {name}, you are {age} years old")
python_task = PythonOperator(
task_id='greet_task',
python_callable=greet,
op_kwargs={'name': 'Alice', 'age': 30}
)
반환값 활용 (XCom)
def get_value():
return "Important Data"
def use_value(**context):
value = context['ti'].xcom_pull(task_ids='get_value')
print(f"Received: {value}")
task1 = PythonOperator(task_id='get_value', python_callable=get_value)
task2 = PythonOperator(task_id='use_value', python_callable=use_value)
task1 >> task2
자세한 내용은 XCom 참조
PythonOperator 주요 파라미터
| 파라미터 | 설명 |
|---|---|
python_callable | 실행할 Python 함수 |
op_args | 위치 인자 (list) |
op_kwargs | 키워드 인자 (dict) |
provide_context | Task Instance 컨텍스트 전달 (기본 True) |
EmailOperator
이메일 전송을 위한 오퍼레이터
기본 사용법
from airflow.operators.email import EmailOperator
email_t1 = EmailOperator(
task_id='email_t1',
to='[email protected]',
subject='Airflow 처리결과',
html_content='정상 처리되었습니다.<br/>'
)
Gmail 설정
사용 전 Gmail 설정 필요:
- Gmail → 설정 → 모든 설정보기 → 전달 및 POP/IMAP → IMAP 사용
- 계정관리 → 보안 → 2단계 인증 → 앱 비밀번호 생성
docker-compose.yaml 설정:
environment:
AIRFLOW__SMTP__SMTP_HOST: 'smtp.gmail.com'
AIRFLOW__SMTP__SMTP_USER: '[email protected]'
AIRFLOW__SMTP__SMTP_PASSWORD: '앱 비밀번호'
AIRFLOW__SMTP__SMTP_PORT: 587
AIRFLOW__SMTP__MAIL_FROM: '[email protected]'
EmptyOperator
실행 없이 DAG의 논리 흐름을 구성할 때 사용 (Airflow 2.x 이상)
사용 용도
- 시작/종료 Task
- 병렬 실행 구분점
- DAG 구조 시각화
기본 사용법
from airflow.operators.empty import EmptyOperator
start = EmptyOperator(task_id='start')
end = EmptyOperator(task_id='end')
task1 = BashOperator(task_id='task_1', bash_command='echo "Task 1"')
task2 = BashOperator(task_id='task_2', bash_command='echo "Task 2"')
start >> [task1, task2] >> end
Sensor Operators
특정 조건이 만족될 때까지 대기하는 Operator
FileSensor
파일 존재 여부를 확인
from airflow.sensors.filesystem import FileSensor
wait_for_file = FileSensor(
task_id='wait_for_file',
filepath='/path/to/file.txt',
poke_interval=30, # 30초마다 체크
timeout=600, # 10분 타임아웃
)
HttpSensor
HTTP 엔드포인트 응답 대기
from airflow.sensors.http_sensor import HttpSensor
wait_for_api = HttpSensor(
task_id='wait_for_api',
http_conn_id='api_connection',
endpoint='/api/status',
poke_interval=60,
)
Sensor 모드
| 모드 | 설명 | 리소스 |
|---|---|---|
poke | 워커가 계속 점유 (정밀도 높음) | 많이 사용 |
reschedule | 지정 시간 후 워커 반환 (리소스 절약) | 적게 사용 |
Database Operators
PostgresOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
create_table = PostgresOperator(
task_id='create_table',
postgres_conn_id='postgres_default',
sql='''
CREATE TABLE IF NOT EXISTS users (
id SERIAL PRIMARY KEY,
name VARCHAR(100)
);
'''
)
Transfer Operators
S3ToGCSOperator
S3에서 GCS로 데이터 전송
from airflow.providers.google.cloud.transfers.s3_to_gcs import S3ToGCSOperator
transfer_data = S3ToGCSOperator(
task_id='transfer_s3_to_gcs',
bucket='my-s3-bucket',
prefix='data/',
dest_gcs_bucket='my-gcs-bucket',
dest_gcs_prefix='imported/',
)
Container Operators
DockerOperator
Docker 컨테이너에서 작업 실행
from airflow.providers.docker.operators.docker import DockerOperator
docker_task = DockerOperator(
task_id='docker_task',
image='python:3.9',
command='python -c "print(\'Hello from Docker\')"',
auto_remove=True,
)
KubernetesPodOperator
Kubernetes Pod에서 작업 실행
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
k8s_task = KubernetesPodOperator(
task_id='k8s_task',
name='airflow-pod',
namespace='default',
image='python:3.9',
cmds=['python', '-c'],
arguments=['print("Hello from K8s")'],
)
사용자 정의 Operator
Custom Operator 생성
from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
class MyCustomOperator(BaseOperator):
@apply_defaults
def __init__(self, my_param, *args, **kwargs):
super().__init__(*args, **kwargs)
self.my_param = my_param
def execute(self, context):
print(f"Executing with param: {self.my_param}")
# 실제 작업 수행
return "Success"
# 사용
custom_task = MyCustomOperator(
task_id='custom_task',
my_param='value'
)
Operator 선택 가이드
언제 어떤 Operator를 사용할까?
| 작업 유형 | 추천 Operator |
|---|---|
| Shell 명령 실행 | BashOperator |
| Python 코드 실행 | PythonOperator |
| 파일 대기 | FileSensor |
| API 호출 대기 | HttpSensor |
| DB 쿼리 | PostgresOperator, MySqlOperator |
| 데이터 전송 | S3ToGCSOperator 등 Transfer |
| Spark 작업 | SparkSubmitOperator (Spark Integration 참조) |
| 컨테이너 실행 | DockerOperator, KubernetesPodOperator |
다음 단계
- Task - Operator로 Task 정의하기
- Task Dependencies - Operator 간 의존성 설정
- XCom - PythonOperator 간 데이터 전달
- Spark Integration - SparkSubmitOperator 사용법