流式输出 Streaming 流式输出让 AI 的回复像打字一样逐字显示,极大地提升了用户体验。LangChain 的 Agent 内置了完善的流式输出支持。如果使用 invoke(),用户需要等待 Agent 完成所有步骤(多次模型调用 + 工具执行)才能看到结果。对于复杂任务,这可能耗时十几秒甚至更长。流式输出解决了这个问题:每生成一个 Token 就立即返回,用户可以实时看到进展。
stream_mode=”messages”——逐 Token 流式 这是最细粒度的流式模式,每个 chunk 对应一个 Token,示例代码如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 from dotenv import load_dotenvload_dotenv() from langchain.agents import create_agentfrom langchain.chat_models import init_chat_modelfrom langchain.messages import HumanMessagemodel = init_chat_model("deepseek:deepseek-v4-flash" ) agent = create_agent( model=model, system_prompt="你是菜鸟教程 RUNOOB 的助手。" , ) print ("实时流式输出:" )for msg_chunk, metadata in agent.stream( {"messages" : [HumanMessage(content="用一句话介绍菜鸟教程 RUNOOB" )]}, stream_mode="messages" , ): if msg_chunk.content: print (msg_chunk.content, end="" , flush=True ) print ()
stream_mode=”updates”——逐步查看 Agent 执行过程 这个模式在构建需要显示”思考过程”的界面时非常有用,示例代码如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 from langchain.tools import tool@tool def search_course (keyword: str ) -> str : """在菜鸟教程 RUNOOB 搜索课程""" courses = { "python" : "Python3 基础教程(30章,20小时)" , "html" : "HTML 基础教程(25章,15小时)" , } return courses.get(keyword.lower(), "未找到相关课程" ) agent = create_agent( model=init_chat_model("deepseek:deepseek-v4-flash" , temperature=0 ), tools=[search_course], system_prompt="你是菜鸟教程 RUNOOB 的课程顾问。" , ) print ("=== Agent 执行过程 ===\n" )for chunk in agent.stream( {"messages" : [HumanMessage(content="帮我查一下 Python 课程" )]}, stream_mode="updates" , ): for node_name, update in chunk.items(): print (f"[{node_name} ]" , end=" " ) if "messages" in update: for msg in update["messages" ]: if msg.type == "ai" : if hasattr (msg, 'tool_calls' ) and msg.tool_calls: calls = [tc['name' ] for tc in msg.tool_calls] print (f"请求调用: {calls} " ) elif msg.content: print (f"回复: {msg.content[:80 ]} " ) elif msg.type == "tool" : print (f"工具返回 [{msg.name} ]: {msg.content} " )
结构化输出 大多数时候,我们需要的不是一段文本,而是结构化的数据,例如 Json 对象。这样做可以省去从文本中解析数据这一步,让 AI 的输出可以直接被程序使用。接下来我们学习几种方法来实现结构化输出。
传入 Pydantic 模型 最简单的做法是将 Pydantic 模型传给 response_format 参数即可,示例代码如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 from dotenv import load_dotenvload_dotenv() from pydantic import BaseModel, Fieldfrom langchain.agents import create_agentfrom langchain.chat_models import init_chat_modelfrom langchain.messages import HumanMessageclass CourseInfo (BaseModel ): """菜鸟教程 RUNOOB 课程提取结果""" course_name: str = Field(description="课程名称" ) difficulty: str = Field(description="难度:入门/进阶/高级" ) estimated_hours: int = Field(description="预计学习时长(小时)" ) is_free: bool = Field(description="是否免费" ) model = init_chat_model("deepseek:deepseek-v4-flash" , temperature=0 ) agent = create_agent( model=model, response_format=CourseInfo, system_prompt="你是菜鸟教程 RUNOOB 的课程助手,从用户描述中提取课程信息。" , ) result = agent.invoke({ "messages" : [HumanMessage( content="我最近在学习 Python3 基础教程,是入门级别的," "大概要学 20 个小时,而且是完全免费的" )] }) if "structured_response" in result: course = result["structured_response" ] print (f"课程名: {course.course_name} " ) print (f"难度: {course.difficulty} " ) print (f"预计时长: {course.estimated_hours} 小时" ) print (f"免费: {'是' if course.is_free else '否' } " ) print (f"对象类型: {type (course)} " )
注意,模型返回的 structured_response 是 Pydantic 模型实例,而不是普通字典,这意味着你可以使用 .course_name 等属性访问。response_format 和 tools 可以同时使用——Agent 在需要时调用工具,最终输出结构化数据。
同时 Pydantic 还支持复杂的嵌套结构,如下:
1 2 3 4 5 6 class LearningPlan (BaseModel ): """学习计划""" goal: str = Field(description="学习目标概述" ) level: Literal ["入门" , "进阶" , "高级" ] = Field(description="难度级别" ) total_hours: float = Field(description="总时长(小时)" ) topics: list [Topic] = Field(description="知识点列表" )
LangChain 中间件(Middleware) LangChain Middleware(中间件)是 LangChain 最强大的特性。它让你在 Agent 执行的各个环节插入自定义逻辑,实现重试、降级、缓存、内容过滤、日志记录等功能——而不需要修改 Agent 本身的代码。
Middleware 是 Agent 执行流程中的钩子(Hook)。每个购自让你在特定的时间点执行自定义代码。如何更加直观的理解 Middleware?假设 Agent 的执行流程是这样的:1. 用户输入 → 2. 模型思考 → 3. 可能调用工具 → 4. 模型再思考 → 5. 输出结果。Middleware 可以让你在这 5 个环节中插入自定义逻辑:
1 2 3 4 5 6 7 8 9 10 11 1. 用户输入 ↓ [before_agent 钩子:日志记录、权限检查] 2. 模型思考 ↓ [before_model 钩子:消息预处理] ↓ [wrap_model_call 钩子:重试、降级、缓存] ↓ [after_model 钩子:内容审核] 3. 工具执行 ↓ [wrap_tool_call 钩子:工具调用重试] 4. 回到模型思考(循环直到完成) ↓ [after_agent 钩子:结果格式化、统计分析] 5. 输出结果
LangChain 的 Middleware 提供了 6 个钩子,按执行时机分为两类:
钩子
执行频率
执行位置
主要用途
before_agent
一次
Agent 开始前
初始化、权限检查、输入预处理
before_model
每次循环
模型调用前
消息预处理,动态上下文注入
wrap_model_call
每次循环
包裹模型调用
重试、降级、缓存、请求改写
after_model
每次循环
模型调用后
内容审核、相应过滤、日志
wrap_tool_call
每次工具调用
包裹工具执行
工具重试、结果缓存、参数改写
after_agent
一次
Agent 结束后
格式化输出、统计、清理资源
Middleware 可以通过类继承或装饰器两种方法使用。 方式一:装饰器,示例代码如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 from langchain.agents.middleware import before_model, after_model@before_model def log_before (state, runtime ): """在每次模型调用前记录日志""" msg_count = len (state.get("messages" , [])) print (f"[before_model] 当前消息数: {msg_count} " ) return None @after_model def log_after (state, runtime ): """在每次模型调用后记录日志""" last_msg = state["messages" ][-1 ] if state.get("messages" ) else None if last_msg and hasattr (last_msg, 'tool_calls' ) and last_msg.tool_calls: print (f"[after_model] 模型请求了工具调用" ) return None
方式二:类继承,适合复杂逻辑,示例代码如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 from langchain.agents.middleware import AgentMiddlewareclass LoggingMiddleware (AgentMiddleware ): """自定义日志中间件""" @property def name (self ) -> str : return "logging" def before_agent (self, state, runtime ): """Agent 开始前的逻辑""" print ("[Logging] Agent 开始执行" ) return None def before_model (self, state, runtime ): """模型调用前的逻辑""" msg_count = len (state.get("messages" , [])) print (f"[Logging] 准备调用模型,当前 {msg_count} 条消息" ) return None def after_model (self, state, runtime ): """模型调用后的逻辑""" print ("[Logging] 模型调用完成" ) return None def after_agent (self, state, runtime ): """Agent 结束后的逻辑""" print ("[Logging] Agent 执行结束" ) return None
@before_model 与 @after_model @before_model 在每次调用模型前执行,可以在这里修改消息、注入上下文条件、或者直接跳过模型调用,下面是一个裁剪消息的示例:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 @before_model def limit_context (state, runtime ): """限制消息历史长度,防止上下文过长""" messages = state.get("messages" , []) MAX_MESSAGES = 6 if len (messages) > MAX_MESSAGES: trimmed = messages[-MAX_MESSAGES:] if trimmed and trimmed[0 ].type != "human" : trimmed = trimmed[1 :] return {"messages" : trimmed} return None
还可以使用 before_model 来屏蔽敏感词。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 @before_model def content_filter (state, runtime ): """检查用户消息是否包含敏感词,如果包含则拦截""" messages = state.get("messages" , []) if not messages: return None last_msg = messages[-1 ] content = str (last_msg.content) if hasattr (last_msg, 'content' ) else "" for word in SENSITIVE_WORDS: if word in content: print (f"[拦截] 检测到敏感词: {word} " ) return { "jump_to" : "end" , "messages" : [ HumanMessage(content=f"抱歉,为了安全,不能处理包含「{word} 」的请求。" ) ] } return None
@after_model 在模型回复后执行,适合审核模型输出、提取关键信息、追加后续指令等。有如下两个场景,分别问响应内容审核以及自动追加提示信息:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 @after_model def response_audit (state, runtime ): """审核模型回复,如果涉及禁止话题则替换""" messages = state.get("messages" , []) if not messages: return None last_msg = messages[-1 ] content = str (last_msg.content) if hasattr (last_msg, 'content' ) else "" for topic in FORBIDDEN_TOPICS: if topic in content: runtime.stream_writer({ "type" : "warning" , "message" : f"检测到回复包含「{topic} 」相关内容,已被替换" }) from langchain.messages import AIMessage return { "messages" : [ AIMessage(content="抱歉,我无法回答这个问题。" "请询问编程学习相关的内容。" ) ] } return None @after_model def append_disclaimer (state, runtime ): """在每次模型回复后自动追加免责声明""" messages = state.get("messages" , []) if not messages: return None last_msg = messages[-1 ] if (last_msg.type == "ai" and last_msg.content and not (hasattr (last_msg, 'tool_calls' ) and last_msg.tool_calls)): return { "messages" : [ AIMessage( content=( last_msg.content + "\n\n---\n*以上内容由菜鸟教程 RUNOOB AI 助手生成,仅供参考。*" ) ) ] } return None
在这两个 hook 中,我们还可以通过 can_jump_to 参数和 jump_to 状态来控制 Agent 的流程:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 from langchain.agents.middleware import before_model@before_model(can_jump_to=["end" ] ) def conditional_exit (state, runtime ): """在特定条件下直接结束 Agent""" messages = state.get("messages" , []) if not messages: return None last_content = str (messages[-1 ].content) if last_content.strip() in ["再见" , "拜拜" , "bye" ]: return { "jump_to" : "end" , "messages" : [{"role" : "assistant" , "content" : "再见!期待下次为您服务。" }] } return None
can_jump_to 支持的值如下:
值
含义
适用场景
[“end”]
可跳转到结束
条件退出、安全拦截
[“model”]
可跳转回模型
需要让模型重新处理
[“tools”]
可跳转到工具节点
跳过模型直接执行工具
[“model”,”end”]
可跳转到模型或结束
多种条件分支
@wrap_model_call @wrap_model_call 是中间件中最强大的钩子,它不像 before 和 after 那样只是观察,而是可以完全控制模型的执行过程 —- 重试、降级、缓存、甚至跳过模型直接用预设回复。
@wrap_model_call 的核心是一个 handler 回调函数。调用 handler(request) 才会真正执行模型;不调用则跳过模型。示例如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 @wrap_model_call def my_middleware (request, handler ): print ("模型即将被调用..." ) response = handler(request) print ("模型调用完成" ) return response
接下来是几个 @wrap_model_call 的常见场景: 场景一:模型调用失败重试
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 @wrap_model_call def retry_on_error (request, handler ): """模型调用失败时自动重试,最多 3 次""" max_retries = 3 last_error = None for attempt in range (max_retries): try : result = handler(request) if attempt > 0 : print (f" [重试成功] 第 {attempt + 1 } 次尝试" ) return result except Exception as e: last_error = e if attempt < max_retries - 1 : wait_time = (attempt + 1 ) * 2 print (f" [重试] 第 {attempt + 1 } 次失败,{wait_time} 秒后重试..." ) import time time.sleep(wait_time) raise last_error
场景二、模型降级/故障转移。当主模型不可用时,自动切换到备用模型:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 @wrap_model_call def fallback_on_error (request, handler ): """主模型失败时自动切换到备用模型""" try : return handler(request) except Exception as e: print (f"[降级] 主模型失败: {e} ,切换到备用模型..." ) request = request.override(model=fallback_model) try : return handler(request) except Exception as e2: print (f"[降级] 备用模型也失败了: {e2} " ) return AIMessage( content="抱歉,服务暂时不可用,请稍后重试。" )
场景三、缓存模型响应。对于重复的查询,可以缓存模型响应以减少 API 调用成本,相当于给 Agent 的模型调用装了一个“缓存门卫”——相同问题第二次来时,门卫直接把上次的答案塞回去,根本不让请求走到大模型那里,从而实现零成本、零延迟的重复查询优化:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 @wrap_model_call def cache_responses (request, handler ): """缓存模型响应,相同问题不重复调用""" messages = request.messages if not messages: return handler(request) last_content = str (messages[-1 ].content) if hasattr (messages[-1 ], 'content' ) else "" cache_key = last_content[:200 ] if cache_key in cache: print (f"[缓存命中] 直接返回缓存结果" ) cached = cache[cache_key] return AIMessage(content=f"{cached} \n\n*(来自缓存)*" ) result = handler(request) if hasattr (result, 'content' ): cache[cache_key] = result.content print (f"[缓存未命中] 已存入缓存,当前 {len (cache)} 条" ) elif hasattr (result, 'model_response' ): pass return result
@wrap_tool_call 让你在工具执行层面实现与 @wrap_model_call 类似的控制能力——重试、缓存、参数改写、结果后处理。 @wrap_tool_call 的结构与 @wrap_model_call 类似,接收 request 和 handler 两个参数:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 from langchain.agents.middleware import wrap_tool_call@wrap_tool_call def my_tool_wrapper (request, handler ): result = handler(request) return result
场景一、工具调用失败时自动重试:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 @wrap_tool_call def retry_tool_on_error (request, handler ): """工具调用失败时自动重试""" max_retries = 3 last_result = None for attempt in range (max_retries): try : result = handler(request) if hasattr (result, 'status' ) and result.status == "error" : if attempt < max_retries - 1 : print (f" [重试] 工具返回错误,第 {attempt + 1 } 次重试..." ) continue if attempt > 0 : print (f" [重试成功] 第 {attempt + 1 } 次尝试" ) return result except Exception as e: if attempt < max_retries - 1 : import time time.sleep((attempt + 1 ) * 2 ) print (f" [重试] 异常 {e} ,{attempt + 1 } 次重试..." ) else : raise return last_result
场景二、修改工具参数。在工具执行之前动态修改参数,可以在不修改工具代码的情况下实现参数转换:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 @wrap_tool_call def normalize_city_name (request, handler ): """自动规范化城市名称(全角转半角、去除多余空格等)""" tool_call = request.tool_call if "city" in tool_call.get("args" , {}): city = tool_call["args" ]["city" ] normalized = city.strip().replace(" " , "" ) new_args = {**tool_call["args" ], "city" : normalized} new_tool_call = {**tool_call, "args" : new_args} request = request.override(tool_call=new_tool_call) return handler(request)
场景三、工具结果缓存,对于重复的工具调用(相同工具 + 相同参数),可以缓存结果:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 @wrap_tool_call def cache_tool_results (request, handler ): """缓存工具执行结果""" tool_name = request.tool_call.get("name" , "unknown" ) tool_args = str (request.tool_call.get("args" , {})) cache_key = f"{tool_name} :{tool_args} " if cache_key in tool_cache: print (f"[工具缓存命中] {tool_name} " ) cached_content = tool_cache[cache_key] return ToolMessage( content=cached_content, tool_call_id=request.tool_call.get("id" , "" ), name=tool_name, ) result = handler(request) if hasattr (result, 'content' ): tool_cache[cache_key] = result.content print (f"[工具缓存写入] {tool_name} ,当前 {len (tool_cache)} 条" ) return result
场景四、工具调用日志与监控。记录所有工具调用的详细信息:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 @wrap_tool_call def monitor_tool_performance (request, handler ): """监控工具调用的性能指标""" tool_name = request.tool_call.get("name" , "unknown" ) tool_args = request.tool_call.get("args" , {}) start_time = time.time() try : result = handler(request) elapsed = time.time() - start_time print (f"[监控] {tool_name} ({tool_args} ) 成功,耗时 {elapsed:.2 f} s" ) return result except Exception as e: elapsed = time.time() - start_time print (f"[监控] {tool_name} ({tool_args} ) 失败,耗时 {elapsed:.2 f} s,错误: {e} " ) raise
场景五、根据结果决定后续流程。你可以根据工具执行结果决定是否继续 Agent 循环。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 @wrap_tool_call def check_empty_result (request, handler ): """如果工具返回空结果,直接结束 Agent,不浪费模型调用""" result = handler(request) if hasattr (result, 'content' ) and ( "未找到" in str (result.content) or "无结果" in str (result.content) or "没有" in str (result.content) ): from langchain.messages import AIMessage return Command(update={ "messages" : [ AIMessage(content="抱歉,没有找到相关信息。请换个关键词试试。" ) ] }) return result
@before_agent 和 @after_agent before_agent 和 after_agent 是 Agent 级别的钩子,分别在 Agent 执行之前和完成之后各执行一次。适合做初始化、预处理、后处理和统计分析。
before_agent 在 Agent 正式开始执行前运行,只运行以此。可以在这里做输入预处理、用户信息验证、资源初始化等等。
场景一、输入预处理——自动修正用户输入
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 @before_agent def preprocess_input (state, runtime ): """在 Agent 开始前处理用户输入""" messages = state.get("messages" , []) if not messages: return None last_msg = messages[-1 ] content = str (last_msg.content) if hasattr (last_msg, 'content' ) else "" greetings = ["你好" , "您好" , "hi" , "hello" , "嗨" ] if content and not any (content.lower().startswith(g) for g in greetings): pass return None
场景二、访问控制、权限检查:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 @before_agent def access_control (state, runtime ): """检查用户是否有权限使用 Agent""" context = runtime.context if context is None : return None user_role = context.get("user_role" , "guest" ) if user_role == "guest" : messages = state.get("messages" , []) if messages: last_content = str (messages[-1 ].content) restricted_keywords = ["删除" , "管理" , "配置" , "admin" ] if any (kw in last_content for kw in restricted_keywords): return { "jump_to" : "end" , "messages" : [HumanMessage( content="您当前的权限不足,无法执行此操作。请登录后重试。" )] } return None
after_agent 在 Agent 完成所有处理后执行(只执行一次)。你可以在这里格式化最终输出、记录统计信息、清理资源等。
场景三、统计分析:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 @after_agent def conversation_stats (state, runtime ): """统计对话信息并追加到结果中""" messages = state.get("messages" , []) model_calls = 0 tool_calls = 0 total_chars = 0 for msg in messages: if msg.type == "ai" : model_calls += 1 if hasattr (msg, 'tool_calls' ) and msg.tool_calls: tool_calls += len (msg.tool_calls) if hasattr (msg, 'content' ) and msg.content: total_chars += len (str (msg.content)) runtime.stream_writer({ "type" : "stats" , "model_calls" : model_calls, "tool_calls" : tool_calls, "total_messages" : len (messages), "total_chars" : total_chars, }) return None
以上就是所有 hook 的基本用法,在实际的 Agent 开发中还会有其他更复杂的场景。