Apache Airflow와 함께하는 지속 실행
Apache Airflow와 함께하는 지속 실행 (Durable Execution with Apache Airflow)
Apache Airflow는 워크플로 오케스트레이터예요. Pydantic AI 통합은 pydantic_ai.durable_exec가 아니라 airflow.providers.common.ai를 통한 apache-airflow-providers-common-ai 패키지가 제공합니다. Airflow의 지속 실행 단위가 Airflow 태스크라는 점이 핵심이에요.
출처: 문서
본문
이 페이지의 래퍼 객체 통합들과 달리, Airflow의 지속 단위는 Airflow 태스크(task) 입니다. 일반 Pydantic AI 에이전트를 작성해서 태스크로 실행하면, Airflow의 재시도 메커니즘과 스텝 수준 캐시가 전체 런을 재생하는 대신 마지막 완료된 모델 요청이나 도구 호출부터 에이전트를 재개합니다.
지속 실행 (Durable Execution)
에이전트가 지속 Airflow 태스크로 실행될 때, Airflow는 완료된 각 모델 요청과 도구 호출을 캐시 항목으로 기록해요. 재시도 시 Airflow는 이 항목들을 재생해 완료된 작업을 건너뜁니다. 각 항목은 그것을 만든 요청의 지문(fingerprint)을 저장하고, 그 지문이 더 이상 일치하지 않으면(이전 시도 이후 대화가 갈라졌으면) 스텝은 낡은 결과를 돌려주는 대신 실시간으로 다시 실행됩니다. 캐시는 단일 DAG 런 태스크의 수명 동안 객체 스토리지(로컬, S3, GCS, Azure)에 있으며, 태스크가 성공하면 삭제됩니다.
예를 들어 에이전트가 모델을 호출해 유용한 응답을 받고, 도구 호출을 시작한 뒤 워커가 크래시했다고 상상해 보세요. 지속 실행이 없으면 Airflow의 일반 재시도가 태스크를 처음부터 다시 시작해 모델 요청과 이후의 사이드 이펙트를 반복합니다. 지속 실행이 있으면 재시도가 런을 재생하고, 이미 완료된 모델 요청의 캐시 결과를 재사용하며, 완료되지 않은 첫 연산부터 계속합니다.
이것은 장기 실행 에이전트와, 반복되는 모델 요청이나 외부 도구 호출이 비용·시간이 들거나 사이드 이펙트를 중복하게 만드는 런에 유용합니다.
지속 에이전트 (Durable Agent)
Airflow의 AgentOperator(또는 @task.agent 데코레이터)로 durable=True를 지정해 에이전트를 지속으로 만들 수 있어요. Airflow와 함께 공급자를 설치합니다:
uv add "apache-airflow-providers-common-ai"
에이전트의 모델과 자격 증명은 Airflow 커넥션에서 옵니다 (예제는 pydanticai_default 사용). 구성 방법은 Pydantic AI connection을 보세요.
지속 실행은 스텝 캐시를 저장할 곳이 필요해요. [common.ai] durable_cache_path를 객체 스토리지 위치로 지정하세요:
[common.ai]
durable_cache_path = s3://my-bucket/airflow-agent-cache
다음은 Airflow 태스크로서의 가장 작은 지속 Pydantic AI 에이전트입니다:
from datetime import timedelta
from airflow.providers.common.ai.operators.agent import AgentOperator
from airflow.sdk import dag
@dag(default_args={"retries": 3, "retry_delay": timedelta(seconds=30)})
def durable_agent_dag():
AgentOperator(
task_id="researcher",
prompt="Summarize quantum error correction.",
llm_conn_id="pydanticai_default",
durable=True,
)
durable_agent_dag()
AgentOperator는 기저의 Pydantic AI 에이전트를 대체하지 않아요. 커넥션과 툴셋에서 에이전트를 만들고, 모델과 툴셋을 감싸서 태스크가 실행되는 동안 회복 가능한 연산을 기록합니다:
- 모델 요청;
- Pydantic AI 도구 호출.
@task.agent 데코레이터도 마찬가지인데, 장식된 함수가 프롬프트를 반환합니다:
from datetime import timedelta
from airflow.providers.common.ai.toolsets.sql import SQLToolset
from airflow.sdk import dag, task
@dag(default_args={"retries": 3, "retry_delay": timedelta(seconds=30)})
def durable_agent_decorator():
@task.agent(
llm_conn_id="pydanticai_default",
system_prompt="You are a data analyst. Use tools to answer questions.",
durable=True,
toolsets=[SQLToolset(db_conn_id="postgres_default", allowed_tables=["orders"])],
)
def analyze(question: str) -> str:
return f"Answer this question about our orders data: {question}"
analyze("What was our total revenue last month?")
durable_agent_decorator()
에이전트의 재시도는 Airflow 태스크 재시도예요. 태스크의 retries와 retry_delay로 구성합니다. 각 재시도에서 캐시된 스텝이 재생되고 런은 완료되지 않은 첫 연산부터 계속합니다. 전체 레퍼런스는 Airflow AgentOperator durable execution docs를 보세요.
도구와 사이드 이펙트 (Tools and side effects)
지속 실행은 각 Pydantic AI 도구 호출의 결과를 캐시하는데, Airflow 툴셋(SQLToolset, HookToolset, MCPToolset 등)이 뒷받침하는 도구들도 포함합니다. 재생 시 도구를 다시 호출하지 않고 캐시된 결과를 반환하므로, 외부 시스템에 쓰는 도구는 런의 모든 재시도 동안 완료된 스텝마다 기껏해야 한 번 실행됩니다.
사람 개입 (Human-in-the-loop)
Airflow의 AgentOperator에는 별도의 사람 개입 검토 모드(enable_hitl_review=True)가 있는데, Airflow의 HITL UI를 통해 에이전트 런을 사람의 승인·거절·변경 요청을 위해 일시 중지합니다.
주의
durable=True와enable_hitl_review=True는 오늘날 결합할 수 없어요. 지속 런은 캐시에서 결정적으로 재생되고 사람 입력을 위해 멈추지 않으며, 사람 개입 런은 상호작용적이고 아직 스텝 캐시에 포착되지 않습니다. 태스크당 하나를 선택하세요.
스트리밍 (Streaming)
스트리밍은 아직 durable=True 아래에서 지원되지 않아요. 지속 모델 래퍼는 완전한 모델 요청을 기록하지 스트리밍된 이벤트를 기록하지 않습니다. 스트리밍 에이전트는 비-지속 태스크로 실행하세요.
요구사항과 제약 (Requirements and Constraints)
Pydantic AI 에이전트를 지속 Airflow 태스크로 실행할 때:
- 지속 단위는 Airflow 태스크이고 회복은 Airflow 태스크 재시도를 통해 일어나므로 태스크에
retries(그리고retry_delay)를 설정하세요. - 구체적인 모델로 에이전트를 정의하세요. 예를 들어
Agent('openai:gpt-5-nano', ...)로 해석되는 커넥션을 씁니다.durable=True일 때 모델이 설정되어 있어야 해요. [common.ai] durable_cache_path를 워커가 읽고 쓸 수 있는 객체 스토리지 위치로 설정하세요.- 재시도 시 저장된 지문이 현재 요청과 일치하면 캐시된 모델 요청·도구 호출이 재생되고, 다르면 실시간으로 다시 실행됩니다. 지문으로 직렬화할 수 없는 요청은 검증되지 않은 위치 재생(positional replay)으로 폴백하므로, 재시도 간에 런을 결정적으로 유지하세요.
durable=True와enable_hitl_review=True는 상호 배타적입니다.durable=True아래에서 스트리밍은 지원되지 않아요.- 스텝 캐시는 하나의 DAG 런 태스크로 한정되며 태스크가 성공하면 삭제됩니다.