Retry policies¶
Note
Requires Airflow >= 3.3.0.
Two policies ask a model about a task failure and turn the answer into a retry
decision. They are two layers of a ladder that runs from fully hardcoded to
fully model-driven, with the SDK’s ExceptionRetryPolicy as the bottom rung:
Layer |
What decides |
How you tune it |
|---|---|---|
Fallback rules
( |
|
Write the rules. Always the floor: whatever no layer above decides lands here. |
Classifier
( |
The model names one of your |
Category descriptions and the confidence bar. No reasoning, no prose. |
LLM
( |
A text model classifies the failure and decides whether to retry and how
long to wait, guided by |
The instructions: your taxonomy, your retry rules, your delays, in prose. |
Each class has only the arguments its layer needs. LLMRetryPolicy is the
policy this guide has always described and is unchanged. The layers chain:
fallback_policy on a ClassifierRetryPolicy names the policy to consult when
the classifier is unsure or unreachable, typically an LLMRetryPolicy on a
text model, so the cheap typed model handles the clear cases and the reasoning
model the rest, and the rules catch what neither decides. See
Escalating to an LLM below. Both work with any LLM provider supported by
pydantic-ai (OpenAI, Anthropic, Bedrock, Vertex, Ollama, etc.).
This page is about deciding whether a failed task should retry. To make an agent’s
retries cheap by replaying the model and tool calls that already succeeded, see
Durable execution; a retried LLMBatchOperator re-attaches to its running batch
instead of paying for it again (Retries re-attach instead of re-submitting).
For the core retry policy concepts, see Tasks. If the task also needs to survive a worker crash without losing its progress, see Retry policies and durable execution.
Setup¶
Install the provider with the LLM backend you need:
pip install 'apache-airflow-providers-common-ai[anthropic]'
Create a connection (
Admin > Connections):Connection Id:
pydanticai_defaultConnection Type:
Pydantic AIPassword: Your API key
Extra:
{"model": "anthropic:claude-haiku-4-5-20251001"}
Usage¶
from airflow.providers.common.ai.policies.retry import LLMRetryPolicy
from airflow.sdk.definitions.retry_policy import RetryAction, RetryRule
from datetime import timedelta
llm_policy = LLMRetryPolicy(
llm_conn_id="pydanticai_default",
timeout=30.0, # max seconds to wait for LLM response
fallback_rules=[ # used when the LLM call fails
RetryRule(exception=ConnectionError, action=RetryAction.RETRY, retry_delay=timedelta(seconds=10)),
RetryRule(exception=PermissionError, action=RetryAction.FAIL),
],
)
@task(retries=5, retry_policy=llm_policy)
def call_external_api(): ...
How it works¶
When a task fails, either policy:
Sends the exception message to the configured LLM. By default, the message is first masked through Airflow’s secrets masker (see
redactorbelow) and truncated tomax_exception_lengthcharacters before it is added to the prompt.With
LLMRetryPolicy, the model returns anErrorClassification: a category, whether to retry, a suggested delay, and its reasoning. WithClassifierRetryPolicy, it picks one of the names incategories; it sees each category’s name and description in the output schema and cannot answer with a name outside the set.The policy returns RETRY or FAIL. For
LLMRetryPolicythat is the model’sshould_retryandsuggested_delay_seconds. ForClassifierRetryPolicyit is the picked category’sretryanddelay, unless the policy has a confidence bar and the answer is under it, in which case the answer is discarded (see Confidence below).The decision is logged in the task logs and written to the task instance’s
retry_reason, on a FAIL as well as a RETRY:<category>: <reasoning>fromLLMRetryPolicy, or one line such ascategory=network confidence=0.91 threshold=0.60 action=retry delay=10sfromClassifierRetryPolicy.
This classification call is a separate model request, made by the policy
itself rather than by an operator – it is not subject to an operator’s
usage_limits, and it runs on every task failure regardless
of any cost cap configured on the failing task. It is bounded by timeout
and max_exception_length, but not by a cost limit.
If the model call fails (provider down, timeout, bad credentials), the policy
falls back to fallback_rules if configured, or to the task’s standard
retry behaviour. ClassifierRetryPolicy does the same when the model cannot
produce one of the categories even after pydantic-ai re-prompts it, or when
its answer is under the confidence bar, after consulting fallback_policy if
set; its fallback decision’s reason then starts with
classifier answer not applied (model_error), (below_threshold) or
(missing_confidence) so a retry_reason read later is not mistaken for a
classifier decision or a plain rule match. LLMRetryPolicy’s fallback
decision is what fallback_rules returned, as it always was.
This policy decides between attempts. Failing over to another vendor within an attempt is a separate mechanism on the connection; see Provider fallback, which also sets out how the two layers compose.
ClassifierRetryPolicy¶
ClassifierRetryPolicy takes the retry decision away from the model. Its
categories maps a category name to an
ErrorCategory: what
failures belong there (the description the model reads), whether it is
retried, after what delay, and how sure the model has to be
(min_confidence, covered below). Everything the model is told about a
category, and everything the policy does with it, sits in that one entry, so
the two cannot drift apart. This is also the policy a classifier model needs:
such a model refuses the free-text fields of ErrorClassification, so an
LLMRetryPolicy pointed at one fails every classification and falls back,
with a log line saying to use ClassifierRetryPolicy.
DEFAULT_CATEGORIES is the
default: the same seven categories the LLM policy’s default instructions
describe, with the same retry/fail split and delays:
Category |
What belongs there |
Action |
|---|---|---|
|
API throttling or a quota exceeded. |
Retry after 60s |
|
Transient connectivity issue: connection reset, DNS, TLS handshake. |
Retry after 10s |
|
Temporary issue likely to resolve on its own. |
Retry after 30s |
|
Credentials invalid, expired, or missing permissions. |
Fail |
|
Schema validation, type mismatch, or bad input data. |
Fail |
|
Resource not found or unavailable, such as a missing table or bucket. |
Fail |
|
Problem that will not resolve without a code or configuration change. |
Fail |
Pass your own to change any of it. It replaces the default rather than merging into it, so include every category you want offered – or spread the default and edit what you need:
from dataclasses import replace
from datetime import timedelta
from airflow.providers.common.ai.policies.retry import (
ClassifierRetryPolicy,
DEFAULT_CATEGORIES,
ErrorCategory,
)
ClassifierRetryPolicy(
llm_conn_id="pydanticai_default",
categories={
**DEFAULT_CATEGORIES,
# The table is created upstream, so a missing one is worth another look.
"resource": replace(DEFAULT_CATEGORIES["resource"], retry=True, delay=timedelta(minutes=5)),
},
)
Or write a taxonomy for your stack. A Snowflake pipeline has kinds of failure the seven defaults can only approximate, and naming them is what lets the policy act on them differently:
snowflake_policy = ClassifierRetryPolicy(
llm_conn_id="pydanticai_default",
categories={
"queued": ErrorCategory(
"Statement queued or a concurrency limit reached; the warehouse is busy.",
delay=timedelta(seconds=120),
),
"warehouse_suspended": ErrorCategory(
"The warehouse is suspended and will auto-resume.", delay=timedelta(seconds=30)
),
"token_expired": ErrorCategory(
"A JWT or session token expired; the token rotates on its own.", delay=timedelta(seconds=30)
),
"schema_drift": ErrorCategory(
"A referenced column, table or view does not exist; a person has to fix the schema.",
retry=False,
),
},
)
Write descriptions as the boundary between categories: what belongs here and what
does not. That is the whole of what the model reads about a category; the name
alone gives it very little, and two categories with similar names and no
descriptions are indistinguishable to it. A pick is relative: the model chooses
the best fit among the categories offered, not whether any fits. If “none of
these” or “not enough in the message to tell” is a real outcome for your
pipeline, add a category for it (retry=True with no delay keeps the task’s
own behaviour) rather than hoping the model refuses. When you change a
description or the set of categories, treat the confidence values you measured
before as stale: the distribution the model returns is over the options it was
given.
A delay is timedelta, not seconds, and is used as-is: a task’s own
max_retry_delay does not clamp it. delay=None (the default) means no
override, so the task’s own retry_delay / retry_exponential_backoff /
max_retry_delay apply instead (see Tasks).
retry=False ends the task straight away even when attempts were left, so a
wrong classification into such a category costs the task the retries it would
otherwise have had; that is what the confidence bar below is for. A bare int
delay, a negative delay, a delay on a retry=False category, an empty
description, or fewer than two categories raises when the policy is
constructed, at Dag parse time, rather than on the first task failure.
Confidence¶
A classifier model reports how sure it is of its answer. min_confidence is
the bar that answer needs for the policy to act on it; under the bar the policy
discards the answer and takes the same path it takes when the model call fails:
fallback_rules if one matches, otherwise the task’s own retry behaviour. It
does not substitute a delay of its own.
Each category can carry its own bar. The stakes differ: a wrong transient
costs one more attempt, while a wrong permanent costs the task every retry it
had left, so the category that ends the task deserves the higher bar.
jev_default is the classifier-model connection from Classifier models.
# A classifier model answers the same question in a few hundred milliseconds and reports
# how sure it is. The categories are this pipeline's own; their descriptions are what the
# model reads. ``permanent`` ends the task and costs it every retry it had left, so it
# demands more certainty than the rest. Under a bar the answer is discarded and the
# fallback rules, then the task's own retry settings, decide instead. The bars come from
# a calibration run on jev-1.13.0: correct picks landed at 0.89 and above, wrong ones
# at 0.47 to 0.69, with one wrong ``permanent`` at 0.90 that no sensible bar catches.
snowflake_policy = ClassifierRetryPolicy(
llm_conn_id="jev_default",
min_confidence=0.8,
categories={
"queued": ErrorCategory(
"Statement queued or a concurrency limit reached; the warehouse is busy.",
delay=timedelta(seconds=120),
),
"warehouse_suspended": ErrorCategory(
"The warehouse is suspended and will auto-resume.", delay=timedelta(seconds=30)
),
"token_expired": ErrorCategory(
"A JWT or session token expired; the token rotates on its own.", delay=timedelta(seconds=30)
),
"schema_drift": ErrorCategory(
"A referenced column, table or view does not exist; a person has to fix the schema.",
retry=False,
),
"permanent": ErrorCategory(
"A code or configuration error that will fail identically on every attempt.",
retry=False,
min_confidence=0.9,
),
},
fallback_rules=[
RetryRule(exception=ConnectionError, action=RetryAction.RETRY, retry_delay=timedelta(seconds=30)),
],
)
A category bar needs a policy bar to inherit from; setting one without the other raises at construction. With no bar the policy acts on every answer and still logs the confidence.
A missing confidence counts as an unsure one. A text model reports no confidence, and so does a response whose metadata was dropped along the way. With no bar configured that changes nothing. With a bar configured, every such answer is discarded and the fallback path decides, so swapping the connection from a classifier model to a text model does not silently switch off a control you set on purpose. To run a text model, remove the bar.
The confidence is a statistic on the shape of the probability distribution the
model returned: concentrated on one category is high, spread out is low. It is
not the probability that the answer is correct. Pick the bar from the confidence
values your own failures produce: run the policy with no bar first, read the
logged confidence= values per category, and set the bar where the wrong
answers start. A bar reduces wrong actions and does not eliminate them: a wrong
pick can arrive with high confidence. Pin the model version
(typesafe:jev-1.13.0, not jev-latest): a bar tuned against one release
is not guaranteed to mean the same thing after the next. See
Classifier models for what these models answer well and badly.
Escalating to an LLM¶
A classifier is cheap and fast, and a text model can reason about a failure it
has never seen a category for. fallback_policy puts one behind the other.
snowflake_policy is the classifier policy from the previous section; the
chain reuses its category table:
# The three layers chained. The classifier answers the clear cases in a few hundred
# milliseconds. When it is unsure or unreachable, a text model reasons about the failure
# and decides retry and delay itself. When that model is unreachable too, the rules decide.
escalating_policy = ClassifierRetryPolicy(
llm_conn_id="jev_default",
min_confidence=0.8,
categories=snowflake_policy.categories,
fallback_policy=LLMRetryPolicy(llm_conn_id="pydanticai_default", timeout=30.0),
fallback_rules=[
RetryRule(exception=ConnectionError, action=RetryAction.RETRY, retry_delay=timedelta(seconds=30)),
],
)
The order of events on a failure:
The classifier names a category. At or above the bar, its category’s action and delay apply and the text model is never called.
Under the bar, with no confidence reported, or if the classifier call fails, the
fallback_policypolicy runs. A text-modelLLMRetryPolicythere classifies the failure with its own instructions and chooses retry and delay itself. Its decision is used, with the reason prefixed by why the classifier’s answer was not:escalated (below_threshold); rate_limit: 429 with a Retry-After header.If that policy returns DEFAULT, whatever reason it attached, it decided nothing: the outer
fallback_rulesapply, then the task’s own retry behaviour. Only a RETRY or FAIL fromfallback_policyends the chain, so an outer rule such asPermissionError -> FAILstill holds when both models are unreachable.
fallback_policy accepts any RetryPolicy. An ExceptionRetryPolicy works
there too; its default is what it returns when none of its rules match, so
default=RetryAction.FAIL fails every unsure classification and the outer
rules never run. Without min_confidence the classifier’s answer is always
acted on, and fallback_policy is consulted only when the classifier call itself
fails. Both model calls run on the worker at failure time, so a task that
escalates pays for two before its retry is scheduled; timeout on each policy
bounds that.
When the connection also carries a fallback chain¶
Either policy builds its hook from llm_conn_id without passing
fallback_conn_ids, so if that connection’s extra configures a chain (see
Provider fallback), the policy inherits it silently – editing the connection changes
retry behaviour with no change to the Dag. Two things follow:
timeoutstops bounding the whole classification call. pydantic-ai applies aModelSettingstimeout to each model in the chain, not to the chain as a whole, so a 30-secondtimeoutacross a three-connection chain is a 90-second worst case before the policy falls back tofallback_rules.If every connection in the chain fails, the classification call raises
pydantic_ai.exceptions.FallbackExceptionGroup.evaluate()still degrades tofallback_rulescorrectly – it catches the broadException, and an exception group is one – so the only cost here is that the classification is wasted.
Separately, and regardless of this policy: if the connection the task itself uses to call
the LLM (for example llm_conn_id on LLMOperator or AgentOperator) carries a fallback
chain, the exception the task raises once that chain is exhausted is
pydantic_ai.exceptions.FallbackExceptionGroup, not the last provider’s own exception.
RetryRule matches with isinstance, so a rule written as
RetryRule(exception=ModelHTTPError, ...) – in fallback_rules here or in a plain
ExceptionRetryPolicy – stops matching. Match pydantic_ai.exceptions.FallbackExceptionGroup
explicitly as well; its only common ancestor with ModelAPIError is Exception, too
broad to write a rule against. The original per-model exceptions are still available on
FallbackExceptionGroup.exceptions, but RetryRule only compares the top-level
exception type, so a rule set that told 429s apart from 400s collapses into one rule once
the chain is in play.
What the model can and cannot do¶
Under either policy the model is given no tools and there is no way to attach
any, so it cannot run code, call an API, read a connection, or reach your data.
Beyond your instructions, it sees only the exception’s class name, the
exception message (after redaction and truncation), how many attempts are left,
and, under ClassifierRetryPolicy, the category names and descriptions. The prompt says
attempt {try_number} of {max_tries}, so the model knows the limit and not
just where it is right now; an instruction like “after two attempts treat an
expired token as auth rather than transient” works because the model
can see which attempt this is.
Under LLMRetryPolicy it answers four fields: category, should_retry,
suggested_delay_seconds and reasoning, and the first two after
category decide the run. A positive delay is used as returned, with no
upper limit; zero or negative means no override, so the task’s own
retry_delay and backoff apply.
category and reasoning become the retry_reason (truncated to 500
characters), recorded on both outcomes. On a RETRY the value is cleared once the next attempt starts running;
a FAIL is terminal, so there is no next attempt to clear it and the reason stays
on the row. Only the model’s own words are stored – attempt counts are left to
whatever displays the reason. Recording on a FAIL requires Airflow 3.4.0; on
earlier versions only the RETRY outcome is recorded.
Under ClassifierRetryPolicy it answers the category name and nothing else. It does not
decide whether to retry, it does not choose the delay, and it does not explain
itself; the first two come from the category’s entry, and the explanation is
the generated line in the task log, which says what mattered (the category,
the confidence, the bar, the action). A model cannot return a category the
policy does not recognize, and it cannot return a category paired with an
action that contradicts it.
RETRY cannot give a task more attempts than retries allows. FAIL ends the
task straight away even when attempts were left, so a wrong classification into
a failing category costs the task the retries it would otherwise have had;
ClassifierRetryPolicy’s confidence bar exists for that.
Custom instructions¶
For LLMRetryPolicy the instructions are the taxonomy, and
DEFAULT_INSTRUCTIONS shows
the shape: name the categories, say which to retry, and give the delays. Override
them to teach the model your stack, and name your own categories if the seven
defaults do not fit; the model returns whatever name you taught it, with its own
retry decision and delay:
SNOWFLAKE_INSTRUCTIONS = (
"You are an error classifier for Snowflake-backed data pipelines. "
"Classify the error into one of: queued, warehouse_suspended, token_expired, "
"schema_drift, permanent.\n\n"
"- 'Statement queued' or 'concurrency limit' -> queued, retry after 120s\n"
"- '000606' or 'is suspended' -> warehouse_suspended, retry after 30s\n"
"- 'JWT token' or 'session token' with 'expired' -> token_expired, retry after 30s\n"
"- '002003' or 'does not exist' -> schema_drift, do NOT retry\n"
"- anything else that will fail identically -> permanent, do NOT retry\n"
"Set suggested_delay_seconds as above and 0 for errors that should not retry."
)
snowflake_policy = LLMRetryPolicy(llm_conn_id="pydanticai_default", instructions=SNOWFLAKE_INSTRUCTIONS)
For ClassifierRetryPolicy the same taxonomy lives in categories, and
the default instructions
(CLASSIFIER_INSTRUCTIONS) say
only that the model is classifying a failed pipeline task and should pick the
best-fitting category. Override instructions there to teach the model your
error strings when a description is not enough on its own, and leave the
category names, actions and delays to the table; a prompt that names them only
spends tokens, because the model is constrained to the table’s names and does
not decide the action:
SNOWFLAKE_HINTS = (
"You are classifying failures from Snowflake-backed data pipelines.\n"
"- 'Statement queued' or 'concurrency limit' -> queued\n"
"- '000606' or 'is suspended' -> warehouse_suspended\n"
"- 'JWT token' or 'session token' with 'expired' -> token_expired\n"
"- '002003' or 'does not exist' -> schema_drift\n"
)
snowflake_policy = ClassifierRetryPolicy(
llm_conn_id="pydanticai_default",
instructions=SNOWFLAKE_HINTS,
categories={...}, # the same names the hints use
fallback_rules=[
RetryRule(
exception=ConnectionError,
action=RetryAction.RETRY,
retry_delay=timedelta(seconds=30),
),
],
)
@task(retries=5, retry_policy=snowflake_policy)
def query_snowflake(): ...
When writing custom instructions:
Be concrete with examples (
"'Warehouse suspended' -> warehouse_suspended") rather than vague rules (“treat warehouse issues as recoverable”).For
LLMRetryPolicy, mention the fourErrorClassificationfields so the model fills them, and keepreasoningconcise:retry_reasonis truncated to 500 characters.For
ClassifierRetryPolicy, use the table’s names as-is. A name you invent cannot come back: a model that insists on one is re-prompted once by pydantic-ai and then gives up, which lands the task onfallback_rulesor on its own retry behaviour, having billed two calls.A classifier model sends
instructionsas the question it scores the exception text against, not as rules it follows step by step, so a long rubric buys less there than a better description on each category does.
Parameters¶
Both policies share every parameter below except categories,
min_confidence and fallback_policy, which exist only on
ClassifierRetryPolicy.
Parameter |
Default |
Description |
|---|---|---|
|
(required) |
Airflow connection ID for the LLM provider. |
|
None |
Override the model from the connection (e.g., |
|
(built-in) |
Custom system prompt for error classification. On |
|
None |
List of |
|
30.0 |
Max seconds to wait for the LLM response before falling back. |
|
|
|
|
None |
|
|
None |
|
|
None (uses |
Callable |
|
True |
Whether to redact the exception’s string representation before it is
added to the classification prompt. Set to |
|
4096 |
Maximum number of characters of the (already redacted) exception
message included in the prompt. Longer messages are truncated with a
trailing |
Custom redactors¶
The default redactor only masks values already registered with Airflow’s
secrets masker via mask_secret(). It does not detect free-text PII –
email addresses, customer names, account numbers – that were never
registered as secrets. If your task’s exception messages can contain that
kind of data, supply your own redactor callable. It replaces the
default masker rather than running in addition to it, so combine your own
logic with redact_registered_secrets()
yourself if you still want known-secret masking too:
import re
from airflow.providers.common.ai.policies.retry import redact_registered_secrets
EMAIL_RE = re.compile(r"[\w.+-]+@[\w-]+\.[\w.-]+")
def redact_emails_and_secrets(message: str) -> str:
return redact_registered_secrets(EMAIL_RE.sub("<email>", message))
llm_policy = LLMRetryPolicy(
llm_conn_id="pydanticai_default",
redactor=redact_emails_and_secrets,
max_exception_length=2048, # keep long tracebacks from inflating token cost
)
To disable redaction entirely (for example, if you are certain your
exception messages contain no sensitive data and need the raw text for
accurate classification), pass redact_exception=False:
LLMRetryPolicy(llm_conn_id="pydanticai_default", redact_exception=False)
Local LLM support¶
By default, the built-in redactor already masks known secrets before the
exception data reaches the LLM provider. For environments where exception
data must not leave your own infrastructure at all – even in masked form –
point to a local model via Ollama or vLLM instead, so the classification
never crosses the network boundary. See Self-hosted models for
general self-hosted connection setup:
LLMRetryPolicy(
llm_conn_id="ollama_local", # host=http://localhost:11434
model_id="ollama:llama3.2",
)