동적 Dag 생성
동적 Dag 생성 (Dynamic Dag Generation)
구조가 동적으로 생성되는 Dags를 만드는 방법을 설명하는 문서예요. 환경 변수를 이용한 동적 구성, 메타데이터가 포함된 Python 코드 생성, 구조화된 데이터 파일의 외부 설정, @dag 데코레이터를 통한 자동 등록, 그리고 파싱 지연 최적화까지 다양한 기법을 살펴볼게요.
출처: 문서
본문
이 문서는 구조가 동적으로 생성되지만 Dag 실행 간에 태스크 수가 변하지 않는 Dags의 생성을 설명해요. 업스트림 태스크의 출력/결과에 따라 태스크(또는 Airflow 2.6부터는 Task Group) 수가 바뀔 수 있는 Dag를 구현하고 싶다면, Dynamic Task Mapping을 참고하세요.
참고 (Note)
태스크와 태스크 그룹의 일관된 생성 순서
Dags를 동적으로 생성하는 모든 경우, Task와 Task Group이 Dag가 생성될 때마다 일관된 순서로 생성되도록 해야 해요. 그렇지 않으면 Grid View에서 페이지를 새로고침할 때마다 Task와 Task Group의 순서가 바뀔 수 있어요. 예를 들어 데이터베이스 쿼리에서 안정적인 정렬 메커니즘을 사용하거나 Python의
sorted()함수를 사용해 달성할 수 있어요.
환경 변수를 이용한 동적 Dags (Dynamic Dags with environment variables)
코드를 구성하기 위해 변수를 사용하고 싶다면, 최상위 레벨 코드에서는 Airflow Variables보다 환경 변수를 항상 사용해야 해요. 최상위 코드에서 Airflow Variables를 사용하면 값을 가져오기 위해 Airflow의 metastore DB에 연결을 만들어 파싱을 느리게 하고 DB에 추가 부하를 줄 수 있어요. Jinja 템플릿을 사용해 Dag에서 Airflow Variables를 가장 잘 활용하는 방법은 Airflow Variables 모범 사례를 참고하세요.
예를 들어 프로덕션과 개발 환경에 대해 DEPLOYMENT 변수를 다르게 설정할 수 있어요. DEPLOYMENT 변수는 프로덕션 환경에서는 PROD로, 개발 환경에서는 DEV로 설정할 수 있어요. 그러면 환경 변수의 값에 따라 프로덕션과 개발 환경에서 Dag를 다르게 구성할 수 있어요.
deployment = os.environ.get("DEPLOYMENT", "PROD")
if deployment == "PROD":
task = Operator(param="prod-param")
elif deployment == "DEV":
task = Operator(param="dev-param")
임베디드 메타데이터가 있는 Python 코드 생성하기 (Generating Python code with embedded meta-data)
메타데이터를 import 가능한 상수로 포함한 Python 코드를 외부에서 생성할 수 있어요. 그런 상수는 Dag가 직접 import해서 객체를 구성하고 의존성을 구축하는 데 사용할 수 있어요. 이렇게 하면 상수에 저장된 메타데이터를 찾고·로드하고·파싱할 필요 없이 여러 Dag에서 그 코드를 쉽게 import할 수 있어요 — Python 인터프리터가 "import" 문을 처리할 때 자동으로 수행해 줘요. 처음에는 이상하게 들리지만, 그런 코드를 생성하고 Dag에서 import할 수 있는 유효한 Python 코드인지 확인하는 것은 놀라울 만큼 쉬워요.
예를 들어 (Dag 폴더에서) my_company_utils/common.py 파일을 동적으로 생성한다고 가정해 보세요:
# This file is generated automatically !
ALL_TASKS = ["task1", "task2", "task3"]
그러면 모든 Dag에서 ALL_TASKS 상수를 이렇게 import해서 사용할 수 있어요:
from my_company_utils.common import ALL_TASKS
with DAG(
dag_id="my_dag",
schedule=None,
start_date=datetime(2021, 1, 1),
catchup=False,
):
for task in ALL_TASKS:
# create your operators and relations here
...
이 경우 my_company_utils 폴더에 빈 __init__.py 파일을 추가해야 하고, .airflowignore 파일에 my_company_utils/* 줄을 추가해(기본 glob 문법 사용) 스케줄러가 Dags를 찾을 때 전체 폴더를 무시하도록 해야 한다는 것을 잊지 마세요.
구조화된 데이터 파일의 외부 설정을 사용한 동적 Dags (Dynamic Dags with external configuration from a structured data file)
Dag 구조를 준비하기 위해 더 복잡한 메타데이터를 사용해야 하고 데이터를 구조화된 비-Python 형식으로 유지하고 싶다면, 부모 문서인 Top level Python Code에서 설명한 이유 때문에 Dag의 최상위 코드로 데이터를 끌어오려 하기보다, 데이터를 파일로 export해서 Dag 폴더에 넣어야 해요.
메타데이터는 Dags와 함께 Dag 폴더의 편리한 파일 형식(JSON, YAML 형식이 좋은 후보)으로 export·저장해야 해요. 이상적으로는 메타데이터를 불러오는 Dag 파일의 모듈과 같은 패키지/폴더에 게시해야 해요. 그러면 Dag에서 메타데이터 파일의 위치를 쉽게 찾을 수 있어요. 읽을 파일의 위치는 Dag를 포함하는 모듈의 __file__ 속성으로 찾을 수 있어요:
my_dir = os.path.dirname(os.path.abspath(__file__))
configuration_file_path = os.path.join(my_dir, "config.yaml")
with open(configuration_file_path) as yaml_file:
configuration = yaml.safe_load(yaml_file)
# Configuration dict is available here
동적 Dags 등록하기 (Registering dynamic Dags)
@dag 데코레이터나 with DAG(..) 컨텍스트 매니저를 사용할 때 Dags를 동적으로 생성할 수 있고, Airflow가 자동으로 등록해 줘요.
from datetime import datetime
from airflow.sdk import dag, task
configs = {
"config1": {"message": "first Dag will receive this message"},
"config2": {"message": "second Dag will receive this message"},
}
for config_name, config in configs.items():
dag_id = f"dynamic_generated_dag_{config_name}"
@dag(dag_id=dag_id, start_date=datetime(2022, 2, 1))
def dynamic_generated_dag():
@task
def print_message(message):
print(message)
print_message(config["message"])
dynamic_generated_dag()
위 코드는 각 config에 대해 Dag를 생성해요: dynamic_generated_dag_config1과 dynamic_generated_dag_config2. 각각은 관련 설정과 함께 별도로 실행될 수 있어요.
Dag가 자동 등록되지 않기를 원한다면, Dag에 auto_register=False를 설정해 그 동작을 비활성화할 수 있어요.
2.4 버전에서 변경: 2.4 버전부터
@dag데코레이터가 붙은 함수를 호출해서(또는with DAG(...)컨텍스트 매니저에서 사용해서) 만든 Dags는 자동으로 등록되며, 더 이상 전역 변수에 저장할 필요가 없어요.
실행 중 Dag 파싱 지연 최적화하기 (Optimizing Dag parsing delays during execution)
2.4 버전에 추가됨 (Added in version 2.4).
때로 단일 Dag 파일에서 많은 동적 Dags를 생성하면, 태스크 실행 중 Dag 파일이 파싱될 때 불필요한 지연이 발생할 수 있어요. 그 영향은 태스크가 시작하기 전의 지연으로 나타나요.
왜 이런 일이 발생할까요? 인지하지 못했을 수 있는데, 태스크가 실행되기 직전에 Airflow는 Dag가 나온 Python 파일을 파싱해요.
Airflow Dag File Processor는 모든 메타데이터를 처리하려면 완전한 Dag 파일을 로드해야 해요. 하지만 태스크 실행은 태스크를 실행하기 위한 단일 Dag 객체만 필요해요. 이를 알면 태스크가 실행될 때 불필요한 Dag 객체 생성을 건너뛰어 파싱 시간을 줄일 수 있어요. 이 최적화는 생성된 Dag 수가 많을 때 가장 효과적이에요.
항상 사용할 수 있는 것은 아니라는 점을 참고하세요(예: 후속 Dag 생성이 이전 Dag에 의존할 때) 또는 Dag 생성에 부작용이 있을 때. 이 해결책을 주의해서 사용하고 철저히 테스트하세요.
얻을 수 있는 성능 개선의 좋은 예는 태스크 실행 중 파싱을 120초에서 200ms로 줄인 방법을 설명한 Airflow's Magic Loop 블로그 포스트에서 볼 수 있어요. (그 예시는 Airflow 2.4 이전에 작성되어 아래 get_parsing_context()가 아닌 Airflow의 문서화되지 않은 동작을 사용해요.)
Airflow 2.4+에서는 대신 get_parsing_context() 메서드를 사용해 문서화되고 예측 가능한 방식으로 현재 컨텍스트를 검색할 수 있어요.
Dag를 생성할 대상 컬렉션을 반복할 때, 컨텍스트를 사용해 모든 Dag 객체를 생성해야 하는지(Dag File processor에서 파싱할 때) 아니면 단일 Dag 객체만 생성해야 하는지(태스크를 실행할 때) 판단할 수 있어요.
get_parsing_context()는 현재 파싱 컨텍스트를 반환해요. 컨텍스트는 AirflowParsingContext이며, 단일 Dag/task만 필요할 때는 dag_id와 task_id 필드가 설정되어 있어요. "전체(full)" 파싱이 필요할 때(예: Dag File Processor에서)는 컨텍스트의 dag_id와 task_id가 None으로 설정돼요.
from airflow.sdk import DAG
from airflow.sdk import get_parsing_context
current_dag_id = get_parsing_context().dag_id
for thing in list_of_things:
dag_id = f"generated_dag_{thing}"
if current_dag_id is not None and current_dag_id != dag_id:
continue # skip generation of non-selected Dag
with DAG(dag_id=dag_id, ...):
...