HITLOperator
HITLOperator (Human-in-the-loop)
HITL(Human-in-the-Loop) 기능을 사용하면 사람의 의사 결정을 워크플로우에 직접 넣을 수 있어요. 워크플로우가 멈추고 사람의 입력을 기다리게 할 수 있어 승인 프로세스, 수동 품질 검사, 사람의 판단이 필수적인 시나리오에 아주 적합한 기능이에요.
출처: 문서
본문
버전 3.1에 추가됨.
Human-in-the-Loop(HITL) 기능을 사용하면 사람의 의사 결정을 워크플로우에 직접 통합할 수 있어요. 이 강력한 기능은 워크플로우가 멈추고 사람의 입력을 기다리게 해서, 승인 프로세스, 수동 품질 검사, 사람의 판단이 필수적인 시나리오에 완벽해요.
버전 3.3 변경: 입력을 기다리는 HITL Task는 이제 triggerer에 defer하는 대신 전용의 스케줄러 관리 awaiting_input Task 상태를 사용해요. 기다리는 동안 Task는 워커 슬롯도 triggerer도 점유하지 않으므로, Task가 응답을 기다리는 동안에도 triggerer는 0으로 스케일될 수 있어요. Task는 사람의 응답이나 스케줄러의 응답 타임아웃 스윕(sweep)에 따라 재개돼요. Airflow 3.1과 3.2에서 HITL Task는 이전의 trigger 기반 deferral을 사용해요.
대기 중인 awaiting_input Task는 pool 슬롯을 점유하지 않아요. 이는 이전 deferral 경로와 다르며, defer된 HITL Task는 include_deferred가 활성화된 pool에 대해 계산되었어요.
이 튜토리얼에서는 워크플로우에서 HITL operators를 사용하는 방법을 살펴보고, Airflow UI에서 어떻게 보이는지 보여드릴게요.
HITL 예시 DAG
HITL이 DAG에서 어떻게 보이는지에 대한 예시예요. 하나씩 쪼개서 자세히 살펴볼게요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
class LocalLogNotifier(BaseNotifier):
"""Simple notifier to demonstrate HITL notification without setup any connection."""
template_fields = ("message",)
def __init__(self, message: str) -> None:
self.message = message
def notify(self, context: Context) -> None:
url = HITLOperator.generate_link_to_ui_from_context(
context=context,
base_url="http://localhost:28080",
)
self.log.info(self.message)
self.log.info("Url to respond %s", url)
hitl_request_callback = LocalLogNotifier(
message="""
[HITL]
Subject: {{ task.subject }}
Body: {{ task.body }}
Options: {{ task.options }}
Is Multiple Option: {{ task.multiple }}
Default Options: {{ task.defaults }}
Params: {{ task.params }}
"""
)
hitl_success_callback = LocalLogNotifier(
message="{% set task_id = task.task_id -%}{{ ti.xcom_pull(task_ids=task_id) }}"
)
hitl_failure_callback = LocalLogNotifier(message="Request to response to '{{ task.subject }}' failed")
with DAG(
dag_id="example_hitl_operator",
start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
catchup=False,
tags=["example", "HITL"],
):
wait_for_input = HITLEntryOperator(
task_id="wait_for_input",
subject="Please provide required information: ",
params={"information": Param("", type="string")},
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
wait_for_option = HITLOperator(
task_id="wait_for_option",
subject="Please choose one option to proceed: ",
options=["option 1", "option 2", "option 3"],
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
wait_for_multiple_options = HITLOperator(
task_id="wait_for_multiple_options",
subject="Please choose option to proceed: ",
options=["option 4", "option 5", "option 6"],
multiple=True,
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
wait_for_default_option = HITLOperator(
task_id="wait_for_default_option",
subject="Please choose option to proceed: ",
options=["option 7", "option 8", "option 9"],
defaults=["option 7"],
response_timeout=datetime.timedelta(seconds=1),
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
valid_input_and_options = ApprovalOperator(
task_id="valid_input_and_options",
subject="Are the following input and options valid?",
body="""
Input: {{ ti.xcom_pull(task_ids='wait_for_input')["params_input"]["information"] }}
Option: {{ ti.xcom_pull(task_ids='wait_for_option')["chosen_options"] }}
Multiple Options: {{ ti.xcom_pull(task_ids='wait_for_multiple_options')["chosen_options"] }}
Timeout Option: {{ ti.xcom_pull(task_ids='wait_for_default_option')["chosen_options"] }}
""",
defaults="Reject",
response_timeout=datetime.timedelta(minutes=5),
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
assigned_users=[{"id": "1", "name": "airflow"}, {"id": "admin", "name": "admin"}],
)
choose_a_branch_to_run = HITLBranchOperator(
task_id="choose_a_branch_to_run",
subject="You're now allowed to proceeded. Please choose one task to run: ",
options=["task_1", "task_2", "task_3"],
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
@task
def task_1(): ...
@task
def task_2(): ...
@task
def task_3(): ...
(
[wait_for_input, wait_for_option, wait_for_default_option, wait_for_multiple_options]
>> valid_input_and_options
>> choose_a_branch_to_run
>> [task_1(), task_2(), task_3()]
)
입력 제공 (Input Provision)
사용자는 이후 Task에 사용되는 params를 사용해 입력을 제공할 수 있어요. 이는 대규모 언어 모델(LLM) 워크플로우 내에서 사람의 안내가 필요한 워크플로우에 유용해요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
wait_for_input = HITLEntryOperator(
task_id="wait_for_input",
subject="Please provide required information: ",
params={"information": Param("", type="string")},
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
Task를 클릭하면 상세 패널에서 Required Actions 탭을 찾을 수 있어요.

옵션 선택 (Option Selection)
입력은 옵션 형태로 제공할 수 있어요. 사용자는 사용 가능한 옵션 중 하나를 선택할 수 있으며, 이를 워크플로우를 안내하는 데 사용할 수 있어요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
wait_for_option = HITLOperator(
task_id="wait_for_option",
subject="Please choose one option to proceed: ",
options=["option 1", "option 2", "option 3"],
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)

여러 옵션도 허용돼요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
wait_for_multiple_options = HITLOperator(
task_id="wait_for_multiple_options",
subject="Please choose option to proceed: ",
options=["option 4", "option 5", "option 6"],
multiple=True,
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)

승인 또는 거부 (Approval or Rejection)
옵션 선택의 특수한 형태로, 옵션이 '승인(Approval)'과 '거부(Rejection)'만 있는 경우예요. HITL operator에 응답할 수 있는 사용자를 제한하기 위해 assigned_users를 설정할 수도 있어요. 이는 사용자 id와 사용자 이름(둘 다 필요)의 목록이어야 해요 (예: [{"id": "1", "name": "user1"}, {"id": "2", "name": "user2"}]. 이 목록에 있는 사용자만 응답할 수 있어요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
valid_input_and_options = ApprovalOperator(
task_id="valid_input_and_options",
subject="Are the following input and options valid?",
body="""
Input: {{ ti.xcom_pull(task_ids='wait_for_input')["params_input"]["information"] }}
Option: {{ ti.xcom_pull(task_ids='wait_for_option')["chosen_options"] }}
Multiple Options: {{ ti.xcom_pull(task_ids='wait_for_multiple_options')["chosen_options"] }}
Timeout Option: {{ ti.xcom_pull(task_ids='wait_for_default_option')["chosen_options"] }}
""",
defaults="Reject",
response_timeout=datetime.timedelta(minutes=5),
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
assigned_users=[{"id": "1", "name": "airflow"}, {"id": "admin", "name": "admin"}],
)
이 코드 조각의 body에서 볼 수 있듯이, XCom을 사용해 사용자가 제공한 정보를 가져올 수 있어요.

브랜치 선택 (Branch Selection)
사용자는 DAG 내에서 따를 브랜치를 선택할 수 있어요. 이는 사람의 판단이 때로 필요한 콘텐츠 검토(moderation) 같은 시나리오에서 흔히 사용돼요.
이것은 옵션 선택과 비슷하지만, 옵션이 Task여야 해요. 그리고 워크플로우에서 그들의 관계를 지정하는 것을 잊지 마세요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
choose_a_branch_to_run = HITLBranchOperator(
task_id="choose_a_branch_to_run",
subject="You're now allowed to proceeded. Please choose one task to run: ",
options=["task_1", "task_2", "task_3"],
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
@task
def task_1(): ...
@task
def task_2(): ...
@task
def task_3(): ...
(
[wait_for_input, wait_for_option, wait_for_default_option, wait_for_multiple_options]
>> valid_input_and_options
>> choose_a_branch_to_run
>> [task_1(), task_2(), task_3()]
)

브랜치를 선택하면 워크플로우는 선택된 경로를 따라 진행돼요.

Notifiers
Notifier는 HITL 이벤트(예: Task가 사람 입력을 기다리거나, 성공하거나, 실패하는 경우)를 처리하기 위한 콜백 메커니즘이에요. 예시에서는 데모 목적으로 메시지를 기록하는 LocalLogNotifier를 사용해요.
HITLOperator.generate_link_to_ui_from_context 메서드는 사용자가 응답해야 하는 UI 페이지로의 직접 링크를 생성하는 데 사용할 수 있어요. 네 개의 인자를 받아요:
context– notifier에 의해notify에 자동으로 전달됨base_url– (선택) Airflow UI의 기본 URL; 제공하지 않으면 구성의api.base_url이 사용됨options– (선택) UI 페이지에 미리 선택된 옵션params_inputs– (선택) UI 페이지에 미리 로드된 입력
이렇게 하면 알림이나 로그에 실행 가능한 링크를 쉽게 포함할 수 있어요. 다른 기능을 제공하기 위해 자신만의 notifier를 구현할 수도 있어요. 자세한 내용은 Notifier 만들기와 Notifications를 참고하세요.
예시 DAG에서 notifier는 다음과 같이 정의돼요:
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
class LocalLogNotifier(BaseNotifier):
"""Simple notifier to demonstrate HITL notification without setup any connection."""
template_fields = ("message",)
def __init__(self, message: str) -> None:
self.message = message
def notify(self, context: Context) -> None:
url = HITLOperator.generate_link_to_ui_from_context(
context=context,
base_url="http://localhost:28080",
)
self.log.info(self.message)
self.log.info("Url to respond %s", url)
hitl_request_callback = LocalLogNotifier(
message="""
[HITL]
Subject: {{ task.subject }}
Body: {{ task.body }}
Options: {{ task.options }}
Is Multiple Option: {{ task.multiple }}
Default Options: {{ task.defaults }}
Params: {{ task.params }}
"""
)
hitl_success_callback = LocalLogNotifier(
message="{% set task_id = task.task_id -%}{{ ti.xcom_pull(task_ids=task_id) }}"
)
hitl_failure_callback = LocalLogNotifier(message="Request to response to '{{ task.subject }}' failed")
notifiers 인자를 사용해 HITL operators에 notifier 목록을 전달할 수 있어요. operator가 사람의 응답을 기다리는 HITL 요청을 만들면, 단일 인자 context로 notify 메서드가 호출돼요.
/opt/airflow/providers/standard/src/airflow/providers/standard/example_dags/example_hitl_operator.py
wait_for_input = HITLEntryOperator(
task_id="wait_for_input",
subject="Please provide required information: ",
params={"information": Param("", type="string")},
notifiers=[hitl_request_callback],
on_success_callback=hitl_success_callback,
on_failure_callback=hitl_failure_callback,
)
HITL DAG 로컬에서 테스트하기
airflow dags test(그리고 그 기반인 dag.test())는 HITL Task를 지원해요. awaiting_input 상태에 도달한 Task는 주차(parked) 상태로 남아요 — 테스트 실행이 그 자체로 해결하지 않으며 — 실행은 어떤 Task가 입력을 기다리는지 로깅하면서, 외부에서 응답이 기록될 때까지 기다려요. 응답은 실제 배포와 같은 채널을 통해 흘러가요: 메타데이터 데이터베이스를 공유하는 api-server의 Required Actions 페이지나 HITL REST API(PATCH .../hitlDetails) (예: airflow standalone, 또는 별도로 시작된 airflow api-server). 응답이 도착하면 테스트 실행은 Task를 재개하고 다운스트림 Task를 계속 진행해요.
또한 이를 통해 AI 에이전트가 HITL 파이프라인을 로컬에서 end-to-end로 구동할 수 있어요: airflow dags test를 실행하고, 기다리는 로그 줄을 보고, 사람에게 물어보고, HITL REST API를 통해 답변을 제출해요. 관련된 두 호출은 (~는 dag_id와 dag_run_id의 와일드카드로 동작):
# Discover pending requests (subject, options, params, run/task identifiers)
GET /api/v2/dags/~/dagRuns/~/hitlDetails?response_received=false
# Submit the response; the test run resumes the task on its next poll.
# map_index is -1 for non-mapped tasks.
PATCH /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}/{map_index}/hitlDetails
{"chosen_options": ["Approve"], "params_input": {}}
Note
response_timeout과 타임아웃 기본값은airflow dags test아래에서는 실행되지 않는 스케줄러에 의해 강제돼요. 따라서 주차된 Task는 응답을 무한정 기다려요. 실행을 끝내려면 UI나 REST API를 통해 응답을 제공하세요.
이점과 일반적인 사용 사례
HITL 기능은 대규모 언어 모델(LLM) 워크플로우에서 가치 있는데, 사람이 제공하는 안내가 더 나은 결과를 얻는 데 필수적일 수 있어요. 또한 자동화된 프로세스를 사람의 검증이 보완하고 강화할 수 있는 엔터프라이즈 데이터 파이프라인에서도 매우 유익해요.