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:
@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 ownRunConfig(max_llm_calls=...)bounds how long a model can keep asking.Any other failure raises
ToolCallErrorout ofrunner.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.