우선순위 가중치
우선순위 가중치 (Priority Weights)
이 페이지는 executor 큐에서 Task의 우선순위를 정하는 priority_weight를 다뤄요. 기본 값은 1이고, 클수록 우선순위가 높아져요. 실제 유효 우선순위 가중치는 weight_rule(가중치 계산 방법)에 따라 결정되는데, 기본 방법은 downstream이에요. 2.9.0부터는 PriorityWeightStrategy 클래스를 상속받아 나만의 커스텀 가중치 규칙을 만들 수도 있어요.
출처: 문서
본문
priority_weight는 executor 큐에서 우선순위를 정의해요. 기본 priority_weight는 1이며, 어떤 정수로도 올릴 수 있고, 숫자가 클수록 우선순위가 높아져요. 또한 각 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 파라미터를 확인해 커스텀 우선순위 전략 규칙을 사용하고 있는지 검증할 수 있어요.
이것은 실험적 기능이에요.