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

airflow/providers/common/ai/example_dags/example_pydantic_ai_agent.py[source]

@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 AgentOperator if 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_id and 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.

Was this entry helpful?