1. Airflow의 스케줄링

각 DAG에 대한 스케줄 간격을 정의하여 파이프라인이 실행되는 시간을 결정할 수 있습니다.

DAG의 스케줄을 잡기 위해서는 start_date (필수), schedule_interval (필수), end_date 등을 설정해야 합니다.

먼저 해당 DAG이 언제부터 시작할지를 정해야 합니다. (start_date 설정)

이후 schedule_interval을 설정해서 실행 간격을 정합니다. 만약 스케줄이 특정 날짜까지만 실행되기를 원하면, end_date를 설정하면 됩니다.

start_date + schedule_interval 때 첫 실행이 일어납니다.

schedule_interval은 Preset을 이용하거나 Cron을 이용해 설정할 수 있습니다. 주로 Cron 방식으로 표현됩니다.

schedule_interval 은 Preset을 이용하거나 cron을 이용해 설정할 수 있습니다. 주로 cron 방식으로 표현합니다.

Preset

Preset 의미
@once 한 번 실행
@hourly 한 시간에 한 번
@daily 하루에 한번, 자정에 실행
@weekly 일주일에 한 번, 일요일 자정에 실행
@monthly 한 달에 한 번, 그 달의 첫 날 자정에 실행
@yearly 1년에 한번씩 1월 1일 자정에 실행

cron

’* * * * * *’
‘분 시간 일 월 요일(0~6)’

Cron 표현식 의미 예시
* * * * * 매 분마다 실행 Every minute
0 * * * * 매 시간 정각에 실행 Hourly at minute 0
0 0 * * * 매일 자정에 실행 Daily at midnight
0 0 * * 0 매주 일요일 자정에 실행 Weekly on Sunday at midnight
0 0 1 * * 매달 첫째 날 자정에 실행 Monthly on the 1st at midnight
0 0 1 1 * 매년 1월 1일 자정에 실행 Yearly on January 1st at midnight

예시)

# 저는 taskflow API를 사용했습니다
# taskflow는 간단하게 데코레이터를 사용해 DAG와 Task를 구성하는 방식입니다.

# 매일 04시30 분 실행
default_args = {
    'owner': 'jinwoo',
    'start_date': datetime(2023, 11, 7, tzinfo=local_tz),  # 시작 날짜
}
@dag(
    schedule_interval="30 4 * * *",
    default_args=default_args,
)

# trigger시 실행
default_args = {
    'owner': 'jinwoo',
    'start_date': datetime(2023, 11, 7, tzinfo=local_tz),  # 시작 날짜
}
@dag(
    schedule_interval= None,
    default_args=default_args,
)


2. backfill, catch up 사용하기

Airflow를 사용하다 보면, 재실행하거나 현재 시점보다 과거의 배치 작업을 실행해야 할 때, 실패 등의 이유로 특정 작업만 재실행하고 싶을 경우가 발생한다. 이 경우 catch up, backfill, clear를 적절하게 사용하면 된다

catch up

python 코드로 dag을 작성할 때 dag 속성 안에 catch up이라는 인수를 둘 수 있습니다. 기본값을 False이고 True를 주게 되면 활성화가 됩니다. catch up은 start_date ~ 현재(또는 end_date) 을 순차적으로 실행시킬때 사용합니다. 예를 들면 start_date = 2024-08-02 이고 현재 날짜(또는 end_date)가 2024-08-05이고 schedule_interval이 매일 새벽 3시라면 2일부터 5일까지 순차적으로 실행시킵니다.


catch up 주의 사항

catch up이 실행될 때 한번에 실행되기 때문에 서버에 과부하가 올 수 있습니다. 따라서 max_active_runs으로 한번에 실행되는 DAGrun의 수를 설정할 수 있습니다

default_args = {
    'owner': 'jinwoo',
    'start_date': datetime(2023, 11, 7, tzinfo=local_tz),  # 시작 날짜
}

@dag(
    schedule_interval="30 4 * * *",
    default_args=default_args,
    max_active_runs=1, # 한번에 실행되는 DAGrun 수 설정
    catchup=True,  # 반복시 설정
)


backfill

Backfill은 DAG가 이미 배포되어 실행 중일 때, DAG의 시작 날짜 이전의 데이터를 처리하길 원할 때 사용합니다. CLI를 사용하여 backfill을 사용할 수 있습니다

airflow dags backfill [-h] [-c CONF] [--delay-on-limit DELAY_ON_LIMIT] [-x]
                      [-n] [-e END_DATE] [-i] [-I] [-l] [-m] [--pool POOL]
                      [--rerun-failed-tasks] [--reset-dagruns] [-B]
                      [-s START_DATE] [-S SUBDIR] [-t TASK_REGEX] [-v] [-y]
                      dag_id
                   
# 예시
airflow dags backfill -s 2024-09-01 -e 2024-09-02 example_dag
인수 설명
dag_id dag의 id
-c Dagrun의 conf 속성을 피클
--continue-on-failure 설정하면 일부 작업이 실패해도 백필 계속 진행
--delay-on-limit dag 실행을 다시 실행하기 전에 최대 활성 Dag실행 제한에 도달했을 때 대기하는 시간
--end-date end-date yyyy-mm-dd
-i 업스트림 작업을 건너뛰고 정규표현식과 일치하는 작업만 실행
-I dependencies_on_past 속성 무시
-l LocalExecutor로 실행
-m 작업을 실행하지 않고 성공한 것으로 표시
--pool 사용할 리소스 풀
--verbose 더 자세한 로깅 출력 만들기
-s start_date yyyy-mm-dd로 정의
--reset-dagruns 설정된 경우 백필은 기존 백필 관련 DAG 실행을 삭제하고 새로 시작
--rerun-failed-tasks 설정된 경우 백필은 예외를 throw하는 대신 백필 날짜 범위에 대해 실패한 모든 작업을 자동으로 다시 실행

3. execution Date

Airflow에서 execution_date는 DAG이 실행되어야 하는 기대값이라고 생각할 수 있습니다.

예를 들어 매일 새벽 5시에 전날 발생했던 로그들을 S3에 이관하는 작업이 있다고 가정해보겠습니다.
이때 실제로 2024-08-23 05:00:00에 실행된 DAG은 2024-08-22의 로그들을 가져옵니다.

보통 사용자는 오늘 날짜에서 -1일 해서 데이터를 가져올 수 있는데, Airflow에서는 execution_date를 사용하여 이를 대신할 수 있습니다.


자주 묻는 질문

처음 Airflow를 이용하는 분들이 종종 이런 질문을 합니다:

“그러면 그냥 오늘 날짜(datetime.today().date())에서 -1일 해서 사용해도 되지 않나요?”

하지만 오늘 날짜로 조건을 설정하고 이틀 뒤나 일주일 뒤에 에러를 발견하고 해당 DAG을 다시 실행하려고 한다면 문제가 발생할 수 있습니다.
이 경우 재실행하는 지금의 날짜가 조건으로 설정되기 때문에, 원하는 데이터를 가져오지 못할 가능성이 큽니다.


execution_date의 장점

execution_date를 사용하면, 재실행하더라도 execution_date처음 실행할 때 할당된 값에서 변경되지 않습니다.
따라서 backfill이나 DAG 재실행 시 매우 유용하게 사용할 수 있습니다.

# 현업에서 사용했던 execution_date 활용
from airflow.operators.python import get_current_context
from dateutil.relativedelta import relativedelta

def get_date():
    context = get_current_context()
    date = context["data_interval_start"] + relativedelta(hours=9) # UTC -> KST
    today = date.strftime("%Y-%m-%d")
    print(today)
    return today