본문 바로가기
Data 엔지니어링/Airflow

[Airflow] Python 오퍼레이터 & 외부 함수 사용하기

by imkus 2025. 1. 9.

1. 파이썬 오퍼레이터

  • 무엇을 하는 오퍼레이터인가?
    정의된 파이썬 함수를 실행시키는 오퍼레이터
  • 가장 많이 쓰이는 오퍼레이터
  • 라이브러리 가져오기
    : from airflow.operators.python import PythonOperator

2. 파이썬 오퍼레이터 종류

패키지 오퍼레이터 중요도 설명
airflow.operators.python PythonOperator ⭐ ⭐ ⭐ 어떤 파이썬 함수를 실행시키기 위한 오퍼레이터
BranchPythonOperator 파이썬 함수 실행 결과에 따라 task를 선택적으로 실행시킬 때 사용되는 오퍼레이터
ShortCircuitOperator   파이썬 함수 실행 결과에 따라 후행 Task를 실행하지 않고 종료시킬 수 있는 오퍼레이터
PythonVirtualenvOperator   파이썬 가상환경 생성후 job 수행하고 마무리되면 가상환경을 삭제해주는 오퍼레이터
ExternalPythonOperator   기존에 존재하는 파이썬 가상환경에서 Job 수행하게 하는 오퍼레이터

 

 

3. 파이썬 오퍼레이터 예제코드

/*dags_python_operator.py*/

from airflow import DAG
import pendulum
import datetime
from airflow.operators.python import PythonOperator
import random


with DAG(
    dag_id = "dags_python_operator",
    schedule = "30 6 * * *",
    start_date = pendulum.datetime(2023, 3, 1, tz="Asia/Seoul"),
    catchup = False
) as dag:
    def select_fruit():
        fruit = ['APPLE', 'BANANA', 'ORANGE', 'AVOCADO']
        rand_int = random.randint(0,3)
        print(fruit[rand_int])

    py_t1 = PythonOperator(
        task_id = 'py_t1',
        python_callable=select_fruit
    )

    py_t1

 

실행 결과

 

 

 

 

4. 파이썬 모듈 경로 이해

from airflow.operators.python import PythonOperator
# Airflow 폴더 아래 operators 폴더 아래 python 파일 아래에서 PytonOperator 클래스를 가져와라
  • Dag에서 우리가 만든 외부 함수를 import 해와야 되는데 import 경로를 어떻게 작성해야 하는지 알려면 파이썬 모듈 경로에 관하여 이해 해야함
  • 파이썬은 sys.path 변수에서 모듈의 위치를 검색

  • sys.path 에 값을 추가하는 방법 (귀찮은 방법)
    1) 명시적으로 추가 (ex: sys.path.append('/home/hjkim'))
    2) OS 환경변수 PYTHONPATH에 값을 추가
  • Airflow는 자동적으로 dags 폴더와 plugins 폴더를 sys.path에 추가함
    (컨테이너에서 airflow info 명령을 수행해보면 아래와 같음)

 

5. plugins 폴더 이용하기

  • 장점
    1) 공통함수 작성
    2) 재활용성 증가
    3) DAG 깔끔해짐

 

 

 

6. 외부함수 코드예제

"""
common_func.py
plugins/common/ 경로에 작성
"""

def get_sftp():
    print('sftp 작업을 시작합니다.')

 

 

7. 외부함수 사용 Dags 코드예제

"""
dags_python_import_func.py
dags/ 경로에 작성
"""

from airflow import DAG
import pendulum
import datetime
from airflow.operators.python import PythonOperator
from common.common_func import get_sftp  #

with DAG(
    dag_id="dags_python_import_func",
    schedule="30 6 * * *",
    start_date=pendulum.datetime(2023, 3, 1, tz="Asia/Seoul"),
    catchup=False
) as dag:
    task_get_sftp = PythonOperator(
        task_id='task_get_sftp',
        python_callable=get_sftp
    )

 

 

8.  실행결과