전체 그래프
Airflow

Operators

data-engineeringtool/airflowoperator

상위: 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 OperatorsDB에서 SQL을 실행하는 오퍼레이터PostgresOperator, MySqlOperator, SnowflakeOperator
Big Data & ML OperatorsSpark, Hive, Dataproc, ML 관련 오퍼레이터SparkSubmitOperator, DataflowOperator
Docker & Kubernetes Operators컨테이너 환경에서 실행DockerOperator, KubernetesPodOperator
Dummy OperatorsTask 흐름을 설정하는 데 사용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_contextTask 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 설정 필요:

  1. Gmail → 설정 → 모든 설정보기 → 전달 및 POP/IMAP → IMAP 사용
  2. 계정관리 → 보안 → 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 사용법

참고 자료