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