创建 Agent2Agent 智能体

借助 Agent Runtime,您可以使用 Agent2Agent (A2A) 协议开发和部署代理。A2A 是一种开放标准,旨在让 AI 代理之间实现无缝通信和协作。

本文档介绍了如何在本地开发和测试 A2A 智能体,包括定义 AgentCardAgentExecutor 等组件。

如需详细了解如何管理已部署的代理,请参阅管理已部署的代理

核心工作流涉及以下步骤:

  1. 定义关键组件
  2. 创建本地智能体
  3. 测试本地智能体

定义智能体组件

如需创建 A2A 智能体,您需要定义以下组件:AgentCardAgentExecutor 和 ADK LlmAgent

  • AgentCard 包含一个描述智能体功能的元数据文档。AgentCard 就像一张商家名片,其他智能体可以使用它来了解您的智能体可以做什么。如需了解详情,请参阅智能体卡片规范
  • AgentExecutor 包含智能体的核心逻辑,并定义了智能体如何处理任务。您可以用该组件实现智能体的行为。如需详细了解,请参阅 A2A 协议规范
  • 可选:LlmAgent 定义 ADK 智能体,包括其系统指令、生成式模型和工具。

定义 AgentCard

以下代码示例为一个货币汇率智能体定义了 AgentCard

from a2a.types import AgentCard, AgentSkill
from vertexai.agent_engines.templates.a2a import create_agent_card

# Define the skill for the CurrencyAgent
currency_skill = AgentSkill(
    id='get_exchange_rate',
    name='Get Currency Exchange Rate',
    description='Retrieves the exchange rate between two currencies on a specified date.',
    tags=['Finance', 'Currency', 'Exchange Rate'],
    examples=[
        'What is the exchange rate from USD to EUR?',
        'How many Japanese Yen is 1 US dollar worth today?',
    ],
)

# Create the agent card using the utility function
agent_card = create_agent_card(
    agent_name='Currency Exchange Agent',
    description='An agent that can provide currency exchange rates',
    skills=[currency_skill]
)

定义 AgentExecutor

以下代码示例定义了一个 AgentExecutor,用于在回答中提供货币汇率。它接受 CurrencyAgent 实例并初始化 ADK 运行程序以执行请求。

import requests
from a2a.server.agent_execution.agent_executor import AgentExecutor
from a2a.server.agent_execution.context import RequestContext
from a2a.server.events.event_queue import EventQueue
from a2a.server.tasks import TaskUpdater
from a2a import types as a2a_types
from a2a.types import Part

from google.adk import Runner
from google.adk.agents import LlmAgent
from google.adk.artifacts.in_memory_artifact_service import InMemoryArtifactService
from google.adk.memory.in_memory_memory_service import InMemoryMemoryService
from google.adk.sessions.in_memory_session_service import InMemorySessionService
from google.genai import types as genai_types

class CurrencyAgentExecutorWithRunner(AgentExecutor):
    """Executor that takes an LlmAgent instance and initializes the ADK Runner internally."""

    def __init__(self, agent: LlmAgent):
        self.agent = agent
        self.runner = None

    def _init_adk(self):
        if not self.runner:
            self.runner = Runner(
                app_name=self.agent.name,
                agent=self.agent,
                artifact_service=InMemoryArtifactService(),
                session_service=InMemorySessionService(),
                memory_service=InMemoryMemoryService(),
            )

    async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        task_id = context.task_id
        updater = TaskUpdater(
            event_queue=event_queue,
            task_id=task_id or "",
            context_id=context.context_id or "",
        )
        await updater.cancel()

    async def execute(
        self,
        context: RequestContext,
        event_queue: EventQueue,
    ) -> None:
        self._init_adk() # Initialize on first execute call

        if not context.message:
            return

        user_id = context.message.metadata.get('user_id') if context.message and context.message.metadata else 'a2a_user'

        updater = TaskUpdater(event_queue, context.task_id, context.context_id)
        
        task = a2a_types.Task(
            id=context.task_id,
            context_id=context.context_id,
            status=a2a_types.TaskStatus(state=a2a_types.TaskState.TASK_STATE_SUBMITTED),
            history=[context.message] if context.message else [],
        )
        await event_queue.enqueue_event(task)

        await updater.start_work()

        query = context.get_user_input()
        content = genai_types.Content(role='user', parts=[genai_types.Part.from_text(text=query)])

        try:
            session = await self.runner.session_service.get_session(
                app_name=self.runner.app_name,
                user_id=user_id,
                session_id=context.context_id,
            ) or await self.runner.session_service.create_session(
                app_name=self.runner.app_name,
                user_id=user_id,
                session_id=context.context_id,
            )

            final_event = None
            async for event in self.runner.run_async(
                session_id=session.id,
                user_id=user_id,
                new_message=content
            ):
                if event.is_final_response():
                    final_event = event

            if final_event and final_event.content and final_event.content.parts:
                response_text = "".join(
                    part.text for part in final_event.content.parts if hasattr(part, 'text') and part.text
                )
                if response_text:
                    await updater.add_artifact(
                        [Part(text=response_text)],
                        name='result',
                        last_chunk=True,
                    )
                    await updater.complete()
                    return

            await updater.update_status(
                a2a_types.TaskState.TASK_STATE_FAILED,
                message=updater.new_agent_message([Part(text='Failed to generate a final response with text content.')]),
            )

        except Exception as e:
            await updater.update_status(
                a2a_types.TaskState.TASK_STATE_FAILED,
                message=updater.new_agent_message([Part(text=f"An error occurred: {str(e)}")]),
            )

定义 LlmAgent

首先,为 LlmAgent 定义要使用的货币汇率工具:

def get_exchange_rate(
    currency_from: str = "USD",
    currency_to: str = "EUR",
    currency_date: str = "latest",
):
    """Retrieves the exchange rate between two currencies on a specified date.
    Uses the Frankfurter API (https://api.frankfurter.app/) to obtain
    exchange rate data.
    """
    try:
        response = requests.get(
            f"https://api.frankfurter.app/{currency_date}",
            params={"from": currency_from, "to": currency_to},
        )
        response.raise_for_status()
        return response.json()
    except requests.exceptions.RequestException as e:
        return {"error": str(e)}

然后,定义使用该工具的 ADK LlmAgent

my_llm_agent = LlmAgent(
    model='gemini-2.0-flash',
    name='currency_exchange_agent',
    description='An agent that can provide currency exchange rates.',
    instruction="""You are a helpful currency exchange assistant.
                   Use the get_exchange_rate tool to answer user questions.
                   If the tool returns an error, inform the user about the error.""",
    tools=[get_exchange_rate],
)

创建本地智能体

定义智能体的组件后,创建一个使用 AgentCardAgentExecutorLlmAgentA2aAgent 类实例,以开始本地测试。

from vertexai.agent_engines.templates.a2a import A2aAgent

a2a_agent = A2aAgent(
    agent_card=agent_card, # Assuming agent_card is defined
    agent_executor_builder=lambda: CurrencyAgentExecutorWithRunner(
        agent=my_llm_agent,
    )
)
a2a_agent.set_up()

A2A 智能体模板可帮助您创建符合 A2A 标准的服务。该服务充当封装容器,可为您抽象出转换层。

测试本地智能体

该货币汇率智能体支持以下三种方法:

  • handle_authenticated_agent_card
  • on_message_send
  • on_get_task

测试 handle_authenticated_agent_card

以下代码会检索智能体经过身份验证的卡片,其中描述了智能体的功能。

# Test the `authenticated_agent_card` endpoint.
response_get_card = await a2a_agent.handle_authenticated_agent_card(request=None, context=None)
print(response_get_card)

测试 on_message_send

以下代码会模拟客户向智能体发送新消息。A2aAgent 可创建新任务并返回任务的 ID。

from a2a.types import SendMessageRequest, Message, Part
from a2a.server.context import ServerCallContext

# 1. Define the message
message = Message(
    role="ROLE_USER",
    message_id="local-test-message-id",
    parts=[Part(text="What is the exchange rate from USD to EUR today?")]
)

# 2. Construct the request
request = SendMessageRequest(message=message)

# 3. Construct context
context = ServerCallContext()

# 4. Call the agent
send_message_response = await a2a_agent.on_message_send(request=request, context=context)

print(send_message_response)

测试 on_get_task

以下代码会检索任务的状态和结果。输出显示任务已完成,并包含“Hello World”响应制品。

from a2a.types import GetTaskRequest

# 1. Provide the task_id from the previous step.
# In a real application, you would store and retrieve this ID.
task_id_to_get = send_message_response.id

# 2. Construct the request
request = GetTaskRequest(id=task_id_to_get)

# 3. Call the agent's handler to get the task status.
# Reusing the context constructed in the previous step
task_status_response = await a2a_agent.on_get_task(request=request, context=context)

print(f"Successfully retrieved status for Task ID: {task_id_to_get}")
print("\nFull task status response:")
print(task_status_response)

后续步骤

指南

了解根据您的开发需求,在 Agent Platform Runtime 上部署代理的五种方式。

指南

将 Agent2Agent 智能体与 Agent Platform Runtime 搭配使用。

指南

创建并部署基本智能体,然后使用 Gen AI Evaluation Service 评估该智能体

问题排查

了解如何解决创建自定义代理时出现的常见错误。

资源

查找 Google Agent Platform 的相关资源和支持。

资源

在 GitHub 上探索 Python 中的 Agent2Agent 示例。