FAQ
FAQ (자주 묻는 질문)
Airflow 사용 중 자주 나오는 질문과 답변을 모아둔 문서예요. 스케줄링·Dag 파일 파싱 문제, Dag 버전 무한 증가 원인과 해결, 스케줄 지연 개선, start_date·execution_date 이해, 태스크 실행 상호작용, UI·API 서버·MySQL 관련 문제, connection 테스트까지 광범위하게 다뤄요.
출처: 문서
본문
스케줄링 / Dag 파일 파싱 (Scheduling / Dag file parsing)
태스크가 왜 스케줄되지 않나요? (Why is task not getting scheduled?)
태스크가 스케줄되지 않을 이유는 매우 많아요. 다음은 흔한 원인들 중 일부예요:
-
스크립트가 "컴파일"되나요, Airflow 엔진이 파싱해서 Dag 객체를 찾을 수 있나요? 이를 테스트하려면
airflow dags list를 실행해 Dag가 목록에 나타나는지 확인할 수 있어요.airflow dags show foo_dag_id를 실행해 태스크가 예상대로 graphviz 형식으로 나타나는지도 확인할 수 있어요. CeleryExecutor를 사용한다면 스케줄러가 실행되는 곳과 worker가 실행되는 곳 양쪽에서 동작하는지 확인하고 싶을 거예요. -
Dag를 포함하는 파일에
airflow와DAG문자열이 내용 어딘가에 있나요? Dag 디렉토리를 검색할 때 Airflow는 사용자 Dag와 함께 있는 모든 python 파일을 DagBag 파싱이 import하지 못하게 하기 위해airflow와DAG를 포함하지 않는 파일은 무시해요. -
start_date가 제대로 설정되었나요? 시간 기반 Dags의 경우 태스크는 start date 다음의 첫 번째 스케줄 간격이 지나기 전까지 트리거되지 않아요. -
schedule인자가 제대로 설정되었나요? 기본값은 하루(datetime.timedelta(1))예요. 다른schedule은 인스턴스화하는 Dag 객체에 직접 지정해야 해요.default_param으로 지정하면 안 되는데, task instance는 부모 Dag의schedule을 오버라이드하지 않기 때문이에요. -
start_date가 UI에서 볼 수 있는 범위 밖인가요?start_date를 3개월 전 같은 시간으로 설정하면 UI의 메인 뷰에서는 볼 수 없지만,Menu -> Browse -> Task Instances에서는 볼 수 있어야 해요. -
태스크의 의존성이 충족되었나요? 태스크의 직접 업스트림인 task instance는
success상태여야 해요. 또한depends_on_past=True를 설정했다면 이전 task instance가 성공했거나 skip되었어야 해요(해당 태스크의 첫 실행이 아닌 경우). 그리고wait_for_downstream=True라면 그 의미를 확실히 이해하세요 — 이전 task instance의 직접 다운스트림 모든 태스크가 성공했거나 skip되었어야 해요. 이러한 속성이 어떻게 설정되어 있는지 태스크의Task Instance Details페이지에서 볼 수 있어요. -
필요한 DagRun이 생성되고 활성 상태인가요? DagRun은 전체 Dag의 특정 실행을 나타내며 상태(running, success, failed, ...)를 가져요. 스케줄러는 앞으로 나아갈 때 새 DagRun을 만들지만, 시간을 거슬러 올라가 만들어지지는 않아요. 스케줄러는
runningDagRun만 평가해 어떤 task instance를 트리거할 수 있는지 봐요. 참고로 task instance를 clear하면(UI나 CLI에서) DagRun의 상태가 다시 running으로 설정돼요. Dag의 schedule 태그를 클릭하면 DagRun 목록을 일괄로 보고 상태를 변경할 수 있어요. -
Dag의
concurrency파라미터에 도달했나요?concurrency는 Dag가 허용하는runningtask instance 수를 정의하며, 그 지점을 넘으면 대기열에 들어가요. -
Dag의
max_active_runs파라미터에 도달했나요?max_active_runs는 Dag의 동시running인스턴스가 허용되는 수를 정의해요.
Scheduler에 대해 읽고 스케줄러 사이클을 완전히 이해하는 것도 좋아요.
Dag 성능을 어떻게 향상시키나요? (How to improve Dag performance?)
더 큰 스케줄링 용량과 빈도를 허용하는 몇 가지 Airflow 구성이 있어요:
Dags에는 효율성을 개선하는 구성이 있어요:
max_active_tasks: max_active_tasks_per_dag를 오버라이드해요.max_active_runs: max_active_runs_per_dag를 오버라이드해요.
Operators나 태스크에도 효율성과 스케줄링 우선순위를 개선하는 구성이 있어요:
max_active_tis_per_dag: 태스크별로dag_runs에 걸친 동시 실행 task instance 수를 제어해요.pool: Pools을 참고하세요.priority_weight: Priority Weights를 참고하세요.queue: CeleryExecutor 배포 전용인 Queues를 참고하세요.
Dag 스케줄링 지연/태스크 지연을 어떻게 줄이나요? (How to reduce Dag scheduling latency / task delay?)
Airflow 2.0은 기본적으로 낮은 Dag 스케줄링 지연을 가져요(특히 Airflow 1.10.x와 비교할 때). 하지만 더 많은 처리량이 필요하다면 여러 스케줄러를 시작할 수 있어요.
다른 태스크의 실패에 기반해 태스크를 어떻게 트리거하나요? (How do I trigger tasks based on another task's failure?)
이는 Trigger Rules로 달성할 수 있어요.
Dag 파일마다 Dag 파일 파싱 타임아웃을 어떻게 제어하나요? (How to control Dag file parsing timeout for different Dag files?)
(Airflow >= 2.3.0에서만 유효)
airflow_local_settings.py에 Dag 파일이 파싱되기 직전에 호출되는 get_dagbag_import_timeout 함수를 추가해요. Dag 파일에 따라 다른 타임아웃 값을 반환할 수 있어요. 반환값이 0보다 작거나 같으면 Dag 파싱 중 타임아웃이 없다는 뜻이에요.
airflow_local_settings.py:
def get_dagbag_import_timeout(dag_file_path: str) -> Union[int, float]:
"""
This setting allows to dynamically control the Dag file parsing timeout.
It is useful when there are a few Dag files requiring longer parsing times, while others do not.
You can control them separately instead of having one value for all Dag files.
If the return value is less than or equal to 0, it means no timeout during the Dag parsing.
"""
if "slow" in dag_file_path:
return 90
if "no-timeout" in dag_file_path:
return 0
return conf.getfloat("core", "DAGBAG_IMPORT_TIMEOUT")
로컬 설정 구성 방법에 대한 자세한 내용은 Configuring local settings를 참고하세요.
Dag 파일이 많을 때(>1000) 새 파일 파싱을 어떻게 빠르게 하나요? (When there are a lot (>1000) of Dag files, how to speed up parsing of new files?)
file_parsing_sort_mode를 modified_time으로 변경하고, min_file_process_interval을 600(10분), 6000(100분) 또는 더 높은 값으로 올려요.
Dag 파서는 파일이 최근에 수정된 경우 min_file_process_interval 체크를 건너뛰어요.
이것은 Dag가 별도 파일에서 import·생성되는 경우에는 동작하지 않을 수 있어요. 예: 실제 Dag 로직이 아래처럼 되어 있는 dag_loader.py를 import하는 dag_file.py. 이 경우 dag_loader.py가 업데이트되었는데 dag_file.py가 업데이트되지 않았다면, Dag Parser가 dag_file.py의 수정 시간을 찾기 때문에 min_file_process_interval에 도달할 때까지 변경 사항이 반영되지 않아요.
dag_file.py:
from dag_loader import create_dag
globals()[dag.dag_id] = create_dag(dag_id, schedule, dag_number, default_args)
dag_loader.py:
from airflow.sdk import DAG
from airflow.sdk import task
import pendulum
def create_dag(dag_id, schedule, dag_number, default_args):
dag = DAG(
dag_id,
schedule=schedule,
default_args=default_args,
pendulum.datetime(2021, 9, 13, tz="UTC"),
)
with dag:
@task()
def hello_world():
print("Hello World")
print(f"This is Dag: {dag_number}")
hello_world()
return dag
UI에서 Dags가 사라지는 것을 보면 어떻게 하나요? (What to do if you see disappearing Dags on UI?)
UI에서 Dags가 사라질 수 있는 이유는 여러 가지가 있어요. 흔한 원인은 다음과 같아요:
-
모든 Dags의 총 파싱이 너무 오래 걸림 — 파싱이 dagbag_import_timeout보다 오래 걸리면 파일이 완전히 처리되지 않을 수 있어요. 이는 종종 Dags가 Dag 작성 모범 사례를 따르지 않을 때 발생해요. 예:
- 과도한 최상위 코드 실행
- 파싱 중 외부 시스템 호출
- 복잡한 동적 Dag 생성
-
일관되지 않은 동적 Dag 생성 — 동적 생성으로 만든 Dags는 파싱 간에 안정적인 Dag ID를 생성해야 해요.
python your_dag_file.py를 반복 실행해 일관성을 검증해요. -
파일 처리 구성 문제 — 특정 파라미터의 조합이 각 루프에서 특정 Dags가 처리될 가능성이 낮아지는 시나리오로 이어질 수 있어요. 다음 파라미터를 확인해요:
- file_parsing_sort_mode — 정렬 방식이 동기화 전략과 일치하는지 확인
- parsing_processes — 병렬 파서 수
- parsing_cleanup_interval — stale Dag 정리 빈도 제어
-
파일 동기화 문제 — git-sync 설정에서 흔함:
- 심볼릭 링크 스왑 지연
- 동기화 중 권한 변경
mtime보존 문제
-
시간 동기화 문제 — 모든 노드(데이터베이스, 스케줄러, 워커)가 NTP를 사용하고 1초 미만의 클럭 드리프트를 갖도록 보장해요.
Dag 버전이 계속 증가하는 이유는 무엇인가요? (Why does my Dag version keep increasing?)
Dag processor가 Dag 파일을 파싱할 때마다 Dag를 직렬화하고 그 결과를 메타데이터 데이터베이스에 저장된 버전과 비교해요. 무엇이든 변경되면 Airflow는 새 Dag 버전을 만들어요.
Dag 버전 무한 증가(version inflation) 는 Dag 작성자가 의도적인 변경을 하지 않았는데 버전 번호가 무한정 증가할 때 발생해요.
무엇이 잘못되나요 (What goes wrong)
의미 있는 변경 없이 Dag 버전이 증가하면:
- 메타데이터 데이터베이스에 불필요한 Dag 버전 레코드가 쌓여 저장·쿼리 오버헤드가 증가해요.
- UI가 오해의 소지가 있는 Dag 변경 이력을 보여줘 실제 수정 사항을 식별하기 어려워져요.
- 스케줄러와 API 서버가 증가하는 Dag 버전을 로드·캐시하면서 더 많은 메모리를 소비할 수 있어요.
흔한 원인 (Common causes)
버전 무한 증가는 파싱 시간(즉 Dag processor가 Dag 파일을 평가할 때마다)에 변하는 값을 Dag나 Task 생성자 인자로 사용해서 발생해요. 가장 흔한 패턴은 다음과 같아요:
1. datetime.now()나 pendulum.now()를 start_date로 사용:
from datetime import datetime
from airflow.sdk import DAG
with DAG(
dag_id="bad_example",
# BAD: datetime.now() produces a different value on every parse
start_date=datetime.now(),
schedule="@daily",
):
...
매번 파싱마다 다른 start_date가 생성되므로, 직렬화된 Dag는 항상 저장된 버전과 달라요.
2. Dag나 Task 인자에 무작위 값 사용:
import random
from airflow.sdk import DAG
from airflow.providers.standard.operators.python import PythonOperator
with DAG(dag_id="bad_random", start_date="2024-01-01", schedule="@daily") as dag:
PythonOperator(
# BAD: random value changes every parse
task_id=f"task_{random.randint(1, 1000)}",
python_callable=lambda: None,
)
3. 런타임에 변하는 값을 생성자에 사용되는 변수에 할당:
from datetime import datetime
from airflow.sdk import DAG
from airflow.providers.standard.operators.python import PythonOperator
# BAD: the variable captures a parse-time value, then is passed to the DAG
default_args = {"start_date": datetime.now()}
with DAG(dag_id="bad_defaults", default_args=default_args, schedule="@daily") as dag:
PythonOperator(task_id="my_task", python_callable=lambda: None)
datetime.now()가 Dag 생성자 안에서 직접 호출되지 않더라도 default_args를 통해 들어와 여전히 매번 파싱마다 다른 직렬화 Dag를 만들 수 있어요.
4. 파싱 사이에 변하는 환경 변수나 파일 내용 사용:
import os
from airflow.sdk import DAG
from airflow.providers.standard.operators.bash import BashOperator
with DAG(dag_id="bad_env", start_date="2024-01-01", schedule="@daily") as dag:
BashOperator(
task_id="echo_build",
# BAD if BUILD_NUMBER changes on every deployment or parse
bash_command=f"echo {os.environ.get('BUILD_NUMBER', 'unknown')}",
)
버전 무한 증가를 피하는 방법 (How to avoid version inflation)
- 고정된
start_date값을 사용하세요. 항상start_date를 정적datetime리터럴로 설정하세요:
import datetime
from airflow.sdk import DAG
with DAG(
dag_id="good_example",
start_date=datetime.datetime(2024, 1, 1),
schedule="@daily",
):
...
-
모든 Dag·Task 생성자 인자를 결정적(deterministic)으로 유지하세요. Dag와 Operator 생성자에 전달되는 인자는 매번 파싱마다 같은 값을 만들어야 해요. 동적 계산은
execute()메서드로 옮기거나 Jinja 템플릿(파싱 시간이 아니라 태스크 실행 시간에 평가됨)을 사용하세요. -
동적 값에는 Jinja 템플릿을 사용하세요:
from airflow.providers.standard.operators.bash import BashOperator
BashOperator(
task_id="echo_date",
# GOOD: the template is resolved at execution time, not parse time
bash_command="echo {{ ds }}",
)
- 최상위 조회 대신 템플릿과 함께 Airflow Variables를 사용하세요:
from airflow.providers.standard.operators.bash import BashOperator
BashOperator(
task_id="echo_var",
# GOOD: Variable is resolved at execution time via template
bash_command="echo {{ var.value.my_variable }}",
)
버전 무한 증가를 직접 진단하는 방법 (How to diagnose version inflation yourself)
Dag 버전이 계속 증가하는데 어떤 생성자 인자가 원인인지 명확하지 않다면, 직렬화된 Dag 행을 직접 검사해 변하는 필드를 식별할 수 있어요.
가장 확실한 방법은 메타데이터 데이터베이스의 serialized_dag 테이블을 쿼리하고, 같은 Dag의 두 연속 버전 사이에 직렬화 페이로드가 어떻게 다른지 비교하는 것이에요. 직렬화된 Dag를 저장하는 컬럼은 Airflow가 새 버전이 필요한지 결정하기 위해 해시를 적용하는 구조화된 데이터를 담고 있으므로, 두 행을 diff하면 어떤 필드가 파싱 사이에 변했는지 정확히 보여줘요.
변하는 필드를 발견하면 Dag 파일에서 그것이 설정된 위치로 거슬러 올라가, 동적 값을 결정적인 값(예: 고정 datetime 리터럴, 실행 시간에 평가되는 Jinja 템플릿, 템플릿으로 해석되는 Airflow Variable)으로 바꾸세요. 구체적인 패턴은 '버전 무한 증가를 피하는 방법'을 참고하세요.
변경 필드를 Dag 코드의 어떤 것과 연결할 수 없고 값이 Airflow 자체에서 오는 것처럼 보인다면(예: 파싱 간에 안정적이어야 하는 내부 속성), https://github.com/apache/airflow/issues에 버그로 보고해 주세요. 두 직렬화된 Dag 페이로드(또는 관련 diff)와 Airflow 버전을 포함해 메인테이너가 문제를 재현할 수 있게 해 주세요.
Dag 버전 무한 증가 감지 (Dag version inflation detection)
Airflow 3.2부터 Dag processor는 Dag 파일을 파싱하기 전에 AST 기반 정적 분석을 수행해 Dag·Task 생성자의 런타임 변동 값을 감지해요. 잠재적 문제가 발견되면 UI에 보이는 Dag 경고(warning) 로 표시돼요.
이 동작은 dag_version_inflation_check_level 구성 옵션으로 제어할 수 있어요:
off— 체크를 완전히 비활성화해요. 오류·경고가 생성되지 않아요.warning(기본값) — Dags는 정상적으로 로드되지만 문제가 감지되면 UI에 경고가 표시돼요.error— 감지된 문제를 Dag import 오류로 취급해 Dag가 로드되지 않게 해요.
추가로 개발 워크플로우에서 정적 린팅의 일부로 Dag·Task 생성자의 동적 값을 감지하는 AIR302 ruff 규칙을 사용해 이런 문제를 더 일찍 잡을 수 있어요. Airflow 특화 규칙으로 ruff를 설정하는 방법은 Code Quality and Linting을 참고하세요.
Dag 구성 (Dag construction)
start_date의 정체는 무엇인가요? (What's the deal with start_date?)
start_date는 부분적으로 pre-DagRun 시대의 유산이지만, 여전히 여러 면에서 관련이 있어요. 새 Dag를 만들 때 태스크에 전역 start_date를 설정하고 싶을 거예요. 이는 DAG() 객체에 start_date를 직접 선언해 할 수 있어요. Dag의 첫 DagRun은 start_date 다음의 첫 번째 완전한 data_interval에 기반해 생성돼요. 예를 들어 start_date=datetime(2024, 1, 1)과 schedule="0 0 3 * *"인 Dag의 첫 Dag run은 data_interval_start=datetime(2024, 1, 3)·data_interval_end=datetime(2024, 2, 3)으로 2024-02-03 자정에 트리거돼요. 그 시점부터 스케줄러는 schedule에 기반해 새 DagRuns를 만들고 의존성이 충족되면서 해당 task instance들이 실행돼요. Dag에 새 태스크를 도입할 때는 start_date에 특별히 주의해야 하고, 새 태스크를 제대로 시작시키기 위해 비활성 DagRuns를 재활성화하고 싶을 수도 있어요.
start_date로 동적 값을 사용하는 것, 특히 datetime.now()는 매우 혼란스러울 수 있어 권장하지 않아요. 태스크는 기간이 닫힌 후 트리거되고, 이론상 now()가 계속 이동하므로 @hourly Dag는 지금으로부터 한 시간 뒤에 절대 도달하지 못할 거예요.
이전에는 Dag의 schedule과 관련해 반올림된 start_date를 사용하는 것도 권장했어요. 이는 @hourly는 00:00 분:초에, @daily job은 자정에, @monthly job은 매월 1일에 하는 것을 의미했어요. 이제는 더 이상 필요하지 않아요. Airflow는 이제 start_date를 검색을 시작하는 순간으로 사용해 start_date와 schedule을 자동으로 정렬해요.
스케줄 간격 내에서 태스크 실행을 지연시키기 위해 어떤 sensor나 TimeDeltaSensor를 사용할 수 있어요. schedule이 datetime.timedelta 객체를 지정하는 것을 허용하지만, 반올림된 스케줄의 아이디어를 강제하므로 매크로나 cron 표현식을 권장해요.
depends_on_past=True를 사용할 때는 start_date에 특별히 주의하는 것이 중요해요. 과거 의존성은 태스크에 지정된 start_date의 특정 스케줄에만 강제되는 것이 아니기 때문이에요. 새 depends_on_past=True를 도입할 때는 새 태스크(들)에 대해 backfill을 실행할 계획이 아니라면 시간에 따른 DagRun 활동 상태를 관찰하는 것도 중요해요.
또한 태스크의 start_date는 backfill에서 무시된다는 점도 중요해요.
시간대 사용하기 (Using time zones)
시간대 인식 datetime(예: Dag의 start_date)을 만드는 것은 매우 간단해요. pendulum을 사용해 시간대 인식 날짜를 제공하기만 하면 돼요. 표준 라이브러리 timezone은 알려진 제한이 있고 Dags에서 사용을 의도적으로 허용하지 않으므로 사용하려 하지 마세요.
execution_date는 무엇을 의미하나요? (What does execution_date mean?)
Execution date 또는 execution_date는 logical date라고 불리는 것의 역사적 이름이고, 보통 Dag run이 나타내는 데이터 구간의 시작이기도 해요.
Airflow는 ETL 요구를 위한 해결책으로 개발됐어요. ETL 세계에서는 보통 데이터를 요약해요. 그래서 2016-02-19의 데이터를 요약하고 싶다면 2016-02-19의 모든 데이터가 사용 가능해진 직후인 2016-02-20 자정 UTC에 하게 돼요. 2016-02-19와 2016-02-20 자정 사이의 이 간격을 data interval이라고 하고, 그것이 2016-02-19 날짜의 데이터를 나타내므로 이 날짜는 run의 logical date, 즉 이 Dag run이 실행되는 날짜, 따라서 execution date라고도 불려요.
하위 호환을 위해 execution_date datetime 값은 여전히 Jinja 템플릿 필드에서 다양한 형식의 Template variables로, 그리고 Airflow의 Python API에서 제공돼요. Operator의 execute 함수에 주어지는 context 딕셔너리에도 포함돼요.
class MyOperator(BaseOperator):
def execute(self, context):
logging.info(context["execution_date"])
다만 가능하면 항상 data_interval_start나 data_interval_end를 사용해야 해요. 그 이름들이 의미상 더 정확하고 오해의 소지가 적기 때문이에요.
ds(즉 data_interval_start의 YYYY-MM-DD 형태)는 일부가 혼동할 수 있는 date start가 아니라 date string을 의미한다는 점을 참고하세요.
팁 (Tip)
logical date에 대한 자세한 내용은 Data Interval과 Running Dags를 참고하세요.
Dags를 동적으로 어떻게 만드나요? (How to create Dags dynamically?)
Airflow는 DAGS_FOLDER에서 전역 네임스페이스에 DAG 객체를 포함하는 모듈을 찾고, 찾은 객체를 DagBag에 추가해요. 이를 알면 전역 네임스페이스에 변수를 동적으로 할당할 방법만 있으면 돼요. 이는 표준 라이브러리의 globals() 함수를 사용해 쉽게 할 수 있는데, 이 함수는 간단한 딕셔너리처럼 동작해요.
def create_dag(dag_id):
"""
A function returning a DAG object.
"""
return DAG(dag_id)
for i in range(10):
dag_id = f"foo_{i}"
globals()[dag_id] = DAG(dag_id)
# or better, call a function that returns a DAG object!
other_dag_id = f"bar_{i}"
globals()[other_dag_id] = create_dag(other_dag_id)
Airflow는 python 파일당 여러 Dag 정의(동적 생성이든 아니든)를 지원하지만, 권장하지는 않아요. Airflow는 장애와 배포 관점에서 Dags 간의 더 나은 격리를 원하며, 같은 파일의 여러 Dags는 그에 반하기 때문이에요.
최상위 Python 코드가 허용되나요? (Are top level Python code allowed?)
Airflow 구성을 정의하는 것 외에 어떤 코드도 쓰는 것이 권장되지는 않지만, Airflow는 Dag 파일 processor를 깨뜨리거나 파일 처리 시간을 dagbag_import_timeout 값보다 늘리지 않는 한 임의의 python 코드를 지원해요.
흔한 예는 보통 데이터베이스 같은 다른 서비스에서 데이터를 쿼리해야 하는 동적 Dag를 만들 때 시간 제한을 위반하는 경우예요. 동시에 요청된 서비스는 파일을 처리하기 위한 데이터를 요청하는 Dag 파일 processor들의 요청에 압도당하고 있어요. 이런 의도하지 않은 상호작용은 서비스가 저하되고 결국 Dag 파일 처리가 실패하게 할 수 있어요.
자세한 내용은 Dag 작성 모범 사례를 참고하세요.
매크로가 다른 Jinja 템플릿에서 해석되나요? (Do Macros resolves in another Jinja template?)
Macros나 어떤 Jinja 템플릿도 다른 Jinja 템플릿 안에서 렌더링하는 것은 불가능해요. 이는 흔히 user_defined_macros에서 시도돼요.
dag = DAG(
# ...
user_defined_macros={"my_custom_macro": "day={{ ds }}"}
)
bo = BashOperator(task_id="my_task", bash_command="echo {{ my_custom_macro }}", dag=dag)
이것은 data_interval_start가 2020-01-01 00:00:00인 Dag run에 대해 "day=2020-01-01" 대신 "day={{ ds }}"를 echo해요.
bo = BashOperator(task_id="my_task", bash_command="echo day={{ ds }}", dag=dag)
ds 매크로를 template_field에서 직접 사용하면 렌더링된 값이 "day=2020-01-01"이 돼요.
next_ds나 prev_ds가 예상 값을 포함하지 않는 이유는 무엇인가요? (Why next_ds or prev_ds might not contain expected values?)
-
Dag를 스케줄할 때
next_ds·next_ds_nodash·prev_ds·prev_ds_nodash는logical_date와 Dag의 스케줄(해당하는 경우)을 사용해 계산돼요.schedule을None이나@once로 설정하면next_ds,next_ds_nodash,prev_ds,prev_ds_nodash값이None으로 설정돼요. -
Dag를 수동으로 트리거할 때는 스케줄이 무시되고
prev_ds == next_ds == ds가 돼요.
태스크 실행 상호작용 (Task execution interactions)
TemplateNotFound는 무엇을 의미하나요? (What does TemplateNotFound mean?)
TemplateNotFound 오류는 보통 Jinja 템플릿을 트리거하는 operator에 경로를 전달할 때 사용자 기대와 어긋나서 발생해요. BashOperator에서 흔히 발생해요.
또 하나 흔히 놓치는 사실은 파일이 파이프라인 파일이 있는 위치 기준으로 해석된다는 것이에요. 다른 비상대적 위치를 허용하려면 Dag 객체의 template_searchpath에 다른 디렉토리를 추가할 수 있어요.
다른 태스크의 실패에 기반해 태스크를 어떻게 트리거하나요? (How to trigger tasks based on another task's failure?)
의존성으로 연결된 태스크의 경우, 태스크 실행이 모든 업스트림 태스크의 실패에 의존한다면 trigger_rule을 TriggerRule.ALL_FAILED로, 업스트림 태스크 중 하나만이라면 TriggerRule.ONE_FAILED로 설정할 수 있어요.
import pendulum
from airflow.sdk import dag, task
from airflow.exceptions import AirflowException
from airflow.utils.trigger_rule import TriggerRule
@task()
def a_func():
raise AirflowException
@task(
trigger_rule=TriggerRule.ALL_FAILED,
)
def b_func():
pass
@dag(schedule="@once", start_date=pendulum.datetime(2021, 1, 1, tz="UTC"))
def my_dag():
a = a_func()
b = b_func()
a >> b
dag = my_dag()
자세한 내용은 Trigger Rules를 참고하세요.
태스크가 의존성으로 연결되지 않았다면 커스텀 Operator를 빌드해야 해요.
Airflow UI
태스크가 UI에 로그 없이 실패한 이유는 무엇인가요? (Why did my task fail with no logs in the UI?)
로그는 태스크가 터미널 상태에 도달할 때 일반적으로 제공됩니다. 때로 태스크의 정상 라이프사이클이 중단되어 태스크의 worker가 태스크 로그를 작성하지 못할 수 있어요. 이는 보통 두 가지 이유 중 하나로 발생해요:
- Task Instance Heartbeat Timeout.
- queued에 갇힌 후 실패한 태스크(Airflow 2.6.0+). scheduler.task_queued_timeout보다 오래 queued 상태인 태스크는 실패로 표시되고, Airflow UI에 태스크 로그가 없어요.
각 태스크에 retries를 설정하면 이 두 문제 중 하나가 워크플로우에 영향을 줄 가능성을 크게 줄여요.
webserver별로 sync perms가 여러 번 발생하는 것을 어떻게 막나요? (How do I stop the sync perms happening multiple times per webserver?)
airflow.cfg의 [fab] update_fab_perms 구성 값을 False로 설정하세요.
pause Dag 토글이 왜 빨간색이 되었나요? (Why did the pause Dag toggle turn red?)
어떤 이유로든 Dag를 pause하거나 unpause하는 데 실패하면 Dag 토글은 이전 상태로 되돌아가고 빨간색이 돼요. 이 동작을 보면 Dag를 다시 pause해 보거나, 문제가 반복되면 콘솔이나 서버 로그를 확인해 보세요.
API 서버 (API Server)
API 서버 메모리 증가를 어떻게 막나요? (How to prevent API server memory growth?)
API 서버는 직렬화된 Dag 객체를 메모리에 캐시해요. 시간이 지나며 Dag 버전이 쌓이면(위 'Dag 버전이 계속 증가하는 이유' 참고) 이 캐시가 커져 수 기가바이트의 메모리를 소비할 수 있어요.
서로 보완적인 두 가지 접근 방식이 있어요:
1. Dag 캐시 축출 (Dag cache eviction, Airflow 3.2.2부터 사용 가능)
API 서버는 크기별, 연령별, 또는 둘 다로 캐시된 직렬화 Dag 버전을 축출할 수 있어요. [api] 섹션에서 구성해요:
[api]
dag_cache_size = 64 ; max cached versions (0 = no size limit)
dag_cache_ttl = 3600 ; seconds before a cached entry expires (0 = no TTL)
dag_cache_size는 메모리의 유일한 하드 상한이에요. 엔트리의 TTL은 매 요청마다가 아니라 [core] min_serialized_dag_update_interval 이후 엔트리가 데이터베이스에 대해 체크될 때만 갱신돼요. TTL이 더 짧으면 자주 요청되는 엔트리도 체크 사이에 만료·다시 로드될 수 있어요. 두 옵션을 모두 0으로 설정하면 축출 없는 무한 dict를 사용하며, 이는 3.2.2 이전의 동작과 일치해요.
캐시는 Dag 버전 ID로 키가 정해져요. Dag가 업데이트된 후 API 서버는 캐시된 엔트리가 만료될 때까지(dag_cache_ttl로 제어) 이전 버전을 제공할 수 있어요.
전체 구성 참조는 dag_cache_size와 dag_cache_ttl을 참고하세요.
2. 롤링 워커 재시작이 있는 Gunicorn (Airflow 3.2.0부터 사용 가능)
Gunicorn은 주기적으로 worker 프로세스를 재활용해 누적된 모든 메모리를 해제해요. 또한 preload + fork를 사용하므로 worker가 copy-on-write를 통해 읽기 전용 메모리 페이지를 공유해, uvicorn의 multiprocess 모드와 비교해 전체 메모리 사용을 40-50% 줄여요.
worker 재활용과 함께 gunicorn을 활성화하려면:
[api]
server_type = gunicorn
# Restart each worker every 12 hours (43200 seconds)
worker_refresh_interval = 43200
worker_refresh_batch_size = 1
이를 위해서는 apache-airflow-core[gunicorn] extra가 설치되어 있어야 해요.
전체 구성 참조는 server_type, worker_refresh_interval, worker_refresh_batch_size를 참고하세요.
참고 (Note)
워커 재활용은 Dag 캐시뿐 아니라 어떤 원천의 메모리 증가도 처리해요. 프로덕션 배포에서는 캐시 축출과 gunicorn 워커 재활용을 모두 사용하는 것이 가장 좋은 결과를 제공해요.
MySQL 및 MySQL 변형 데이터베이스 (MySQL and MySQL variant Databases)
"MySQL Server has gone away"는 무엇을 의미하나요? (What does "MySQL Server has gone away" mean?)
때때로 "MySQL Server has gone away" 메시지와 함께 OperationalError가 발생할 수 있어요. 이는 connection pool이 연결을 너무 오래 열어두어 만료된 이전 연결을 받기 때문이에요. 유효한 연결을 보장하려면 sql_alchemy_pool_recycle을 설정해 그 초가 지나면 연결이 무효화되고 새 연결이 생성되게 해요.
Airflow가 확장 ASCII나 유니코드 문자를 지원하나요? (Does Airflow support extended ASCII or unicode characters?)
Airflow에서 확장 ASCII나 Unicode 문자를 사용하려면 MySQL 데이터베이스에 적절한 connection 문자열을 제공해야 해요. 이는 charset을 명시적으로 정의하기 때문이에요.
sql_alchemy_conn = mysql://airflow@localhost:3306/airflow?charset=utf8
WTForms 템플릿과 다른 Airflow 모듈이 아래처럼 UnicodeDecodeError를 발생시키는 것을 보게 될 거예요.
'ascii' codec can't decode byte 0xae in position 506: ordinal not in range(128)
Exception: Global variable explicit_defaults_for_timestamp needs to be on (1)?를 어떻게 고치나요?
이는 mysql 서버에서 explicit_defaults_for_timestamp가 비활성화되어 있다는 뜻이며 다음으로 활성화해야 해요:
my.cnf파일의mysqld섹션 아래에explicit_defaults_for_timestamp = 1을 설정하세요.- Mysql 서버를 재시작하세요.
Connections
connection을 어떻게 테스트하거나 Canary Dag를 사용하나요? (How can I test a connection or use a Canary Dag?)
보안상의 이유로 test connection 기능은 Airflow UI, API, CLI 전반에서 기본적으로 비활성화돼 있어요. 이는 config:core__test_connection을 설정해 수정할 수 있어요.
Dag를 활용해 connections를 정기적으로 테스트할 수 있어요. 이것을 "Canary Dag"라고 하며, Dags가 의존하는 외부 시스템의 실패를 감지하고 알릴 수 있어요. 다음 Airflow 3 예시처럼 connections를 테스트하는 간단한 Dag를 만들 수 있어요:
from airflow import DAG
from airflow.sdk import task
with DAG(dag_id="canary", schedule="@daily", doc_md="Canary Dag to regularly test connections to systems."):
@task(doc_md="Test a connection by its Connection ID.")
def test_connection(conn_id):
from airflow.hooks.base import BaseHook
ok, status = BaseHook.get_hook(conn_id=conn_id).test_connection()
if ok:
return status
raise RuntimeError(status)
for conn_id in [
# Add more connections here to create tasks to test them.
"aws_default",
]:
test_connection.override(task_id="test_" + conn_id)(conn_id)