오퍼레이터, 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
'Data 엔지니어링 > Airflow' 카테고리의 다른 글
| [Airflow] @task 데코레이터 (0) | 2025.01.09 |
|---|---|
| [Airflow] Python 오퍼레이터 & 외부 함수 사용하기 (0) | 2025.01.09 |
| [Airflow] Email Operator로 메일 전송하기 (1) | 2025.01.09 |
| [Airflow] Bash Operator & 외부 쉘파일 수행하기 (0) | 2025.01.09 |
| [Airflow] Airflow Docker 설치 방법 (0) | 2025.01.06 |