우선순위 가중치

우선순위 가중치 (Priority Weights)

이 페이지는 executor 큐에서 Task의 우선순위를 정하는 priority_weight를 다뤄요. 기본 값은 1이고, 클수록 우선순위가 높아져요. 실제 유효 우선순위 가중치는 weight_rule(가중치 계산 방법)에 따라 결정되는데, 기본 방법은 downstream이에요. 2.9.0부터는 PriorityWeightStrategy 클래스를 상속받아 나만의 커스텀 가중치 규칙을 만들 수도 있어요.

출처: 문서

본문

priority_weight는 executor 큐에서 우선순위를 정의해요. 기본 priority_weight1이며, 어떤 정수로도 올릴 수 있고, 숫자가 클수록 우선순위가 높아져요. 또한 각 Task에는 weight_rule에 기반해 계산되는 진짜 priority_weight가 있는데, 이 weight_rule이 Task의 유효 총 우선순위 가중치에 사용되는 가중치 계산 방법을 정의해요.

아래는 가중치 계산 방법들이에요. 기본적으로 Airflow의 가중치 계산 방법은 downstream이에요.

downstream

Task의 유효 가중치는 모든 하위(다운스트림) 후손들의 합계예요. 그 결과 양수 가중치 값을 사용할 때 상위(업스트림) Task가 더 높은 가중치를 갖고 더 적극적으로 스케줄링돼요. 이는 여러 Dag run 인스턴스가 있고, 각 DAG가 하위 Task를 계속 처리하기 전에 모든 run의 모든 상위 Task가 완료되길 원할 때 유용해요.

upstream

유효 가중치는 모든 상위(업스트림) 조상들의 합계예요. 이는 반대라서, 양수 가중치 값을 사용할 때 하위 Task가 더 높은 가중치를 갖고 더 적극적으로 스케줄링돼요. 이는 여러 Dag run 인스턴스가 있고, 다른 Dag run의 업스트림 Task를 시작하기 전에 각 DAG가 완료되길 선호할 때 유용해요.

absolute

유효 가중치는 추가 계산 없이 지정된 priority_weight 그대로예요. 각 Task가 가져야 할 정확한 우선순위 가중치를 알고 있을 때 이 방법을 쓰면 돼요. 또한 absolute로 설정하면 매우 큰 DAG에서 Task 생성 과정이 크게 빨라지는 보너스 효과도 있어요.

priority_weight 파라미터는 Pools와 함께 사용할 수 있어요.

Note

대부분의 데이터베이스 엔진이 정수를 32비트로 사용하므로, 계산되거나 정의된 priority_weight의 최댓값은 2,147,483,647이고 최솟값은 -2,147,483,648이에요.

커스텀 가중치 규칙

버전 2.9.0에 추가됨.

PriorityWeightStrategy 클래스를 상속받아 나만의 커스텀 가중치 계산 방법을 구현하고, 플러그인에 등록할 수 있어요.

airflow/example_dags/plugins/decreasing_priority_weight_strategy.py[source]

class DecreasingPriorityStrategy(PriorityWeightStrategy):
    """A priority weight strategy that decreases the priority weight with each attempt of the DAG task."""

    def get_weight(self, ti: TaskInstance) -> int:
        try_number = ti.try_number or 0
        return max(3 - try_number + 1, 1)

class DecreasingPriorityWeightStrategyPlugin(AirflowPlugin):
    name = "decreasing_priority_weight_strategy_plugin"
    priority_weight_strategies = [DecreasingPriorityStrategy]

커스텀 우선순위 가중치 전략이 Airflow에 이미 있는지 확인하려면 bash 명령 airflow plugins를 실행하면 돼요. 그다음 사용하려면 커스텀 클래스의 인스턴스를 만들어 Task의 weight_rule 파라미터로 제공하거나, 커스텀 클래스의 경로를 제공하면 돼요:

airflow/example_dags/example_custom_weight.py[source]


with DAG(
    dag_id="example_custom_weight",
    schedule="0 0 * * *",
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    catchup=False,
    dagrun_timeout=datetime.timedelta(minutes=60),
    tags=["example", "example2"],
) as dag:
    start = EmptyOperator(
        task_id="start",
    )

    # provide the class instance
    task_1 = BashOperator(task_id="task_1", bash_command="echo 1", weight_rule=DecreasingPriorityStrategy())

    # or provide the path of the class
    task_2 = BashOperator(
        task_id="task_2",
        bash_command="echo 1",
        weight_rule="airflow.example_dags.plugins.decreasing_priority_weight_strategy.DecreasingPriorityStrategy",
    )

    task_non_custom = BashOperator(task_id="task_non_custom", bash_command="echo 1", priority_weight=2)

    start >> [task_1, task_2, task_non_custom]

DAG가 실행된 후에는 Task의 priority_weight 파라미터를 확인해 커스텀 우선순위 전략 규칙을 사용하고 있는지 검증할 수 있어요.

이것은 실험적 기능이에요.

더 알아보기 (Learn more)