Mastra框架实战:构建生产级多智能体协作系统
1. 项目概述:为什么是Mastra?
最近在探索AI Agent开发,从基础的LangChain、AutoGPT一路摸过来,发现了一个挺有意思的框架——Mastra。它不像一些“大而全”的框架那样试图包办一切,而是把重点放在了“多智能体协作”和“生产级部署”上。简单来说,如果你想让多个AI Agent像一支训练有素的团队一样,各司其职又紧密配合,共同完成一个复杂的任务,并且你还希望这个系统能稳定、高效地跑在真实的生产环境里,那Mastra就值得你花时间研究一下。
我最初接触它,是因为想做一个智能客服的升级版。传统的单Agent客服,遇到稍微复杂点的问题,比如用户想“查询订单状态,然后根据物流信息推荐一个类似商品,最后再问问有没有优惠券”,一个Agent就容易手忙脚乱,逻辑链一长就容易出错或遗忘。而Mastra的思路是:拆!让一个Agent专门负责理解用户意图和拆解任务(Orchestrator),一个Agent去查数据库(DB_Query),一个Agent去调用商品推荐API(Recommender),最后再有一个Agent负责组织回复语言(Response_Formatter)。它们之间通过清晰的消息传递机制协作,整个流程可控、可观测,也更容易调试和优化。
所以,这篇笔记就围绕“Mastra的基本用法”展开。我不会只停留在API调用层面,而是会结合我搭建那个多Agent客服系统的实际经历,拆解Mastra的核心概念、工作流设计、代码实现细节,以及过程中踩过的那些坑。无论你是想入门多Agent系统,还是正在为现有单Agent应用的复杂度和稳定性发愁,希望这些实战经验能给你一些直接的参考。
2. 核心概念与架构拆解
在动手写代码之前,必须得先理解Mastra是怎么想问题的。它有几个核心的构建块(Building Blocks),理解了它们,就等于看懂了这张设计蓝图。
2.1 智能体(Agent):不再是孤胆英雄
在Mastra里,Agent是一个具有特定角色、能力和记忆的独立单元。它和我们常说的“LLM调用”不太一样。一个Mastra Agent通常由几个部分组成:
- 角色指令(Instruction): 用自然语言清晰定义这个Agent是谁、负责什么。比如:“你是一个专业的订单查询助手,只负责从数据库中准确、快速地检索订单信息,不回答与订单无关的问题。”
- 模型(Model): 背后驱动的LLM,比如GPT-4、Claude 3或者开源的Llama 3。Mastra支持灵活配置。
- 工具(Tools): Agent可以调用的函数,比如
search_database(order_id)、call_recommendation_api(product_category)。这是Agent与外部世界交互的手和脚。 - 记忆(Memory): 分为会话记忆(记住当前对话上下文)和长期记忆(可以向量化存储和检索的历史信息)。这保证了Agent不是“金鱼脑”。
关键点: Mastra鼓励你创建小而专的Agent。一个Agent最好只做好一件事。这符合软件工程的“单一职责原则”,使得每个Agent都更容易开发、测试和维护。在我的客服系统里,就有IntentClassifier(意图分类)、OrderLookup(订单查询)、FaqAnswer(常见问答)等多个专职Agent。
2.2 工作流(Workflow):定义团队的协作剧本
多个Agent如何一起工作?靠Workflow来编排。你可以把Workflow看作一个有向无环图(DAG),节点是Agent(或人工判断节点),边定义了数据流动的方向和条件。
Mastra提供了几种方式来定义Workflow:
- 顺序流(Sequential): 最简单的A->B->C,上一个Agent的输出是下一个Agent的输入。
- 条件流(Conditional): 根据某个Agent的输出结果,决定下一步走哪个分支。比如,
IntentClassifier输出是“查询订单”,则路由到OrderLookup;如果是“投诉”,则路由到ComplaintHandler。 - 并行流(Parallel): 多个Agent同时执行,最后汇总结果。比如,用户问“A产品和B产品哪个好?”,可以同时启动
ProductAAgent和ProductBAgent去收集信息。
实操心得: 在设计Workflow时,画图!先用纸笔或绘图工具把Agent之间的交互图画出来。明确每个环节的输入输出是什么,异常情况(比如某个Agent调用工具失败)该如何处理(是重试、降级还是转人工)。这一步的思考深度,直接决定了后续开发是否顺畅。
2.3 状态(State)与消息(Message):协作的基石
这是Mastra框架里非常关键的两个概念,保证了Agent之间能有效沟通。
- 状态(State): 可以理解为整个工作流执行过程中的“共享白板”或“上下文容器”。它是一个字典(dict)结构,随着工作流执行而不断更新。比如,
State里可能包含{“user_input”: “我的订单123456到哪里了?”, “parsed_intent”: “query_order”, “order_id”: “123456”, “order_details”: {…}}。每个Agent都可以读取和修改State中自己负责的部分。 - 消息(Message): Agent之间、Agent与用户之间通信的基本单位。通常包含
role(如user,assistant,system)和content。在Mastra中,Workflow会管理一个Message列表,记录完整的对话历史,每个Agent都能看到相关的历史消息来决定自己的行动。
为什么这么设计?这种State+Message的模式,将结构化数据(State)和非结构化对话(Message)分开了。工具调用的结果、提取的实体等结构化数据放在State里,便于程序处理;而自然的语言交互则放在Message里,供LLM理解。这种分离让系统更清晰、更健壮。
3. 环境搭建与第一个“Hello Agent”
理论说得再多,不如跑通一行代码。我们从一个最简单的例子开始,感受一下Mastra的脉搏。
3.1 安装与初始化
Mastra是一个Python框架,安装很简单。建议使用虚拟环境。
# 创建并激活虚拟环境(可选但推荐) python -m venv mastra-env source mastra-env/bin/activate # Linux/Mac # mastra-env\Scripts\activate # Windows # 安装Mastra pip install mastra安装时,它会自动安装一些核心依赖。如果你需要用到特定的功能,比如向量数据库记忆、更复杂的工具包,可能还需要额外安装对应的库(如pymongo,chromadb等)。
3.2 构建你的第一个智能体:翻译官
我们来创建一个简单的翻译Agent,它可以将中文翻译成英文。
import asyncio from mastra.agents import Agent from mastra.models import OpenAIModel # 假设使用OpenAI from mastra.memories import InMemoryMemory # 使用简单的内存记忆 from dotenv import load_dotenv import os # 1. 加载环境变量(你的OPENAI_API_KEY) load_dotenv() # 2. 定义模型 model = OpenAIModel( model="gpt-4o-mini", # 或 "gpt-4", "gpt-3.5-turbo" api_key=os.getenv("OPENAI_API_KEY") ) # 3. 创建Agent translator_agent = Agent( name="Translator", instruction=""" 你是一名专业的翻译官。你的任务是将用户输入的中文文本,准确、流畅地翻译成英文。 只输出翻译后的英文内容,不要添加任何解释或额外信息。 """, model=model, memory=InMemoryMemory() # 使用内存记忆,会话结束后消失 ) # 4. 运行Agent async def main(): user_input = "今天天气真好,我们一起去公园散步吧。" response = await translator_agent.run(task=user_input) print(f"用户输入: {user_input}") print(f"翻译结果: {response.final_output}") if __name__ == "__main__": asyncio.run(main())代码解读与注意事项:
- 异步(Async): Mastra的核心API是异步的(使用
async/await),这是为了高效处理可能并发的Agent调用或网络IO。所以主入口需要用asyncio.run()。 - Agent.run(): 这是触发Agent工作的主要方法。它接收一个
task(字符串),Agent会根据其instruction和memory来处理这个任务。 - response对象:
run()方法返回的是一个复杂的响应对象,其中final_output包含了Agent返回的主要文本内容。你还可以从中查看中间步骤、使用的工具等,这对调试非常有帮助。 - 记忆(Memory): 这里用了
InMemoryMemory,它只存在于程序运行时。如果你希望Agent能记住跨会话的信息,需要配置更持久的记忆存储,比如向量数据库。
运行这段代码,你应该能得到类似"The weather is so nice today, let's go for a walk in the park together."的输出。恭喜,你的第一个Mastra Agent已经跑起来了!但这只是一个单兵作战,接下来我们让它学会使用工具。
4. 为智能体装备工具(Tools)
没有工具的Agent,就像没有手脚的专家,空有知识却无法行动。让翻译官不仅能翻译,还能在翻译前先查询一下某个专业术语的解释。
4.1 创建自定义工具
假设我们有一个简单的函数,可以“查询术语表”(这里模拟为本地字典)。
from mastra.tools import tool from typing import Dict, Any # 模拟一个简单的术语表 TERM_GLOSSARY = { "碳中和": "Carbon neutrality", "元宇宙": "Metaverse", "内卷": "Involution", "躺平": "Lie flat" } # 使用 @tool 装饰器将普通函数转换为Mastra可识别的工具 @tool def lookup_glossary(term: str) -> Dict[str, Any]: """ 查询术语表中对应术语的英文翻译。 Args: term: 要查询的中文术语。 Returns: 一个字典,包含‘found’布尔值和‘translation’翻译结果(如果找到)。 """ translation = TERM_GLOSSARY.get(term) if translation: return {"found": True, "translation": translation} else: return {"found": False, "translation": None} # 创建新的、装备了工具的翻译Agent translator_with_tools = Agent( name="TranslatorWithTools", instruction=""" 你是一名专业的翻译官,尤其擅长翻译包含专业术语的文本。 你的工作流程是: 1. 分析用户输入的中文文本。 2. 识别文本中可能存在的专业术语(如‘碳中和’、‘元宇宙’)。 3. 对于识别出的每个术语,使用‘lookup_glossary’工具查询其标准英文翻译。 4. 综合工具查询结果和你的知识,将整段文本准确、流畅地翻译成英文。 5. 如果工具没有查到某个术语,则依靠你自己的理解进行翻译。 输出最终翻译结果。 """, model=model, tools=[lookup_glossary], # 将工具列表传给Agent memory=InMemoryMemory() )4.2 观察工具的使用
现在运行这个升级版的Agent:
async def main_with_tools(): # 测试1:包含已知术语 test_text_1 = "在追求碳中和的背景下,许多企业感到内卷。" response1 = await translator_with_tools.run(task=test_text_1) print(f"测试1 - 输入: {test_text_1}") print(f"测试1 - 输出: {response1.final_output}") # 可以查看Agent思考过程(如果模型支持) # print(response1.messages) # 会看到它调用工具的请求和结果 # 测试2:包含未知术语 test_text_2 = "今天数字孪生技术发展很快。" response2 = await translator_with_tools.run(task=test_text_2) print(f"\n测试2 - 输入: {test_text_2}") print(f"测试2 - 输出: {response2.final_output}") if __name__ == "__main__": asyncio.run(main_with_tools())你会观察到:对于测试1,Agent在输出最终翻译前,会先自动调用lookup_glossary工具去查询“碳中和”和“内卷”。你可以在response.messages或Mastra的日志中看到类似[Tool Call] lookup_glossary(term='碳中和')和[Tool Result] {'found': True, 'translation': 'Carbon neutrality'}的记录。对于测试2,“数字孪生”不在我们的术语表中,Agent会依赖自己的知识进行翻译。
注意:工具的描述(Docstring)非常重要!LLM(尤其是GPT-4这类模型)主要依靠函数名和描述来决定何时、如何使用工具。务必清晰、准确地描述工具的功能、参数和返回值。
5. 构建多智能体工作流(Workflow)
现在,让我们进入Mastra最精彩的部分:让多个Agent协作。我们构建一个简单的“研究助手”工作流,包含两个Agent:一个Researcher负责搜索网络信息,一个Writer负责整理成报告。
5.1 定义参与协作的智能体
首先,定义两个Agent。为了简化,我们用一个模拟的搜索工具。
from mastra.agents import Agent from mastra.tools import tool # 模拟一个网络搜索工具 @tool def web_search(query: str) -> str: """ 根据查询词模拟搜索网络信息。返回模拟的搜索结果摘要。 """ # 这里只是一个模拟,真实情况可以接入SerperAPI、Google Search API等 mock_data = { "大语言模型的发展": "大语言模型(LLM)如GPT-4在自然语言理解和生成上取得突破,推动了AI普及。其发展依赖于算力增长和海量数据训练。", "Python编程入门": "Python是一种解释型、高级别的通用编程语言,以语法简洁清晰著称,广泛应用于Web开发、数据分析、人工智能等领域。" } return mock_data.get(query, f"未找到关于'{query}'的详细信息。") # 研究员Agent researcher = Agent( name="Researcher", instruction=""" 你是一个网络研究员。根据用户提出的主题,使用‘web_search’工具查找相关信息。 你需要从搜索结果中提取关键事实、数据和观点,并以清晰、有条理的要点形式总结出来。 将总结好的信息提供给报告撰写员。 """, model=model, tools=[web_search], memory=InMemoryMemory() ) # 报告撰写员Agent writer = Agent( name="Writer", instruction=""" 你是一名技术文档撰写员。你将收到研究员提供的关于某个主题的研究摘要。 你的任务是将这些摘要信息,组织成一篇结构完整、语言流畅、易于理解的简短报告。 报告应包含简介、核心内容分点阐述和简要总结。 直接输出报告正文,不需要标题。 """, model=model, # 这个Agent不需要工具 memory=InMemoryMemory() )5.2 设计并实现顺序工作流
现在,我们用最基础的SequentialWorkflow把这两个Agent串联起来。
from mastra.workflows import SequentialWorkflow from mastra import State # 定义一个顺序工作流 research_and_write_workflow = SequentialWorkflow( name="ResearchAndWrite", agents=[researcher, writer] # 按顺序执行 ) async def run_workflow(): user_topic = "大语言模型的发展" print(f"用户查询主题: {user_topic}") # 初始化工作流状态,将用户输入放入State initial_state = State(user_query=user_topic) # 运行工作流 final_state = await research_and_write_workflow.run(initial_state) # 从最终State中获取结果 # 通常,最后一个Agent(Writer)的输出会被放在State的某个字段或final_output中 # 具体取决于Workflow和Agent的配置。这里我们假设最终报告在State的‘final_report’里。 # 我们需要查看Workflow是如何传递数据的。一个常见模式是每个Agent的输出会更新State。 # 为了简单演示,我们直接打印最终Agent(Writer)的响应消息。 # 在实际复杂Workflow中,你需要设计清晰的状态传递路径。 print("\n--- 工作流执行报告 ---") # 查看最终状态的所有键 print(f"最终State内容: {final_state.data}") # 更实际的做法:在Workflow定义中指定输入输出映射。 # 下面展示一种更可控的方式:使用自定义的Workflow步骤函数。 if __name__ == "__main__": asyncio.run(run_workflow())运行上面的代码,你可能会发现一个问题:Researcher的总结输出,如何自动成为Writer的输入?在简单的SequentialWorkflow中,默认会将上一个Agent的final_output作为下一个Agent的task输入。但为了更精细的控制,Mastra提供了更强大的方式:使用Workflow类并自定义steps。
5.3 实现可控的协作流
让我们用更底层、更灵活的方式来构建这个工作流,明确控制数据的流动。
from mastra.workflows import Workflow from mastra import State class ResearchWriteWorkflow(Workflow): """ 自定义研究-撰写工作流。 1. Researcher接收用户查询,进行搜索并总结。 2. Writer接收Researcher的总结,撰写报告。 """ def __init__(self): super().__init__(name="ControlledResearchWrite") self.researcher = researcher # 使用之前定义的Agent self.writer = writer async def run(self, initial_state: State) -> State: """ 执行工作流步骤。 """ state = initial_state # 步骤1:研究员执行研究 print(f"[Workflow] 步骤1: Researcher开始研究主题‘{state.data.get('user_query')}‘...") research_result = await self.researcher.run(task=state.data['user_query']) # 将研究员的结果存入state state.update(research_summary=research_result.final_output) print(f"[Workflow] Researcher总结完成。") # 步骤2:撰写员基于总结撰写报告 print(f"[Workflow] 步骤2: Writer开始撰写报告...") # 为Writer构造任务,可以包含更多上下文 writer_task = f"请根据以下研究摘要,撰写一份简短的技术报告:\n\n{state.data['research_summary']}" report_result = await self.writer.run(task=writer_task) state.update(final_report=report_result.final_output) print(f"[Workflow] Writer报告完成。") return state # 使用自定义工作流 async def main_custom_workflow(): workflow = ResearchWriteWorkflow() initial_state = State(user_query="Python编程入门") final_state = await workflow.run(initial_state) print("\n" + "="*50) print("最终生成的报告:") print("="*50) print(final_state.data.get('final_report', '报告生成失败。')) if __name__ == "__main__": asyncio.run(main_custom_workflow())通过这种自定义Workflow类的方式,我们完全掌控了执行流程和状态管理。你可以清晰地看到:
Researcher接收user_query并运行。- 将其输出
research_summary存入state。 Writer读取state中的research_summary,并以此为基础任务运行。- 最终报告
final_report也被存入state。
这是构建复杂多Agent系统的推荐模式。你可以在此基础上轻松地添加条件判断、并行执行、循环等逻辑。
6. 状态管理、记忆与持久化
在真实应用中,Agent系统往往需要处理多轮对话,并记住关键信息。这就需要有效的状态管理和记忆持久化。
6.1 深入理解State的数据流
State是Workflow的“中央数据总线”。好的状态设计能让工作流清晰易懂。
- 状态结构规划: 在开始编码前,规划好你的
State字典里应该有哪些键。例如,对于一个客服工作流,可能有:session_id,user_input,parsed_intent,extracted_entities(如order_id,product_name),current_step,agent_responses(列表),final_answer等。 - 状态更新策略: 是覆盖还是追加?比如
agent_responses通常用列表追加,而current_step则是覆盖。在自定义Workflow的run方法中,通过state.update()或直接修改state.data来管理。 - 状态序列化: 如果Workflow执行时间很长或需要暂停恢复,你需要能将
State对象序列化成JSON等格式存储起来,下次再反序列化加载。Mastra的State基类通常支持dict转换。
6.2 为智能体配置长期记忆
之前的例子用了InMemoryMemory,这是会话级记忆。要实现跨会话记忆,需要使用向量数据库(Vector Store)作为记忆后端。
以下是一个使用ChromaDB(轻量级向量数据库)作为记忆存储的示例:
from mastra.memories import VectorStoreMemory from mastra.embeddings import OpenAIEmbeddings import chromadb from chromadb.config import Settings # 1. 初始化嵌入模型(用于将文本转换为向量) embedding_model = OpenAIEmbeddings( model="text-embedding-3-small", api_key=os.getenv("OPENAI_API_KEY") ) # 2. 初始化Chroma客户端(持久化模式) chroma_client = chromadb.PersistentClient(path="./chroma_memory_db") # 3. 创建向量存储记忆 vector_memory = VectorStoreMemory( embedding_model=embedding_model, vector_store_client=chroma_client, collection_name="agent_long_term_memory", # 指定集合名 search_kwargs={"k": 5} # 每次检索最相关的5条记忆 ) # 4. 创建带有长期记忆的Agent personal_assistant = Agent( name="PersonalAssistant", instruction=""" 你是用户的个人助理。你需要记住关于用户的个人信息和偏好,并在对话中自然地运用这些信息。 例如,如果用户告诉过你他喜欢咖啡,那么当他问‘今天有什么推荐?’时,你可以联想到咖啡。 """, model=model, memory=vector_memory # 使用向量记忆 ) async def main_with_memory(): # 第一次对话:告诉助理一个信息 print("【第一轮对话】") response1 = await personal_assistant.run(task="记住,我最喜欢的水果是芒果。") print(f"用户: 记住,我最喜欢的水果是芒果。") print(f"助理: {response1.final_output}") # 此时,“最喜欢的水果是芒果”这个信息会被存储到向量数据库中。 # 模拟一点时间间隔或新会话 print("\n... 一段时间后 ...\n") # 第二次对话:询问相关推荐 print("【第二轮对话】") # 在运行前,记忆系统会自动从向量库中检索与当前任务相关的历史记忆,并注入到对话上下文中。 response2 = await personal_assistant.run(task="今天下午茶有什么建议吗?") print(f"用户: 今天下午茶有什么建议吗?") print(f"助理: {response2.final_output}") # 理想的回答应该包含:“既然你喜欢芒果,可以试试芒果布丁或芒果冰沙。” if __name__ == "__main__": asyncio.run(main_with_memory())关键点:
- 记忆检索是自动的: 当Agent运行时,其配置的
memory对象会基于当前对话上下文(或任务),自动从向量库中检索相关的历史记忆片段,并将它们作为系统提示词或上下文的一部分提供给LLM。 - 记忆存储也是自动的: 默认情况下,Agent的输入和输出会被自动存储到记忆中。你也可以手动控制存储哪些内容。
- 记忆的关联性: 向量检索的核心是“相关性”,而不是“时间顺序”。这保证了Agent能想起最相关的事情,而不是最近的事情。
7. 生产级考量与部署实践
让Agent在本地跑通只是第一步,要真正投入使用,必须考虑生产环境的要求。
7.1 错误处理与鲁棒性
多Agent系统链路长,任何一个环节出错(LLM API超时、工具调用异常、网络问题)都可能导致整个流程失败。必须添加健壮的错误处理。
import asyncio import logging from mastra.workflows import Workflow from mastra import State from openai import APITimeoutError, RateLimitError logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class RobustResearchWorkflow(Workflow): async def run(self, initial_state: State) -> State: state = initial_state max_retries = 3 # 步骤1:研究员Agent(带重试机制) for attempt in range(max_retries): try: research_result = await asyncio.wait_for( self.researcher.run(task=state.data['user_query']), timeout=30.0 # 设置超时 ) state.update(research_summary=research_result.final_output) break # 成功则跳出重试循环 except (APITimeoutError, asyncio.TimeoutError) as e: logger.warning(f"Researcher调用超时,第{attempt+1}次重试...") if attempt == max_retries - 1: state.update(research_summary="【信息】网络查询暂时不可用,将基于现有知识生成报告。") logger.error("Researcher最终失败,使用降级方案。") except RateLimitError as e: logger.error(f"触发速率限制: {e}") # 可以在这里添加指数退避等待 await asyncio.sleep(2 ** attempt) if attempt == max_retries - 1: raise # 重试后仍失败,向上抛出 except Exception as e: logger.exception(f"Researcher执行发生未知错误: {e}") state.update(research_summary="【信息】研究阶段遇到问题。") break # 非重试性错误,直接退出 # 步骤2:撰写员Agent try: if state.data.get('research_summary'): writer_task = f"基于以下信息撰写报告:{state.data['research_summary']}" else: writer_task = f"请直接为‘{state.data['user_query']}‘这个主题撰写一份简介报告。" report_result = await self.writer.run(task=writer_task) state.update(final_report=report_result.final_output) except Exception as e: logger.exception(f"Writer执行失败: {e}") state.update(final_report="抱歉,报告生成过程出现异常。") return state要点:
- 重试机制: 对于暂时的网络或API错误(超时、限流)进行有限次重试。
- 超时控制: 使用
asyncio.wait_for防止某个Agent调用无限期挂起。 - 降级方案: 当核心步骤(如研究)失败时,提供备选路径(如让撰写员直接基于知识生成),保证系统仍有输出,而非完全崩溃。
- 日志记录: 详细的日志是排查生产问题的生命线。记录关键步骤、输入输出和异常。
7.2 性能监控与成本控制
- Token消耗监控: 每次LLM调用都会消耗Token。在生产环境中,需要记录每个Agent每次调用的输入/输出Token数,并关联到具体用户或会话,以便进行成本分析和优化。可以在自定义Agent或Workflow的
run方法前后拦截请求和响应,进行计算和上报。 - 延迟监控: 记录每个Agent和工作流整体的执行时间,设立性能基线。慢的环节可能是优化重点(比如是否检索了过多记忆导致上下文过长)。
- 缓存策略: 对于频繁出现的、结果固定的查询(如某些FAQ),可以在工具层或Agent层加入缓存(如Redis),避免重复调用LLM或外部API,显著降低成本和延迟。
7.3 部署模式
- 异步Web服务: 使用
FastAPI或Django Ninja等异步Web框架,将你的Mastra Workflow封装成API端点。这是最常见的部署方式。from fastapi import FastAPI, HTTPException from pydantic import BaseModel import uvicorn app = FastAPI() # 初始化你的工作流 workflow = ResearchWriteWorkflow() class QueryRequest(BaseModel): topic: str @app.post("/research-report/") async def generate_report(request: QueryRequest): try: initial_state = State(user_query=request.topic) final_state = await workflow.run(initial_state) return {"report": final_state.data.get('final_report')} except Exception as e: logger.error(f"API处理失败: {e}") raise HTTPException(status_code=500, detail="报告生成失败") # if __name__ == "__main__": # uvicorn.run(app, host="0.0.0.0", port=8000) - 队列与后台任务: 对于耗时较长的复杂工作流,不要同步阻塞HTTP请求。可以将任务推入消息队列(如Celery + Redis/RabbitMQ,或Dramatiq),由后台Worker异步处理,并通过WebSocket或轮询接口向客户端返回结果。
- 容器化: 使用Docker将你的Agent应用及其依赖打包成镜像,确保环境一致性。结合Kubernetes或Docker Compose进行编排和管理。
8. 常见问题与调试技巧实录
在开发和调试Mastra应用时,我遇到了不少典型问题,这里汇总一下,希望能帮你快速排雷。
8.1 Agent不调用工具
- 问题现象: 明明给Agent配置了工具,但它总是忽略,选择自己“脑补”回答。
- 排查思路:
- 检查工具描述: 这是最常见的原因。工具的函数名和文档字符串(Docstring)必须清晰、准确。LLM根据这些描述来决定是否调用。尝试将描述写得更加具体、场景化,例如“必须使用此工具来查询XXX信息,否则无法回答用户问题”。
- 检查模型能力: 有些较弱的模型(如早期的GPT-3.5-turbo)工具调用能力不稳定。升级到GPT-4、Claude 3或最新的高性能模型。
- 提供示例(Few-Shot): 在Agent的
instruction中,加入一两个工具调用的示例,演示在什么情况下应该调用工具以及如何解析结果。 - 开启详细日志: 查看Agent思考的中间过程(如果模型支持并开启了相关设置),看它是否在考虑工具但最终放弃了,这有助于理解其“思维”过程。
8.2 工作流状态传递混乱
- 问题现象: 后一个Agent拿不到前一个Agent的数据,或者数据格式不对。
- 解决方案:
- 设计清晰的状态Schema: 在Workflow类开头用注释明确写出
State中每个字段的名称、类型和含义。 - 使用类型提示和验证: 考虑使用Pydantic模型来定义
State的数据结构,这样可以在运行时进行验证,提前发现字段缺失或类型错误。 - 打印调试: 在每个Agent执行前后,打印
state.data的内容,确认数据是否正确写入和读取。 - 避免直接操作
state.data: 尽量使用state.get()、state.update()等方法,它们比直接操作字典更安全,有时框架会在此基础上添加额外逻辑。
- 设计清晰的状态Schema: 在Workflow类开头用注释明确写出
8.3 记忆检索不准确或无关
- 问题现象: Agent总是回忆起不相关的历史对话,干扰当前回答。
- 优化方法:
- 调整检索数量(k): 减少
search_kwargs中的k值(比如从5调到3),只取最相关的几条记忆。 - 优化记忆存储内容: 不要存储所有对话回合。只存储重要的、需要被长期记住的事实(如用户偏好、关键决策)。可以在存储前对文本进行摘要。
- 使用元数据过滤: 高级的向量存储(如Chroma、Weaviate)支持为每条记忆附加元数据(如
session_id,topic,importance)。在检索时,可以结合元数据过滤,例如只检索同一topic下的记忆。 - 提升嵌入模型质量: 使用更强大的嵌入模型(如
text-embedding-3-large)能获得更好的语义表示,从而提高检索相关性。
- 调整检索数量(k): 减少
8.4 系统响应速度慢
- 瓶颈分析:
- LLM API延迟: 这是主要瓶颈。监控每个Agent调用的耗时。
- 工具调用延迟: 如果工具需要访问外部API或数据库,可能很慢。
- 记忆检索延迟: 向量检索在数据量大时可能变慢。
- 顺序执行: 如果Workflow是严格的A->B->C顺序,总耗时就是各步骤之和。
- 优化策略:
- 并行化: 将没有依赖关系的Agent改为并行执行。Mastra支持
ParallelWorkflow或在自定义Workflow中使用asyncio.gather。 - 缓存: 对LLM响应(针对固定输入)和工具结果实施缓存。
- 流式输出: 如果最终输出是文本,可以考虑使用LLM的流式响应,让用户边生成边看到部分结果,提升体验感。
- 精简上下文: 检查每次调用LLM时传入的上下文(历史消息、记忆)是否过长。移除无关历史,对必要的历史进行摘要。
- 并行化: 将没有依赖关系的Agent改为并行执行。Mastra支持
8.5 调试与日志记录最佳实践
- 结构化日志: 使用
structlog或logging模块的DictFormatter,输出JSON格式的日志,方便被ELK(Elasticsearch, Logstash, Kibana)等日志系统收集和查询。在日志中记录session_id、agent_name、step、input_snapshot、output_snapshot、duration等关键字段。 - 可视化工作流执行: 对于复杂工作流,可以开发一个简单的可视化工具,将
State的变迁和Agent的调用序列以图表或时间线形式展示出来,这对理解系统行为和排查问题有奇效。 - 单元测试: 为每个独立的
Tool和Agent编写单元测试,模拟输入验证输出。对于Workflow,可以编写集成测试,使用固定的Mock数据来验证整个数据流是否正确。
踩过这些坑之后,我的体会是,开发多Agent系统更像是在设计一个微型社会的运行规则。框架(如Mastra)提供了基础设施和语法,但真正的挑战在于如何清晰地定义角色、设计协作协议、处理异常情况以及保障系统整体的稳定和高效。从一个小而专的Agent开始,逐步构建和测试工作流,持续监控和迭代,是通往成功最实在的路径。
