Single prompts: LLMOperator and @task.llm¶
Use LLMOperator for
general-purpose LLM calls: summarization, extraction, classification,
structured output, or any prompt-based task.
The operator sends a prompt to an LLM via
PydanticAIHook and
returns the output as XCom.
See also
Basic usage¶
Provide a prompt and the operator returns the LLM’s response as a string:
@dag(tags=["example"])
def example_llm_operator():
LLMOperator(
task_id="summarize",
prompt="Summarize the key findings from the Q4 earnings report.",
llm_conn_id="pydanticai_default",
system_prompt="You are a financial analyst. Be concise.",
)
Structured output¶
Set output_type to a Pydantic BaseModel subclass and the model instance is pushed to
XCom unchanged, so downstream tasks can type-hint the class and use attribute access.
Structured output and XCom explains the XCom deserialization rules, the cross-Dag gap and
serialize_output.
# Pydantic output classes must be defined at module scope so they survive
# XCom serialization (their qualname is used to re-import them downstream).
class Entities(BaseModel):
"""Named entities extracted from a text."""
names: list[str]
locations: list[str]
@dag(tags=["example"])
def example_llm_operator_structured():
LLMOperator(
task_id="extract_entities",
prompt="Extract all named entities from the article.",
llm_conn_id="pydanticai_default",
system_prompt="Extract named entities.",
output_type=Entities,
)
Agent parameters¶
Pass additional keyword arguments to the pydantic-ai Agent constructor
via agent_params: for example retries, model_settings, or tools.
See the pydantic-ai Agent docs for
the full list of supported parameters.
@dag(tags=["example"])
def example_llm_operator_agent_params():
LLMOperator(
task_id="creative_writing",
prompt="Write a haiku about data pipelines.",
llm_conn_id="pydanticai_default",
system_prompt="You are a creative writer.",
agent_params={"model_settings": {"temperature": 0.9}, "retries": 3},
)
Usage limits¶
Set usage_limits to a
pydantic-ai UsageLimits
to fail the task when the run exceeds a configured budget: request count,
input/output tokens, or tool calls. The check happens inside pydantic-ai’s
run loop, so the limit applies even when retries triggers multiple model
calls within a single task.
@dag(tags=["example"])
def example_llm_operator_usage_limits():
LLMOperator(
task_id="capped_summary",
prompt="Summarize the attached design doc in three bullet points.",
llm_conn_id="pydanticai_default",
system_prompt="You are a concise technical reviewer.",
# Fail the task if the run exceeds 5 model requests, 4_000 input
# tokens, or 1_000 output tokens. Useful for guardrails on shared
# connections or untrusted prompts.
usage_limits=UsageLimits(
request_limit=5,
input_tokens_limit=4_000,
output_tokens_limit=1_000,
# Fail the task if the run's estimated USD cost exceeds $0.50.
# See docs/operators/llm.rst for caveats (not a hard guarantee;
# not enforced for models pydantic-ai can't price, which log a
# warning instead of failing the run).
cost_limit=Decimal("0.50"),
),
)
A plain dict can be passed instead of a UsageLimits instance, which lets
Jinja template individual fields – e.g. a per-run cost cap driven by an Airflow
Variable so the budget can change per environment without editing the Dag:
@dag(tags=["example"])
def example_llm_operator_templated_usage_limits():
LLMOperator(
task_id="capped_summary",
prompt="Summarize the trade-offs of a message queue vs. direct HTTP calls in three bullet points.",
llm_conn_id="pydanticai_default",
system_prompt="You are a concise technical reviewer.",
# A plain dict lets every UsageLimits field be templated -- e.g. driven by
# an Airflow Variable so the budget can change per environment without
# editing the Dag. This caps a single task run, not a day's total spend --
# each run gets the full budget again. Use var.value.get() with a default
# so the example doesn't fail outright if the Variable isn't set.
usage_limits={
"cost_limit": "{{ var.value.get('llm_cost_cap_per_task', '0.50') }}",
"request_limit": 5,
},
)
Each dict value is rendered by Jinja like any other template_fields entry,
then coerced to that field’s type (Decimal, int, or bool). A value
that doesn’t parse – a Variable that exists but is empty renders to "", a
typo renders to a non-numeric string – fails the task with a ValueError
naming the field and the rendered value, instead of silently disabling the
limit. A UsageLimits instance passed directly is used as-is and is not
templated or validated.
Common knobs on UsageLimits:
request_limit: max model requests per run (caps retry/tool-loop blow-ups). pydantic-ai applies a default of50whenUsageLimits()is constructed without an explicit value, so passingUsageLimits(input_tokens_limit=4_000)(or the dict form{"input_tokens_limit": 4_000}) silently inherits that 50-request cap. Setrequest_limit=Noneexplicitly when you only want a token cap.input_tokens_limit/output_tokens_limit: per-run token caps.total_tokens_limit: combined input + output cap.tool_calls_limit: max tool invocations (AgentOperatoronly).cost_limit: aDecimalcap on the run’s estimated USD cost. This is not a hard guarantee against overspend: the response that crosses the limit has already been produced and billed: pydantic-ai checks the accumulated cost after each response and then fails the run withUsageLimitExceeded. It protects you from further spend, not from the request that broke the budget; even a single-request run fails as soon as that request’s cost pushes the total over the limit. Pricing is looked up by model name, not by endpoint: a self-hosted deployment serving a model pydantic-ai recognizes is still priced, at that model’s public list rates rather than at what the deployment actually costs you. That covers vLLM, whose only working prefix isopenai:<model>(see Self-hosted models). A model pydantic-ai cannot price (ollama:llama3.2, a private fine-tune) reports no cost at all, socost_limitis not enforced there – aCostNotFoundWarningis emitted instead of failing the run. And like the other knobs above, settingcost_limitalone still inherits therequest_limit=50default; see therequest_limitnote above. Note thatcost_limitonly caps the operator’s own LLM calls – the meta-agent thatLLMRetryPolicyruns to classify a failed task is a separate, uncapped LLM call; see Retry policies. It counts model spend only, so compute billed by aSandboxToolsetbackend is outside it; see Cost and operations.
When the limit is hit pydantic-ai raises UsageLimitExceeded, which
propagates to Airflow as a task failure, so Airflow’s standard retry policy
applies on top. Every limit here bounds a single agent run, not a task: each
Airflow task retry re-renders usage_limits and starts a fresh count, and for
AgentOperator so does each HITL regeneration. A cost_limit of
Decimal("0.50") caps one run, so it is not a bound on what the task spends in
total.
TaskFlow decorator¶
The @task.llm decorator wraps LLMOperator. The function returns the
prompt string; all other parameters are passed to the operator:
@dag(tags=["example"])
def example_llm_decorator():
@task.llm(llm_conn_id="pydanticai_default", system_prompt="Summarize concisely.")
def summarize(text: str):
return f"Summarize this article: {text}"
summarize("Apache Airflow is a platform for programmatically authoring...")
With structured output:
@dag(tags=["example"])
def example_llm_decorator_structured():
@task.llm(
llm_conn_id="pydanticai_default",
system_prompt="Extract named entities.",
output_type=Entities,
)
def extract(text: str):
return f"Extract entities from: {text}"
extract("Alice visited Paris and met Bob in London.")
Multimodal prompts¶
@task.llm accepts the same prompt shape as @task.agent – the callable
may return either a str or a non-empty Sequence[UserContent] (e.g.,
["Describe this:", ImageUrl(url="...")]) for vision, audio, or document
inputs. See @task.agent multimodal prompts for
the full example. require_approval=True is not supported with a Sequence
prompt: the approval session model expects a string, and the task raises at the
approval boundary.
Classification with Literal¶
Set output_type to a Literal to constrain the LLM to a fixed set of
labels, useful for classification tasks:
@dag(tags=["example"])
def example_llm_classification():
@task.llm(
llm_conn_id="pydanticai_default",
system_prompt=(
"Classify the severity of the given pipeline incident. "
"Use 'critical' for data loss or complete pipeline failure, "
"'high' for significant delays or partial failures, "
"'medium' for degraded performance, "
"'low' for cosmetic issues or minor warnings."
),
output_type=Literal["critical", "high", "medium", "low"],
)
def classify_incident(description: str):
# Pre-process the description before sending to the LLM
return f"Classify this incident:\n{description.strip()}"
classify_incident(
"Scheduler heartbeat lost for 15 minutes. "
"Multiple DAG runs stuck in queued state. "
"No new tasks being scheduled across all DAGs."
)
Multi-task pipeline with dynamic mapping¶
Combine @task.llm with upstream and downstream tasks. Use .expand()
to process a list of items in parallel:
# Pydantic output classes must be defined at module scope so they can be
# imported by name when downstream tasks deserialize the XCom payload.
class TicketAnalysis(BaseModel):
"""Structured analysis of a single support ticket."""
priority: str
category: str
summary: str
suggested_action: str
@dag(tags=["example"])
def example_llm_analysis_pipeline():
"""Triage a queue of support tickets: one model call per ticket, typed results, ready for a schedule."""
@task
def get_support_tickets():
"""Fetch unprocessed support tickets."""
return [
(
"Our nightly ETL pipeline has been failing for the past 3 days. "
"The error shows a connection timeout to the Postgres source database. "
"This is blocking our daily financial reports."
),
(
"We'd like to add a new connection type for our internal ML model registry. "
"Is there documentation on creating custom hooks?"
),
(
"After upgrading to the latest version, the Grid view takes over "
"30 seconds to load for DAGs with more than 500 tasks. "
"Previously it loaded in under 5 seconds."
),
]
@task.llm(
llm_conn_id="pydanticai_default",
system_prompt=(
"Analyze the support ticket and extract: "
"priority (critical/high/medium/low), "
"category (bug/feature_request/question/performance), "
"a one-sentence summary, and a suggested next action."
),
output_type=TicketAnalysis,
)
def analyze_ticket(ticket: str):
return f"Analyze this support ticket:\n\n{ticket}"
@task
def store_results(analyses: list[TicketAnalysis]):
"""Store ticket analyses. In production, this would write to a database or ticketing system."""
for analysis in analyses:
print(f"[{analysis.priority.upper()}] {analysis.category}: {analysis.summary}")
tickets = get_support_tickets()
analyses = analyze_ticket.expand(ticket=tickets)
store_results(analyses)
See also
Dynamic System Prompt –
system_prompt is templated identically on @task.llm, so the same
upstream-XCom pattern applies here.
Human-in-the-loop approval¶
Set require_approval=True to pause the task after the model answers and wait for a human
reviewer. The full guide, including timeouts, notifiers and assigned reviewers, is
Approval gates for LLM operators.
Reviewing uncertain output¶
A decision_policy with a confidence bar sends output the model is unsure about to the
same review flow. See Approval gates for LLM operators.
Parameters¶
prompt: The prompt to send to the LLM (operator) or the return value of the decorated function (decorator).llm_conn_id: Airflow connection ID for the LLM provider.fallback_conn_ids: Connection IDs to fail over to, in order, when the primary provider is unavailable. Overrides the list in the connection’s extra; an explicit[]disables a chain configured there. See Provider fallback.model_id: Model identifier (e.g."openai:gpt-5"). Overrides the connection’s extra field.system_prompt: System-level instructions for the agent. Supports Jinja templating.output_type: Expected output type (default:str). Set to a PydanticBaseModelfor structured output.decision_policy: ADecisionPolicy(min_confidence=..., on_uncertain=...).min_confidenceis the confidence the least confident reporting field needs for the operator to return the output without a person, from 0 to 1; no reported confidence counts as uncertain.on_uncertainis"review"(default) or"fail". DefaultNone: no gate.agent_params: Additional keyword arguments passed to the pydantic-aiAgentconstructor (e.g.retries,model_settings,tools). Supports Jinja templating.usage_limits: Optional pydantic-aiUsageLimitsenforced on the run, or adictof the same fields (templated via Jinja, then coerced per field type). Fails the task when token / request / tool-call budgets are exceeded, or when a templated dict value cannot be coerced. DefaultNone.require_approval: IfTrue, the task defers after generating output and waits for human review. DefaultFalse. Needs Airflow 3.1+.approval_timeout: Maximum time to wait for a review (timedelta).Nonemeans wait indefinitely. DefaultNone.on_approval_timeout: Outcome whenapproval_timeoutexpires without a review:"fail"(default),"approve", or"reject". Requires a review path (require_approval=Trueor adecision_policythat reviews) and a positiveapproval_timeout.allow_modifications: IfTrue, the reviewer can edit the output before approving. DefaultFalse.approval_notifiers: Notifier, or list of notifiers, called once the review is open. DefaultNone.approval_assigned_users: Users allowed to answer the review, as{"id": ..., "name": ...}dicts whereidis the auth manager’s user id.None(default) lets any user with the permission respond. Fixed at first run. Needs Airflow 3.1+.
Logging¶
After each LLM call, the operator logs a summary with model name, token usage, and request count at INFO level. At DEBUG level, the LLM output is also logged (truncated to 500 characters). See AgentOperator logging for details on the log format.