Setup and Teardown
Setup and Teardown
데이터 워크플로우에서 리소스(예: 컴퓨트 리소스)를 만들고, 사용하고, 정리하는 흔한 패턴을 지원하는 Airflow의 setup/teardown 태스크를 설명하는 문서예요. as_setup()과 as_teardown() 메서드, 컨텍스트 매니저 사용법, setup "scope" 의미, Dag run 상태 제어, 태스크 그룹과 병렬 실행까지 자세히 살펴볼게요.
출처: 문서
본문
데이터 워크플로우에서는 리소스(예: 컴퓨트 리소스)를 만들고, 그것으로 일부 작업을 수행하고, 정리하는 것이 흔해요. Airflow는 이 필요를 지원하기 위해 setup과 teardown 태스크를 제공해요.
setup/teardown 태스크의 핵심 기능:
태스크를 clear하면 그 setup과 teardown도 clear돼요. 기본적으로 teardown 태스크는 Dag run 상태를 평가하는 목적에서는 무시돼요. teardown 태스크는 setup이 성공했다면 작업 태스크가 실패했더라도 실행돼요. 하지만 setup이 skip되었다면 생략돼요. Teardown 태스크는 태스크 그룹에 대한 의존성을 설정할 때 무시돼요. Dag run이 수동으로 "failed"나 "success"로 설정되어도 리소스가 정리되도록 teardown도 수행돼요.
setup과 teardown이 동작하는 방식 (How setup and teardown works)
기본 사용법 (Basic usage)
클러스터를 만들고, 쿼리를 실행하고, 클러스터를 삭제하는 Dag가 있다고 가정해 봐요. setup/teardown 태스크를 사용하지 않으면 이런 관계를 설정할 수 있어요:
create_cluster >> run_query >> delete_cluster
create_cluster와 delete_cluster를 setup/teardown 태스크로 활성화하려면, 그것들을 as_setup과 as_teardown 메서드로 표시하고 그 사이에 업스트림/다운스트림 관계를 추가해요:
create_cluster.as_setup() >> run_query >> delete_cluster.as_teardown()
create_cluster >> delete_cluster
편의상 create_cluster를 as_teardown 메서드에 전달해 한 줄로 이 작업을 할 수 있어요:
create_cluster >> run_query >> delete_cluster.as_teardown(setups=create_cluster)
이 Dag의 그래프는 다음과 같아요:
관찰:
run_query를 다시 실행하려고 clear하면, create_cluster와 delete_cluster가 둘 다 clear돼요.run_query가 실패하면 delete_cluster는 여전히 실행돼요. Dag run의 성공은 오직run_query의 성공에만 의존해요.
추가로, 감싸야 할 태스크가 여러 개라면 teardown을 컨텍스트 매니저로 사용할 수 있어요:
with delete_cluster().as_teardown(setups=create_cluster()):
[RunQueryOne(), RunQueryTwo()] >> DoSomeOtherStuff()
WorkOne() >> [do_this_stuff(), do_other_stuff()]
이렇게 하면 create_cluster가 컨텍스트의 태스크들보다 먼저 실행되고, delete_cluster가 그 뒤에 실행되도록 설정돼요.
그래프로 보면 다음과 같아요:
이미 인스턴스화된 태스크를 setup 컨텍스트에 추가하려고 한다면 명시적으로 해야 한다는 점을 참고하세요:
with my_teardown_task as scope:
scope.add_task(work_task) # work_task was already instantiated elsewhere
Setup "scope"
setup과 teardown 사이의 태스크들은 setup/teardown 쌍의 "scope" 안에 있어요.
예시를 볼게요:
s1 >> w1 >> w2 >> t1.as_teardown(setups=s1) >> w3
w2 >> w4
그리고 그래프:
위 예시에서 w1과 w2는 s1과 t1 "사이"에 있으므로 s1을 요구한다고 간주돼요. 따라서 w1이나 w2가 clear되면 s1과 t1도 clear돼요. 하지만 w3나 w4가 clear되면 s1도 t1도 clear되지 않아요.
단일 teardown에 여러 setup 태스크를 연결할 수 있어요. teardown은 setup 중 적어도 하나가 성공적으로 완료되면 실행돼요.
teardown 없이 setup만 있을 수도 있어요:
create_cluster >> run_query >> other_task
이 경우 create_cluster의 모든 다운스트림이 그것을 요구한다고 간주돼요. 따라서 other_task를 clear하면 create_cluster도 clear돼요. run_query 뒤에 create_cluster의 teardown을 추가한다고 가정해 봐요:
create_cluster >> run_query >> other_task
run_query >> delete_cluster.as_teardown(setups=create_cluster)
이제 Airflow는 other_task가 create_cluster를 요구하지 않는다고 추론할 수 있어서, other_task를 clear하면 create_cluster도 함께 clear되지 않아요.
그 예시에서 (가상의 문서 세계에서) 사실 클러스터를 삭제하고 싶었어요. 하지만 그렇지 않고 단지 "other_task가 create_cluster를 요구하지 않는다"를 말하고 싶다면, EmptyOperator를 사용해 setup의 scope를 제한할 수 있어요:
create_cluster >> run_query >> other_task
run_query >> EmptyOperator(task_id="cluster_teardown").as_teardown(setups=create_cluster)
암묵적 ALL_SUCCESS 제약 (Implicit ALL_SUCCESS constraint)
setup의 scope에 있는 모든 태스크는 그 setup에 암묵적인 "all_success" 제약이 있어요. 이는 간접 setup이 있는 태스크가 clear되면 그것들이 완료될 때까지 기다리도록 보장하기 위해 필요해요. setup이 실패하거나 skip되면, 그것들에 의존하는 작업 태스크는 실패 또는 skip으로 표시돼요. 또한 setup의 다운스트림에 직접 연결된 비-teardown은 ALL_SUCCESS 트리거 규칙을 가져야 해요.
Dag run 상태 제어하기 (Controlling Dag run state)
setup/teardown 태스크의 또 다른 기능은 teardown 태스크가 Dag run 상태에 영향을 미칠지 여부를 선택할 수 있다는 것이에요. teardown 태스크가 수행하는 "cleanup" 작업이 실패해도 상관없고, "작업(work)" 태스크가 실패할 때만 Dag run을 실패로 간주할 수도 있어요. 기본적으로 teardown 태스크는 Dag run 상태에 고려되지 않아요.
위 예시를 이어서, run의 성공이 delete_cluster에 의존하길 원한다면 delete_cluster를 teardown으로 설정할 때 on_failure_fail_dagrun=True를 설정해요. 예를 들어:
create_cluster >> run_query >> delete_cluster.as_teardown(setups=create_cluster, on_failure_fail_dagrun=True)
태스크 그룹으로 작성하기 (Authoring with task groups)
태스크 그룹에서 태스크 그룹으로, 또는 태스크 그룹에서 task로 의존성을 추가할 때는 teardown을 무시해요. 이렇게 하면 teardown을 병렬로 실행할 수 있고, teardown 태스크가 실패해도 Dag 실행이 계속 진행될 수 있어요.
이 예시를 고려해 봐요:
with TaskGroup("my_group") as tg:
s1 = s1()
w1 = w1()
t1 = t1()
s1 >> w1 >> t1.as_teardown(setups=s1)
w2 = w2()
tg >> w2
그래프:
t1이 teardown 태스크가 아니었다면 이 Dag는 사실상 s1 >> w1 >> t1 >> w2가 될 거예요. 하지만 t1을 teardown으로 표시했으므로 tg >> w2에서 무시돼요. 따라서 Dag는 다음과 동등해요:
s1 >> w1 >> [t1.as_teardown(setups=s1), w2]
이제 중첩이 있는 예시를 고려해 봐요:
with TaskGroup("my_group") as tg:
s1 = s1()
w1 = w1()
t1 = t1()
s1 >> w1 >> t1.as_teardown(setups=s1)
w2 = w2()
tg >> w2
dag_s1 = dag_s1()
dag_t1 = dag_t1()
dag_s1 >> [tg, w2] >> dag_t1.as_teardown(setups=dag_s1)
그래프:
이 예시에서 s1은 dag_s1의 다운스트림이므로 dag_s1이 성공적으로 완료될 때까지 기다려야 해요. 하지만 t1과 dag_t1은 동시에 실행될 수 있어요. t1이 tg >> dag_t1 표현식에서 무시되기 때문이에요. w2를 clear하면 dag_s1과 dag_t1이 clear되지만 태스크 그룹 안의 것은 clear되지 않아요.
setup과 teardown 병렬 실행하기 (Running setups and teardowns in parallel)
setup 태스크를 병렬로 실행할 수 있어요:
(
[create_cluster, create_bucket]
>> run_query
>> [delete_cluster.as_teardown(setups=create_cluster), delete_bucket.as_teardown(setups=create_bucket)]
)
그래프:
시각적으로 그룹에 넣는 것도 좋을 수 있어요:
with TaskGroup("setup") as tg_s:
create_cluster = create_cluster()
create_bucket = create_bucket()
run_query = run_query()
with TaskGroup("teardown") as tg_t:
delete_cluster = delete_cluster().as_teardown(setups=create_cluster)
delete_bucket = delete_bucket().as_teardown(setups=create_bucket)
tg_s >> run_query >> tg_t
그리고 그래프:
teardown에 대한 트리거 규칙 동작 (Trigger rule behavior for teardowns)
Teardown은 ALL_DONE_SETUP_SUCCESS라는 (구성 불가능한) 트리거 규칙을 사용해요. 이 규칙에서는 모든 업스트림이 완료되고 직접 연결된 setup이 적어도 하나 성공하면 teardown이 실행돼요. teardown의 모든 setup이 skip되거나 실패했다면 그 상태들이 teardown으로 전파돼요.
수동 Dag 상태 변경에 대한 부작용 (Side-effect on manual Dag state changes)
teardown 태스크는 리소스를 정리하는 데 자주 사용되므로 Dag가 수동으로 종료되어도 실행되어야 해요. 조기 종료를 위해 사용자는 Dag run을 수동으로 "success"나 "failed"로 표시할 수 있는데, 이는 완료 전에 모든 태스크를 중단시켜요. Dag에 teardown 태스크가 있다면 그래도 실행돼요. 따라서 부작용으로, 사용자가 요청해도 teardown 태스크가 스케줄될 수 있게 허용되는 Dag는 즉시 터미널 상태로 설정되지 않아요.