Pydantic AI¶
This provider’s operators run Pydantic AI agents, and
AgentOperator builds one for you
from a connection, a prompt and a list of toolsets. If you already have a Pydantic AI
agent, with its own instructions, output types, capabilities or history processing, you
do not have to rebuild it as an AgentOperator. Run it in a @task and give it what
Airflow has: a model from a connection, and toolsets bound to your connections.
Run your own agent in a task¶
@dag(tags=["example"])
def example_pydantic_ai_agent():
"""Answer a question across a database and a reports bucket with your own agent."""
@task
def run_pydantic_ai_agent(question: str = DEFAULT_QUESTION) -> str:
from airflow.providers.common.ai.toolsets.sql import SQLToolset
agent = Agent(
PydanticAIHook.get_hook(LLM_CONN_ID).get_conn(),
instructions=(
"You check reports against the warehouse. Read the report files, query the "
"orders table, and say whether the numbers agree."
),
toolsets=[
SQLToolset(db_conn_id=DB_CONN_ID, allowed_tables=["orders"]),
ObjectStorageToolset(FILES, conn_id=FILES_CONN_ID),
],
)
return agent.run_sync(question).output
run_pydantic_ai_agent()
PydanticAIHook.get_hook returns the hook for the connection’s type, and its
get_conn() returns
the Pydantic AI model the connection describes, with the same vendor prefixes, fallback
connections and self-hosted endpoints as AgentOperator; see Supported model providers.
Every toolset this provider ships is a Pydantic AI toolset, so it goes into
toolsets= as it is, next to any toolsets and tools of your own.
For the OpenTelemetry spans AgentOperator emits, build the agent with the hook’s
create_agent()
instead of Agent(...). It takes the same arguments, and under [common.ai]
otel_export_enabled it sends the agent’s spans through Airflow’s tracing, without prompt
text unless capture_content is on; see Observability (OpenTelemetry tracing). Pydantic AI’s own
instrumentation records prompt and completion text by default.
What you get, and what you do not¶
The SQL, hook, object storage, DataFusion, MCP, sandbox and managed-agent toolsets behave
as they do inside AgentOperator: SQL validation, allowed_tables, result bounds,
object-storage path checks, and the secret masker on everything a tool returns, including
the text of a failure the model is asked to correct. Their calls are counted in the
common_ai.tool_calls metric described in Observability (OpenTelemetry tracing). The Agent Skills
toolset is the exception: its results are masked only inside AgentOperator, and its
calls are not counted.
AgentOperator adds things on top of the agent that a task running your own agent
does not get:
Durable replay of model and tool steps across task retries (
durable=True).Human review of the output (
enable_hitl_review), and a pause for a person to approve marked tool calls before they run (Approve an agent’s tool calls).Masking of the results of the Agent Skills toolset and of toolsets you wrote yourself. There is no public masking wrapper, so run the agent through
AgentOperatorif their results can carry a secret.Rendering of templated connection IDs in toolsets, such as
SQLToolset("{{ ... }}"). In your own task, pass the connection ID itself.Tool call logging, and the task’s identity (
airflow.dag_id,airflow.task_idand the rest) on the agent’s spans.
If you find yourself rebuilding one of these, that is a sign AgentOperator fits:
its agent_params passes any other argument through to the Pydantic AI Agent.