클러스터 정책

클러스터 정책 (Cluster Policies)

이 페이지는 클러스터 전체 차원에서 DAG나 Task를 검사·변경하는 클러스터 정책(Cluster Policy)을 다뤄요. dag_policy, task_policy, task_instance_mutation_hook 세 가지 주요 타입이 있으며, airflow_local_settings.py 파일이나 setuptools entrypoint로 정책 함수를 정의할 수 있어요. 정책을 통해 DAG·Task 표준 확인, 기본 인자 설정, 커스텀 라우팅 로직을 수행할 수 있어요.

출처: 문서

본문

DAG나 Task를 클러스터 전체 차원에서 확인하거나 변경하고 싶다면, 클러스터 정책(Cluster Policy)으로 할 수 있어요. 또한 dag_id나 다른 속성에 기반해 DAG에 클러스터 전체 설정을 적용할 수도 있어요.

일반적인 사용 사례는 다음과 같아요:

  • DAG/Task가 일정 표준을 충족하는지 확인
  • DAG/Task에 기본 인자 설정하기
  • 커스텀 라우팅 로직 수행

클러스터 정책에는 세 가지 주요 타입이 있어요:

  • dag_policy: dag라는 DAG 파라미터를 받아요. DagBag의 DAG가 로드될 때 실행돼요.
  • task_policy: task라는 BaseOperator 파라미터를 받아요. 이 정책은 로드 시 DagBag에서 Task를 파싱할 때 Task가 생성되면 실행돼요. 즉 task definition 전체를 task policy에서 변경할 수 있어요. 특정 DagRun에서 실행되는 특정 Task와는 관련이 없어요. 정의된 task_policy는 미래에 실행될 모든 task instance에 적용돼요.
  • task_instance_mutation_hook: task_instance라는 TaskInstance 파라미터와 선택적 dag_run이라는 DagRun 파라미터를 받아요. task_instance_mutation_hook은 Task가 아니라 특정 DagRun과 관련된 Task의 인스턴스에 적용돼요. task instance가 생성되거나 조정(reconcile)될 때 scheduler 쪽에서 실행돼요( Dag file processor나 worker에서가 아니라). 정책은 그 Task의 현재 실행되는 run(인스턴스)에만 적용돼요. dag_run 인자는 정책이 run 설정(dag_run.conf)에 라우팅할 수 있게 해요. 초기 task-instance 구성에서는 None일 수 있으며, task_instance만 선언하는 훅은 변경 없이 계속 작동해요. dag_run.conf는 수동 트리거나 API 트리거된 run에만 채워지고, 스케줄링된 run은 빈 conf를 가진다는 점을 명심하세요.

Warning

task_instance_mutation_hook은 세션 커밋을 금지하는 scheduler 트랜잭션 안에서 실행돼요. 훅 안에서 새 데이터베이스 세션을 열거나 커밋하지 마세요 — 특히 활성 세션을 전달하지 않고 task_instance.get_dagrun()을 호출하지 마세요. 결과 커밋이 scheduler를 크래시시켜요. run 설정을 읽으려면 dag_run 인자를 사용해요.

Dag·Task 클러스터 정책은 AirflowClusterPolicyViolation 예외를 발생시켜 전달된 Dag/task가 준수되지 않아 로드되면 안 된다는 것을 나타낼 수 있어요.

또한 AirflowClusterPolicySkipDag 예외를 발생시켜 그 DAG를 의도적으로 건너뛸 수도 있어요. AirflowClusterPolicyViolation과 달리 이 예외는 Airflow 웹 UI에 표시되지 않아요(내부적으로 meta database의 import_error 테이블에 기록되지 않아요).

클러스터 정책이 설정한 추가 속성은 Dag 파일에서 정의한 것보다 우선해요. 예를 들어 Dag 파일의 Task에 sla를 설정했고 클러스터 정책도 sla를 설정한다면, 클러스터 정책의 값이 우선해요.

정책 함수 정의 방법

클러스터 정책을 구성하는 방법은 두 가지가 있어요:

  1. Python 검색 경로 어딘가에 airflow_local_settings.py 파일을 만들고( $AIRFLOW_HOME 아래의 config/ 폴더가 좋은 "기본" 위치예요), 파일에 위 클러스터 정책 이름 중 하나 이상과 일치하는 callable을 추가해요 (예: dag_policy).

로컬 설정을 구성하는 방법에 대한 자세한 내용은 로컬 설정 구성을 참고해요.

  1. 커스텀 모듈에서 Pluggy 인터페이스를 사용하는 setuptools entrypoint를 사용해요.

버전 2.6에 추가됨.

이 방법은 더 고급이며 Python 패키징에 이미 능숙한 사람들을 위한 것이에요.

먼저 모듈에 정책 함수를 만들어요:

from airflow.policies import hookimpl

@hookimpl
def task_policy(task):
    ...

그런 다음 프로젝트 사양에 entrypoint를 추가해요. 예를 들어 pyproject.tomlsetuptools를 사용해요:

[project.entry-points."airflow.policy"]
my_policy = "my_package.my_policies"

Entrypoint group은 airflow.policy여야 하며, 이름은 항목별로 고유해야 해요. 그렇지 않으면 중복 항목이 pluggy에 의해 무시돼요. 값은 @hookimpl 마커로 장식된 모듈(또는 클래스)이어야 해요.

그렇게 한 뒤 분산(distribution)을 Airflow 환경에 설치하면, 정책 함수가 다양한 Airflow 컴포넌트에 의해 호출돼요. (정확한 호출 순서는 정의되어 있지 않으므로, 여러 플러그인이 있다면 특정 호출 순서에 의존하지 마세요.)

한 가지 중요한 점은 (정책 함수를 정의하는 두 방법 모두) 인자 이름이 아래 문서화된 것과 정확히 일치해야 한다는 거예요.

사용 가능한 정책 함수

airflow.policies.task_policy(task)[source] : DagBag에 로드된 후 Task를 변경할 수 있게 해요.

관리자가 Task의 일부 파라미터를 다시 연결(rewire)할 수 있게 해요. 또는 DAG 실행을 멈추기 위해 `AirflowClusterPolicyViolation` 예외를 발생시킬 수 있어요.

이것이 어떻게 유용할 수 있는지 몇 가지 예시:

- `SparkOperator`를 사용하는 Task에 특정 큐(예: `spark` 큐)를 강제해 그 Task가 올바른 worker에 연결되도록 할 수 있어요.
- Task 타임아웃 정책을 강제해 어떤 Task도 48시간 이상 실행되지 않게 할 수 있어요.

Parameters:
:   **task** – 변경할 task

airflow.policies.dag_policy(dag)[source] : DagBag에 로드된 후 DAG를 변경할 수 있게 해요.

관리자가 DAG의 일부 파라미터를 다시 연결할 수 있게 해요. DAG 실행을 멈추기 위해 `AirflowClusterPolicyViolation` 예외를 발생시킬 수도 있어요.

이것이 어떻게 유용할 수 있는지 몇 가지 예시:

- DAG에 기본 사용자를 강제할 수 있어요.
- 모든 DAG가 태그를 구성했는지 확인할 수 있어요.

Parameters:
:   **dag** – 변경할 dag

airflow.policies.task_instance_mutation_hook(task_instance, dag_run)[source] : Airflow scheduler가 큐에 넣기 전에 task instance를 변경할 수 있게 해요.

예를 들어 재시도 중에 task instance를 수정하거나, run 구성(`dag_run.conf`)에 기반해 Task를 다른 큐로 라우팅하는 데 사용할 수 있어요.

이 훅은 task instance가 생성되거나 조정될 때 scheduler 쪽에서, 세션 커밋을 금지하는 트랜잭션 안에서 실행돼요. 구현은 따라서 새 데이터베이스 세션이나 커밋을 열면 안 되며 — 특히 활성 세션을 전달하지 않고 `task_instance.get_dagrun()`을 호출하면 안 돼요. 결과 커밋이 `RuntimeError("UNEXPECTED COMMIT ...")`를 발생시키고 scheduler를 크래시시켜요. `dag_run` 인자를 대신 사용해요.

`dag_run`은 추가 데이터베이스 접근 없이 사용 가능한 곳 어디든 제공돼요. 초기 task-instance 구성에서는 `None`일 수 있어요 (예: 인스턴스가 run에 바인딩되기 전에 `TaskInstance.refresh_from_task`에서 훅이 재적용될 때). 구현은 그 경우를 처리해야 해요. `dag_run.conf`는 수동 트리거나 API 트리거된 run에만 채워지고, 스케줄링된 run은 빈 `conf`를 가진다는 점을 명심하세요.

Parameters:
:   - **task_instance** (*airflow.models.taskinstance.TaskInstance*) – 변경할 task instance
    - **dag_run** (*airflow.models.dagrun.DagRun* *|* *None*) – task instance가 속한 DagRun, 아직 사용할 수 없으면 `None`

airflow.policies.pod_mutation_hook(pod)[source] : 스케줄링 전에 pod를 변경해요.

이 설정은 `kubernetes.client.models.V1Pod` 객체가 스케줄링을 위해 Kubernetes 클라이언트에 전달되기 전에 변경할 수 있게 해요.

예를 들어 KubernetesExecutor나 KubernetesPodOperator가 실행하는 모든 워커 pod에 sidecar나 init 컨테이너를 추가하는 데 사용할 수 있어요.

airflow.policies.get_airflow_context_vars(context)[source] : airflow context vars를 기본 airflow context vars에 주입해요.

이 설정은 key-value 쌍인 airflow context vars를 가져올 수 있게 해요. 그런 다음 기본 airflow context vars에 주입되며, 결국 Task 실행 시 환경 변수로 사용할 수 있어요. dag_id, task_id, logical_date, dag_run_id, try_number는 예약된 키예요.

Parameters:
:   **context** – 관심 있는 task_instance의 context

예제

Dag 정책

이 정책은 각 DAG에 태그가 하나 이상 정의되어 있는지 확인해요:

from airflow.models.dag import DAG
from airflow.exceptions import AirflowClusterPolicyViolation
from airflow.policies import hookimpl

@hookimpl
def dag_policy(dag: DAG):
    if not dag.tags:
        raise AirflowClusterPolicyViolation(
            f"Dag {dag.dag_id} must have at least one tag"
        )

Note

import 순환을 피하려면, 클러스터 정책에서 타입 애노테이션에 DAG를 사용한다면 airflow가 아니라 airflow.models에서 import해야 해요.

Note

Dag 정책은 DAG가 완전히 로드된 후 적용되므로, default_args 파라미터를 덮어써도 효과가 없어요. 기본 operator 설정을 덮어쓰려면 task 정책을 사용해요.

Task 정책

모든 Task에 최대 타임아웃 정책을 강제하는 예제예요:

from airflow.exceptions import AirflowClusterPolicyViolation
from airflow.models.baseoperator import BaseOperator
from airflow.policies import hookimpl

@hookimpl
def task_policy(task: BaseOperator):
    if task.task_id == "DONT_RUN":
        raise AirflowClusterPolicyViolation("Task should never run")

기술적 보안 통제가 아니라 일반적인 오류로부터 보호하기 위해 구현할 수도 있어요. 예를 들어 Airflow owner가 없는 Task를 실행하지 않게요:

@hookimpl
def task_policy(task: BaseOperator):
    if not task.owner:
        raise AirflowClusterPolicyViolation("Task must have non-None (non-empty) owner")

적용할 검사가 여러 개라면, 이러한 규칙을 별도 Python 모듈에 정리하고 단일 정책/task mutation 훅이 여러 커스텀 검사를 수행해 다양한 오류 메시지를 집계해 단일 AirflowClusterPolicyViolation을 UI(및 데이터베이스의 import errors 테이블)에 보고할 수 있게 하는 것이 모범 사례예요.

예를 들어 airflow_local_settings.py는 이런 패턴을 따를 수 있어요:

TASK_RULES: list[Callable[[BaseOperator], None]] = [
    task_must_have_owners,
]

def _check_task_rules(current_task: BaseOperator):
    """Check task rules for given task."""
    notices = []
    for rule in TASK_RULES:
        try:
            rule(current_task)
        except AirflowClusterPolicyViolation as e:
            notices.append(str(e))
    if notices:
        raise AirflowClusterPolicyViolation("\n".join(notices))

로컬 설정을 구성하는 방법에 대한 자세한 내용은 로컬 설정 구성을 참고해요.

Task instance mutation

두 번째(또는 그 이후) 재시도 중인 Task를 다른 큐로 재라우팅하는 예제예요:

from airflow.models.taskinstance import TaskInstance
from airflow.models.dagrun import DagRun
from airflow.policies import hookimpl

@hookimpl
def task_instance_mutation_hook(task_instance: TaskInstance, dag_run: DagRun | None):
    if not dag_run:
        return
    if task_instance.try_number > 1:
        task_instance.queue = "retry_queue"

우선순위 가중치는 weight rules를 사용해 동적으로 결정되므로, mutation 훅 안에서 task instance의 priority_weight를 변경할 수 없다는 점을 명심해요.

메타데이터 엔진 훅

클러스터 정책에 더해 airflow_local_settings.py는 Airflow가 메타데이터 데이터베이스 엔진을 만드는 방법을 오버라이드할 수 있어요. 이는 정적 설정으로는 표현할 수 없는 per-connection 로직이 필요할 때 유용해요 — 예를 들어 SQLAlchemy do_connect 이벤트 핸들러를 통해 수명이 짧은 JWT 토큰이나 IAM 자격 증명을 주입하는 경우예요.

두 개의 함수를 오버라이드할 수 있어요:

  • create_metadata_engine(sql_alchemy_conn, *, engine_args, connect_args) -> Engineconfigure_orm()이 동기 메타데이터 엔진을 만드는 데 호출해요.
  • create_async_metadata_engine(sql_alchemy_conn_async, *, connect_args, engine_args) -> AsyncEngine_configure_async_session()이 비동기 메타데이터 엔진을 만드는 데 호출해요.

기본 구현은 Airflow가 항상 사용해 온 인자로 sqlalchemy.create_engine / sqlalchemy.ext.asyncio.create_async_engine을 호출하므로, 오버라이드를 제공하지 않는 한 동작 변화는 없어요.

예제: 새 물리 연결마다 JWT 토큰을 새로고침하는 do_connect 핸들러를 등록해요:

from sqlalchemy import event
from airflow.settings import create_metadata_engine

orig = create_metadata_engine
def create_metadata_engine(sql_alchemy_conn, *, engine_args, connect_args):
    engine = orig(sql_alchemy_conn, engine_args=engine_args, connect_args=connect_args )
    def refresh_token(dbapi_conn, connection_record):
        from myauth import get_jwt
        dbapi_conn._jwt = get_jwt()
    event.listen(engine, "do_connect", refresh_token)
    return engine

더 알아보기 (Learn more)