Google ADK

You have an agent built with the Agent Development Kit and you want it to run as an Airflow task, reading data through connections your deployment already manages. Keep the agent as it is, and add Airflow’s toolsets to its tools with the AirflowTools toolset.

Note

AirflowTools and the framework-neutral tool interface it builds on are experimental. They may change in a minor release of this provider. See Stable and experimental features.

Run an ADK agent in a task

The agent below answers a question about a database with a SQLToolset. Its model is Claude through ADK’s AnthropicLlm, with the API key from an Airflow connection:

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

@dag(tags=["example"])
def example_adk_agent():
    """Answer a question about a database with a Google ADK agent."""

    @task
    def run_adk_agent(question: str = DEFAULT_QUESTION) -> str:
        from anthropic import AsyncAnthropic
        from google.adk.agents import LlmAgent
        from google.adk.models.anthropic_llm import AnthropicLlm
        from google.adk.runners import InMemoryRunner
        from google.genai import types

        from airflow.providers.common.ai.tools.adk import AirflowTools
        from airflow.providers.common.ai.tools.tracing import agent_framework_tracing
        from airflow.providers.common.ai.toolsets.sql import SQLToolset

        llm = BaseHook.get_connection(LLM_CONN_ID)
        agent = LlmAgent(
            name="analyst",
            model=AnthropicLlm(
                model=LLM_MODEL,
                client=AsyncAnthropic(api_key=llm.password, base_url=llm.host or None),
            ),
            instruction=(
                "You are a SQL analyst. Use list_tables and get_schema to explore "
                "the database, then run read-only queries to answer the question."
            ),
            tools=[AirflowTools(SQLToolset(db_conn_id=DB_CONN_ID))],
        )

        async def ask() -> str:
            runner = InMemoryRunner(agent=agent, app_name="airflow")
            session = await runner.session_service.create_session(app_name="airflow", user_id="airflow")
            message = types.Content(role="user", parts=[types.Part(text=question)])
            answer = ""
            async for event in runner.run_async(
                user_id="airflow", session_id=session.id, new_message=message
            ):
                if event.is_final_response() and event.content and event.content.parts:
                    answer = "".join(part.text or "" for part in event.content.parts)
            return answer

        # Spans carry the task's identity and no prompt text; see the tracing section of the guide.
        with agent_framework_tracing():
            return asyncio.run(ask())

    run_adk_agent()


Everything except AirflowTools is plain ADK: the agent, its runner and session service, callbacks and your own function tools work as the ADK documentation describes. AirflowTools is an ADK BaseToolset, so it goes in tools= next to any other tool or toolset.

What the toolset gives the agent

AirflowTools takes the same arguments as the Strands plugin (see Strands Agents): toolsets that implement ToolProvider, such as SQLToolset and HookToolset, and individual AirflowTool objects. MCPToolset is the exception: use ADK’s own McpToolset for an MCP server. Each tool keeps the toolset tool’s name, description and argument schema, runs through the toolset, and passes its result through Airflow’s secret masker.

  • A result reaches the model as {"result": ...}.

  • A failure the model can correct, such as a rejected SQL statement or an invalid argument, reaches it as {"error": ...}, ADK’s own convention for a failed tool, and is marked as an error in ADK’s telemetry. The tool’s retry limit bounds how many times in a row the model can try again.

  • A refusal the tool reports as final, such as a file that does not exist, also comes back as {"error": ...} without counting against that limit. ADK’s own RunConfig(max_llm_calls=...) bounds how long a model can keep asking.

  • Any other failure raises ToolCallError out of runner.run_async, with the message masked. The task fails and Airflow’s own retry takes over, which is ADK’s normal behaviour for a tool that raises.

Outside AgentOperator, a toolset’s connection ID is used as written: it is not rendered as a template.

Tracing

The example runs the agent inside agent_framework_tracing(), so ADK’s OpenTelemetry spans carry the task’s identity and leave out prompts, completions and tool inputs and outputs unless [common.ai] capture_content is on. See Observability (OpenTelemetry tracing).

Differences from AgentOperator

  • A task retry runs the agent again from the start. ADK’s session services can persist a session across runs, but resuming one after a worker is lost is not something this provider has tested, so tools that change data must be safe to repeat.

  • There is no human review step like AgentOperator(enable_hitl_review=True).

Installation

pip install apache-airflow-providers-common-ai "google-adk>=2.9.1"

ADK caps opentelemetry-api and opentelemetry-sdk at 1.42.1, which is lower than the version pinned in Airflow’s constraints file, so the command above fails when it is run with that constraints file. Airflow itself accepts OpenTelemetry 1.42.1: build the worker image without the constraints file for this step, or with a copy of it that pins OpenTelemetry to 1.42.1. For AnthropicLlm, install anthropic through this provider’s anthropic extra.

Was this entry helpful?