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

[Airflow] 오퍼레이터 & Cron스케줄 & Task 연결

by imkus 2025. 1. 8.

오퍼레이터, Cron, Task 관련하여 핵심내용만 간추려 작성하였습니다.

 

오퍼레이터

DAG의 구성

  • 오퍼레이터
    특정 행위를 할 수 있는 기능을 모아 놓은 클래스, 설계도
  • Task
    오퍼레이터에서 객체화(인스턴스화)되어 DAG에서 실행가능한 오브젝트
  • Bash 오퍼레이터
    쉘 스크립트 명령을 수행하는 오퍼레이터
  • Python 오퍼레이터
    Python 함수를 실행시켜주는 오퍼레이터
  • S3 오퍼레이터
    AWS S3 솔루션을 컨트롤할 수 있도록 하는 오퍼레이터
  • GCS 오퍼레이터
    구글클라우드 GCS를 다룰 수 있게 해주는 오퍼레이터

 

 

Task의 수행 주체

  • 스케줄러
    : DAG Parsing 후 DB에 정보저장
    : DAG 시작시간 결정
  • 워커
    : 실제 작업 수행

 

 

DAG 오퍼레이터 생성예시

/*dags_bash_operator.py*/

from __future__ import annotations

import datetime

import pendulum   # datetime을 좀더 편하게 사용할 수 있게해주는 라이브러리

from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.empty import EmptyOperator



# DAG에 대한 정의
with DAG(
    dag_id="dags_bash_oprator", # 화면에 보이는 id, dag_id와 파일명은 왠만하면 같게 해주는 것이 좋음(필수는 아님)
    schedule="0 0 * * *", # 분 시 일 월 요일 정의
    start_date=pendulum.datetime(2021, 1, 1, tz="Asia/Seoul"), # DAG이 언제부터 돌도록 할 것인지 정의, tz는 UTC 기준이 아닌 Asia/Seoul로 해야함
    catchup=False, # start_date 부터 현재일자 까지 누락된 것들을 모두 돌릴 것인지에 대한 여부 (단, 누락된 일자만큼 모두 동시에 돌아가므로 일반적으로 false로 두는 것이 좋음)
    dagrun_timeout=datetime.timedelta(minutes=60), # DAG이 몇분 이상 돌면 timeOut 되는지 설정, 필요 없을 경우 삭제해도 됨
    tags=["example", "example2"], # Airflow 웹서비스 화면에 dag_id 밑에 보이는 태그들을 정의 (태그만 눌렀을 때 같은 태그를 같은 DAG만 보이도록 하는 용도)
    params={"example_key": "example_value"}, # Task가 공통적으로 넘겨줄 파라미터가 있다면 여기서 작성해주면 됨
) as dag:
    bash_t1 = BashOperator(
        task_id = "bash_t1", # airflow graph에 보이는 task id ,객체명과 task_id는 같게 하는 것이 좋음
        bash_command = "echo whoami", # 어떤 스크립트를 수행할 것이지 정의
    )


    bash_t2 = BashOperator(
        task_id = "bash_t2", # airflow graph에 보이는 task id ,객체명과 task_id는 같게 하는 것이 좋음
        bash_command = "echo $HOSTNAME", # 어떤 스크립트를 수행할 것이지 정의
    )


    # task가 동작하는 순서 정의
    bash_t1 >> bash_t2

생성한 DAG 파일은 docker-compose.yaml 파일에 명시되어 있는 로컬상의 dags 디렉토리에 넣으면 됩니다.

 

":" 문자를 기준으로 우측은 컨테이너 내부의 경로이고 좌측은 로컬상의 경로입니다.

 

 

 

 

Cron 스케줄

Cron 스케줄 개념

  • task가 실행되어야 하는 시간(주기)을 정하기 위한 다섯개의 필드로 구성된 문자열

 

 

 

 

 

Task 연결

Task 연결 방법

  • Task 연결 방법 종류
    1)  >> , << 사용하기 (Airflow 공식 추천방식)
    2) 함수 사용하기

 

복잡한 Task 연결 방법

 

Task 연결 예제 코드

/*dags_conn_test.py*/


from airflow import DAG
import pendulum
import datetime
from airflow.operators.empty import EmptyOperator


with DAG(
    dag_id = "dags_conn_test",
    schedule = None,
    start_date = pendulum.datetime(2023, 3, 1, tz = "Asia/Seoul"),
    catchup = False
) as dag:
    t1 = EmptyOperator(
        task_id = "t1"
    )

    t2 = EmptyOperator(
        task_id = "t2"
    )

    t3 = EmptyOperator(
        task_id = "t3"
    )

    t4 = EmptyOperator(
        task_id = "t4"
    )

    t5 = EmptyOperator(
        task_id = "t5"
    )

    t6 = EmptyOperator(
        task_id = "t6"
    )

    t7 = EmptyOperator(
        task_id = "t7"
    )

    t8 = EmptyOperator(
        task_id = "t8"
    )


    t1 >> [t2, t3] >> t4
    t5 >> t4
    [t4, t7] >> t6 >> t8

 

Task 연결 예제 코드로 각 태스크를 연결한 결과

 

 

DAGs 관련 자세한 내용

https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/dags.html

 

DAGs — Airflow Documentation

 

airflow.apache.org