构建AI Agent中间件系统:提升可控性与可观测性的工程实践
1. 项目概述为什么我们需要在Agent执行流中“插一脚”最近和几个做AI应用落地的朋友聊天大家不约而同地提到了同一个痛点Agent智能体好用但不好管。无论是基于LangChain、AutoGen还是其他框架构建的Agent一旦把复杂的业务逻辑交给它去自主执行就像把车钥匙交给了一个驾驶技术一流但对本地路况一无所知的新司机。它能从A点开到B点但中途会不会超速会不会误入单行道遇到突发路况比如业务规则临时变更能不能灵活处理我们心里完全没底。这就是“中间件系统”要解决的核心问题。它不是一个独立运行的服务而是一套精巧的“拦截器”或“插件”机制允许我们在Agent既定的执行流水线中像设置检查站和加工站一样插入我们自定义的逻辑。想象一下你给那位新司机配了一个实时在线的副驾这个副驾不干涉方向盘但能在司机做出每个动作起步、转弯、加速前后进行安全提醒、路况播报甚至根据你的指令临时修改目的地。这个“副驾系统”就是Agent执行流中的中间件。它的价值远不止于“监控”。对于企业级应用而言它关乎可控性、可观测性和可扩展性。你可以用它来注入上下文在执行前自动从数据库或缓存中拉取最新的用户会话历史、产品信息丰富Agent的“记忆”。实施合规审查在Agent调用工具如发送邮件、查询数据库前校验操作是否符合安全策略比如检查邮件内容是否包含敏感词数据库查询是否越权。实现业务熔断当检测到Agent连续多次尝试失败或陷入循环时主动中断流程转由人工或备用流程处理避免资源浪费和不良体验。增强可观测性以非侵入的方式收集每个步骤的输入输出、耗时、Token消耗为性能分析和成本核算提供精准数据。进行后处理对Agent生成的结果进行格式化、翻译、敏感信息脱敏再返回给用户。简单说中间件系统让你在不重写Agent核心大脑的前提下为它的每一次“思考”和“行动”戴上符合你业务需求的“手套”与“眼镜”。下面我就结合具体的实现思路和踩过的坑来拆解如何构建这样一个系统。2. 核心设计思路洋葱模型与责任链模式在设计中间件系统时最忌讳的就是写成一堆散乱、耦合的if-else判断硬编码在Agent的主循环里。那样会导致代码迅速腐化难以维护和扩展。业界对此有非常成熟的设计模式可以借鉴其中最贴切的就是洋葱模型Onion Model和责任链模式Chain of Responsibility。2.1 理解执行流的“洋葱”结构把Agent的一次完整调用例如处理用户问题“帮我查一下上季度的销售额”想象成一个洋葱。最核心的“芯”是Agent真正的执行逻辑理解问题、规划步骤、调用工具、合成答案。而中间件就是包裹在这个芯外面的一层层“皮”。每一层皮中间件都有机会在请求向内传递到下一层/核心前和响应向外返回到上一层/客户端前执行自己的逻辑。这个顺序至关重要请求阶段Inbound从最外层中间件开始依次向内执行。适合做参数校验、上下文注入、身份认证。核心处理Agent执行核心任务。响应阶段Outbound从最内层中间件开始依次向外执行。适合做结果格式化、日志记录、错误包装。这种结构保证了中间件之间的解耦。一个负责日志的中间件完全不需要知道前面有没有权限校验中间件。2.2 实现责任链的两种常见方式在代码层面我们常用责任链模式来实现这个洋葱模型。假设我们有一个基础的Agent执行函数async def execute_agent(input: str) - str:。方式一装饰器Decorator链这是最Pythonic、最直观的方式。每个中间件都是一个装饰器函数它接收一个执行函数可能是核心函数也可能是已被其他装饰器包装过的函数返回一个新的包装函数。def logging_middleware(next_fn): async def wrapper(input_text): start_time time.time() print(f[LOG] 开始执行输入: {input_text}) try: result await next_fn(input_text) # 调用链中的下一个函数可能是核心也可能是其他中间件 print(f[LOG] 执行成功耗时: {time.time() - start_time:.2f}s) return result except Exception as e: print(f[LOG] 执行失败错误: {e}) raise return wrapper def validation_middleware(next_fn): async def wrapper(input_text): if not input_text or len(input_text.strip()) 0: raise ValueError(输入不能为空) # 可以在这里加入更复杂的业务规则校验 return await next_fn(input_text) return wrapper # 组合中间件注意顺序先应用的装饰器在外层 logging_middleware validation_middleware async def execute_agent(input: str) - str: # 这里是Agent的核心逻辑 return f处理结果: {input} # 调用 result await execute_agent(查询销售额)方式二显式的中间件处理器列表这种方式更结构化适合中间件数量多、需要动态配置的场景。定义一个中间件基类或协议然后在一个执行器中顺序调用。from abc import ABC, abstractmethod from typing import Any, Awaitable, Callable class Middleware(ABC): abstractmethod async def handle(self, request: Any, next_fn: Callable[[Any], Awaitable[Any]]) - Any: pass class LoggingMiddleware(Middleware): async def handle(self, request: str, next_fn): print(f[MW-LOG] 请求入参: {request}) response await next_fn(request) # 将控制权交给下一个中间件或核心逻辑 print(f[MW-LOG] 响应结果: {response}) return response class AgentCore: async def execute(self, input_text: str) - str: return f核心处理结果: {input_text} class MiddlewareExecutor: def __init__(self, core: AgentCore, middlewares: list[Middleware]): self.core core self.middlewares middlewares async def execute(self, input_text: str) - str: # 构建一个最终的“下一个函数”链 def create_handler(index: int) - Callable[[str], Awaitable[str]]: async def handler(current_input: str) - str: if index len(self.middlewares): # 调用当前中间件并传入下一个handler return await self.middlewares[index].handle(current_input, create_handler(index 1)) else: # 所有中间件执行完毕调用核心逻辑 return await self.core.execute(current_input) return handler # 从第一个中间件开始执行 return await create_handler(0)(input_text) # 使用 core AgentCore() executor MiddlewareExecutor(core, [LoggingMiddleware()]) result await executor.execute(测试输入)实操心得对于中小型项目或中间件逻辑相对简单的场景装饰器链更轻量、易读。而对于需要从配置文件加载中间件、中间件有复杂状态或依赖注入需求的企业级应用显式的处理器列表模式更具优势结构清晰易于测试和动态编排。3. 五大核心中间件场景与实现详解理解了设计模式我们来看几个必须实现的、具有代表性的中间件场景。我会给出关键代码和配置要点。3.1 上下文注入中间件让Agent拥有“记忆”Agent的每次调用通常是独立的但它需要历史对话、用户画像、业务数据作为上下文才能做出精准判断。硬编码在Prompt里不灵活每次都去数据库查又低效。实现要点在Inbound阶段工作在执行核心Agent前从外部源Redis、数据库、向量库获取数据。数据拼接将获取的数据以特定格式如## 用户历史\n{history}\n## 当前问题拼接到原始用户输入之前形成新的、增强版的Prompt。缓存策略为减轻外部源压力引入本地缓存如functools.lru_cache或分布式缓存Redis并设置合理的过期时间。import redis.asyncio as redis from functools import lru_cache class ContextEnrichmentMiddleware: def __init__(self, redis_client: redis.Redis): self.redis redis_client async def handle(self, request: dict, next_fn) - dict: request 结构: {user_id: 123, input_text: ..., session_id: ...} user_id request.get(user_id) session_id request.get(session_id) # 1. 获取用户历史对话带缓存 history await self._get_conversation_history(session_id) # 2. 获取用户画像如偏好、权限 profile await self._get_user_profile(user_id) # 3. 获取实时业务数据如产品库存 biz_data await self._get_business_data(request[input_text]) # 构建增强的Prompt enriched_prompt f 用户信息{profile} 最近对话历史 {history} 当前问题{request[input_text]} 相关业务数据{biz_data} 请根据以上信息回答。 # 更新请求 request[enriched_prompt] enriched_prompt # 将增强后的请求传递给下一个环节 return await next_fn(request) lru_cache(maxsize128) async def _get_conversation_history(self, session_id: str) - str: # 这里可以使用缓存避免频繁查询数据库 cache_key fchat_history:{session_id} history await self.redis.get(cache_key) if history: return history.decode(utf-8) # 缓存未命中查询数据库... # db_history await db.query(...) # await self.redis.setex(cache_key, 300, db_history) # 缓存5分钟 # return db_history return 暂无历史对话注意事项上下文注入要避免“信息过载”。给Agent太多无关信息会干扰其判断增加Token消耗和成本。最佳实践是动态相关性检索例如使用向量数据库只检索与当前问题最相关的几条历史记录或文档片段而不是全部加载。3.2 合规与安全审查中间件守住红线这是企业应用的生死线。必须在Agent调用外部工具Tool Call或最终输出答案前进行审查。实现要点双检查点工具调用前Pre-tool-call检查工具名、参数是否被允许。例如是否允许该用户调用“send_email”工具参数中的收件人是否在白名单内最终输出后Post-output对Agent生成的文本进行敏感词过滤、个人隐私信息PII脱敏、内容安全审核。异步与同步审查对于简单的关键词过滤可以同步进行对于复杂的AI内容审核如调用内容安全API应采用异步非阻塞方式避免影响主流程响应速度。熔断机制当连续多次审查不通过时应触发熔断暂时禁用某些功能或转人工。class SecurityAndComplianceMiddleware: def __init__(self, forbidden_tools: set, pii_detector): self.forbidden_tools forbidden_tools # 例如 {delete_database, format_system} self.pii_detector pii_detector # 一个用于检测手机号、邮箱等的正则或模型 async def handle(self, request: dict, next_fn) - dict: # 假设request中包含了Agent规划的行动步骤 action_plan request.get(action_plan, []) for step in action_plan: if step[type] tool_call: tool_name step[name] # 1. 工具权限校验 if tool_name in self.forbidden_tools: raise PermissionError(f无权调用工具: {tool_name}) # 2. 参数安全检查示例检查SQL注入风险 if tool_name query_database: sql step[args].get(query, ) if self._has_sql_injection_risk(sql): raise SecurityError(检测到潜在的SQL注入风险) # 安全校验通过执行后续逻辑 response await next_fn(request) # 3. 对最终输出进行PII脱敏 if final_output in response: response[final_output] self.pii_detector.redact(response[final_output]) return response def _has_sql_injection_risk(self, sql: str) - bool: # 简单的关键词检查生产环境应使用更专业的SQL解析库或ORM的参数化查询来根本性避免 dangerous_patterns [DROP, DELETE FROM, INSERT INTO, --, ;, UNION SELECT] import re for pattern in dangerous_patterns: if re.search(rf\b{pattern}\b, sql, re.IGNORECASE): return True return False3.3 可观测性与监控中间件看清每一个环节没有监控的系统就是“盲开”。这个中间件负责无侵入地收集全链路数据。实现要点结构化日志不仅打印文本更要以JSON等结构化格式记录方便后续接入ELKElasticsearch, Logstash, Kibana或时序数据库。记录内容应包括request_id,user_id,stage阶段,input,output,latency延迟,token_usage,error如果有。指标埋点使用像Prometheus这样的工具记录计数器如请求总数、错误数、直方图如响应时间分布。分布式追踪如果系统是分布式的例如Agent、工具、模型服务分开部署集成OpenTelemetry来追踪一个请求在所有服务间的流转路径。import time import json import logging from contextvars import ContextVar from prometheus_client import Counter, Histogram # 定义指标 REQUEST_COUNT Counter(agent_requests_total, Total agent requests, [endpoint, status]) REQUEST_LATENCY Histogram(agent_request_duration_seconds, Request latency in seconds, [endpoint]) # 用于传递请求ID的上下文变量 request_id_ctx: ContextVar[str] ContextVar(request_id, defaultunknown) class ObservabilityMiddleware: def __init__(self, logger: logging.Logger): self.logger logger async def handle(self, request: dict, next_fn) - dict: request_id request.get(request_id, default_id) request_id_ctx.set(request_id) start_time time.time() endpoint request.get(path, execute) # 记录请求开始日志 self.logger.info(json.dumps({ request_id: request_id, event: request_start, input: request.get(input_text, )[:200], # 截断避免日志过大 user: request.get(user_id) })) try: response await next_fn(request) status success # 记录成功日志和指标 self.logger.info(json.dumps({ request_id: request_id, event: request_end, status: status, latency_sec: time.time() - start_time, token_used: response.get(usage, {}).get(total_tokens, 0) })) REQUEST_COUNT.labels(endpointendpoint, statusstatus).inc() REQUEST_LATENCY.labels(endpointendpoint).observe(time.time() - start_time) return response except Exception as e: status error # 记录错误日志和指标 self.logger.error(json.dumps({ request_id: request_id, event: request_error, error_type: e.__class__.__name__, error_msg: str(e), latency_sec: time.time() - start_time }), exc_infoTrue) REQUEST_COUNT.labels(endpointendpoint, statusstatus).inc() raise # 重新抛出异常让上层处理 # 在其他地方可以通过 request_id_ctx.get() 获取当前请求ID3.4 流式输出中间件提升用户体验对于生成时间较长的内容让用户干等着是一种糟糕的体验。中间件可以拦截Agent的流式生成结果如果底层LLM支持并实时推送给前端。实现要点拦截生成器核心Agent的execute方法需要改为异步生成器async generator逐步yieldtokens或chunks。中间件也需支持生成器中间件的handle方法需要能够接收并包装一个生成器在yield每个片段前后执行逻辑例如计算已生成字数、进行部分结果的安全扫描。与Web框架集成对于FastAPI或类似框架可以将这个生成器作为StreamingResponse的内容返回。from typing import AsyncGenerator class StreamingMiddleware: 处理流式输出的中间件示例 async def handle_stream(self, request: dict, next_fn) - AsyncGenerator[str, None]: next_fn 现在返回一个异步生成器 (AsyncGenerator) # Inbound 逻辑可以在这里做一些初始化 buffer [] print(f开始流式处理请求: {request[input_text][:50]}...) # 调用下一个环节获取流式生成器 stream_generator await next_fn(request) # 注意这里next_fn可能返回一个生成器 # 处理流式响应 async for chunk in stream_generator: # 这里可以加入后处理逻辑例如 # 1. 缓存片段用于后续完整内容分析 buffer.append(chunk) # 2. 简单的实时敏感词检测生产环境可能需要更复杂的缓冲窗口检测 if 敏感词 in chunk: chunk chunk.replace(敏感词, ***) # 3. 添加SSE (Server-Sent Events) 格式包装 # formatted_chunk fdata: {chunk}\n\n yield chunk # 将处理后的片段返回给上游 # 流结束后可以做一些清理或最终分析 full_content .join(buffer) print(f流式处理结束总长度: {len(full_content)}) # 例如将完整内容存入审计日志 # await audit_log(request, full_content) # 假设核心Agent的流式执行函数 async def agent_execute_stream(request: dict) - AsyncGenerator[str, None]: # 模拟LLM流式生成 simulated_response 这是一个流式生成的回答它被分成多个部分返回。 for word in simulated_response.split(): await asyncio.sleep(0.1) # 模拟生成延迟 yield word 3.5 故障熔断与降级中间件保障系统韧性当依赖的下游服务如LLM API、向量数据库、工具服务不稳定时不能让Agent无限重试或长时间等待需要有快速失败和降级策略。实现要点集成熔断器模式使用如pybreaker库为每个关键依赖设置熔断器。当失败次数超过阈值熔断器“打开”短时间内直接拒绝请求避免雪崩经过一段时间后进入“半开”状态试探性放行少量请求。超时控制为每个中间件环节或工具调用设置独立的超时时间使用asyncio.wait_for防止单个慢请求阻塞整个链路。优雅降级当核心LLM服务不可用时可以降级到规则引擎返回一个预设的兜底答案或者返回一个友好的错误提示引导用户稍后重试。import pybreaker import asyncio from datetime import datetime # 为LLM调用定义一个熔断器 llm_breaker pybreaker.CircuitBreaker(fail_max5, reset_timeout60) # 5次失败后熔断60秒 class CircuitBreakerMiddleware: def __init__(self): self.breaker llm_breaker async def handle(self, request: dict, next_fn) - dict: # 尝试通过熔断器执行 try: # 使用熔断器包装调用 response await self.breaker.call_async(self._call_with_timeout, request, next_fn) return response except pybreaker.CircuitBreakerError as e: # 熔断器已打开直接拒绝请求执行降级逻辑 self.logger.warning(f熔断器已打开请求被快速失败: {e}) return await self._fallback_response(request) except asyncio.TimeoutError: self.logger.error(请求超时) # 超时也算作一次失败熔断器会计数 raise except Exception as e: # 其他异常也会被熔断器捕获并计数 self.logger.error(f请求执行失败: {e}) raise async def _call_with_timeout(self, request: dict, next_fn) - dict: 为关键调用添加超时控制 try: # 设置10秒超时 return await asyncio.wait_for(next_fn(request), timeout10.0) except asyncio.TimeoutError: self.logger.error(下游服务响应超时) raise async def _fallback_response(self, request: dict) - dict: 优雅降级返回预设的兜底答案 # 这里可以根据请求类型返回不同的兜底内容 # 例如从缓存中获取旧答案或返回一个静态提示 return { final_output: 当前服务暂时繁忙您可以稍后重试或尝试描述其他问题。, source: fallback, is_fallback: True }4. 中间件系统的配置、测试与部署实践设计好了中间件如何管理它们并确保它们稳定可靠地工作4.1 动态配置与编排我们不应该把中间件的启用顺序和参数硬编码在代码里。理想的方式是通过配置文件如YAML或配置中心来管理。# middleware_config.yaml middlewares: - name: request_logging class: app.middleware.ObservabilityMiddleware enabled: true params: log_level: INFO - name: context_enrichment class: app.middleware.ContextEnrichmentMiddleware enabled: true params: redis_url: redis://localhost:6379/0 cache_ttl: 300 - name: security_check class: app.middleware.SecurityAndComplianceMiddleware enabled: true params: forbidden_tools: [send_email, execute_code] - name: circuit_breaker class: app.middleware.CircuitBreakerMiddleware enabled: true # 注意顺序很重要先执行的中间件在外层。然后在应用启动时根据这个配置动态加载和实例化中间件构建执行链。这样在不停机的情况下通过修改配置就能启用、禁用或调整某个中间件。4.2 中间件的单元测试与集成测试中间件逻辑必须被充分测试尤其是安全、合规相关的。单元测试针对每个中间件的handle方法模拟输入request和next_fn通常是一个返回固定值的简单异步函数断言其输出或副作用如日志记录、指标增长是否符合预期。import pytest from unittest.mock import AsyncMock, MagicMock pytest.mark.asyncio async def test_security_middleware_blocks_forbidden_tool(): middleware SecurityAndComplianceMiddleware(forbidden_tools{dangerous_tool}) mock_next AsyncMock(return_value{output: ok}) malicious_request { action_plan: [{type: tool_call, name: dangerous_tool}] } with pytest.raises(PermissionError): await middleware.handle(malicious_request, mock_next) # 确保下一个函数没有被调用 mock_next.assert_not_called()集成测试将多个中间件与核心Agent串联起来进行端到端测试。使用测试专用的LLM Mock如unittest.mock.patch掉OpenAI客户端和工具Mock验证整个链路的输入输出。4.3 性能考量与最佳实践中间件会增加开销设计不当会成为性能瓶颈。异步Async是所有I/O操作的必须任何涉及网络调用数据库、Redis、API的中间件逻辑都必须使用异步库如aioredis,aiohttp和async/await语法避免阻塞事件循环。避免在中间件中进行重型同步计算如复杂的字符串处理、大文件解析。如果必须考虑将其放入线程池执行asyncio.to_thread防止阻塞其他并发请求。设置合理的超时和超时传递为每个可能耗时的中间件环节设置独立的超时并且超时时间应从外向内逐层递减确保整体响应时间可控。中间件应保持无状态Stateless尽可能让中间件本身不保存请求相关的状态。状态应存储在请求上下文ContextVar或外部存储中。这有利于水平扩展和避免并发问题。谨慎使用全局变量在多线程/协程环境下修改全局变量是危险的。如果需要共享配置使用依赖注入或单例模式。5. 常见问题排查与实战技巧在实际部署和运行中你肯定会遇到各种问题。这里记录了几个最典型的坑和解决办法。5.1 中间件执行顺序错乱或未生效问题现象日志中间件没记录到错误或者安全校验在工具调用后才执行。排查步骤检查注册顺序确认中间件加载到执行链MiddlewareExecutor的顺序是否正确。最外层的中间件最先执行handle的“入站”逻辑但最后执行“出站”逻辑。画一个简单的执行顺序图来验证。验证中间件是否被正确包装在装饰器模式下检查装饰器应用的顺序。a b def fn()等价于fn a(b(fn))所以a是最外层。检查中间件的next_fn调用确保每个中间件的handle方法内部都正确调用了await next_fn(request)这是责任链向下传递的关键。如果某个中间件因为条件判断提前返回了响应链就会在此中断。5.2 内存泄漏或性能下降问题现象服务运行一段时间后内存持续增长响应变慢。可能原因与解决中间件中未释放资源例如在中间件中打开了数据库连接或文件句柄但在异常情况下没有正确关闭。使用try...finally块或异步上下文管理器async with来确保资源释放。缓存失控上下文注入中间件如果使用了无过期时间或过大的缓存如lru_cache缓存了过大的对象会导致内存堆积。为缓存设置大小限制和合理的TTL。日志中间件记录全量数据如果日志中间件无差别地记录完整的请求和响应体特别是包含大文件Base64编码会迅速耗尽磁盘和内存。务必进行截断或采样记录。使用asyncio.create_task后未管理如果在中间件内部分发了后台任务如异步发送审计日志但没有妥善收集或等待它们可能导致任务堆积。考虑使用后台任务管理器。5.3 错误处理与异常传播问题现象某个中间件抛出异常后整个请求无声无息地失败了或者返回了不友好的错误。最佳实践中间件应捕获自身逻辑的异常如果中间件的辅助功能如发送监控指标失败不应中断主业务流程。应使用try...except记录错误后继续调用next_fn。async def handle(self, request, next_fn): try: # 非核心的监控上报 await self._report_metrics(request) except Exception as e: self.logger.error(f监控上报失败但不中断流程: {e}) # 核心流程继续 return await next_fn(request)区分业务异常和系统异常安全中间件校验不通过应抛出明确的业务异常如PermissionError并由全局异常处理器转换为对用户友好的错误信息。而网络超时等系统异常则应触发熔断或重试机制。设计统一的错误响应格式在所有中间件和核心逻辑的最外层有一个全局异常捕获中间件负责将任何未处理的异常转换为结构化的错误响应如{code: 500, msg: 系统内部错误, request_id: xxx}并确保记录到日志。5.4 调试困难请求在哪个中间件卡住了技巧为每个请求生成唯一的request_id并在每个中间件的入口和出口都打印带有该ID的日志。这样在分布式日志系统中你可以通过request_id轻松串联起一个请求在所有中间件和核心服务中的完整生命周期快速定位延迟或错误的环节。这就是前面ObservabilityMiddleware和ContextVar所做的事情。构建一个健壮的Agent中间件系统初期会花费一些设计精力但一旦建成它将成为你AI应用架构中不可或缺的“神经系统”让你对Agent的每一次运行都了如指掌、掌控自如。从简单的日志记录开始逐步引入安全、监控、流式处理等中间件你的Agent系统会在这个过程中变得越来越强大和可靠。