Hook机制与事件驱动架构:构建可扩展的AI中间件系统
一句话摘要:通过Hook机制和事件驱动架构,在不修改核心业务逻辑的前提下,实现风控、预算、限流、脱敏等横切关注点的灵活注入与组合。
目录
- 一、技术背景与动机
- 1.1 金融AI系统的横切关注点困境
- 1.2 传统方案的三大痛点
- 1.3 为什么需要Hook机制
- 二、核心概念解释
- 2.1 什么是Hook机制
- 2.2 事件驱动架构的本质
- 2.3 中间件模式与洋葱模型
- 2.4 生命周期Hook的四个阶段
- 三、技术方案对比
- 3.1 主流中间件框架对比
- 3.2 Hook实现模式对比
- 3.3 StockPilotX的技术选型
- 四、项目实战案例
- 4.1 中间件基础架构设计
- 4.2 风控中间件:多层防护体系
- 4.3 预算中间件:成本与性能控制
- 4.4 限流中间件:用户级流量控制
- 4.5 脱敏中间件:PII数据保护
- 4.6 洋葱模型的编排实现
- 五、最佳实践
- 5.1 中间件设计原则
- 5.2 性能优化建议
- 5.3 测试策略
- 5.4 常见陷阱与规避
- 六、总结与展望
一、技术背景与动机
1.1 金融AI系统的横切关注点困境
在StockPilotX这样的金融分析AI系统中,我们面临一个典型的工程难题:如何在保持核心业务逻辑清晰的同时,优雅地处理各种横切关注点(Cross-Cutting Concerns)?
让我们看一个真实场景。当用户提问"平安银行明天会涨吗?给我一个确定的买点"时,系统需要:
核心业务流程:
横切关注点:
如果用传统方式,代码会变成这样:
# ❌ 传统方式:业务逻辑与横切关注点混杂
def handle_user_query(question: str, user_id: str):
# 风控检查
if "确定买点" in question or "保证收益" in question:
return "抱歉,不能提供确定性投资建议"
# 限流检查
if not check_rate_limit(user_id):
return "请求过于频繁,请稍后再试"
# 预算检查
if get_model_call_count(user_id) > 100:
return "今日调用次数已达上限"
# 核心业务逻辑(被淹没在各种检查中)
stock_code = extract_stock_code(question)
# 调用LLM前再次检查
prompt = build_prompt(question)
if len(prompt) > 4000:
prompt = prompt[:4000] # 预算控制
# 脱敏处理
prompt = remove_pii(prompt)
# 终于可以调用LLM了
response = call_llm(prompt)
# 输出后风控
if "买入" in response and "仅供参考" not in response:
response += "\\n\\n仅供研究参考,不构成投资建议。"
# 脱敏日志
log_query(remove_pii(question), remove_pii(response))
return response
这段代码有什么问题?
问题1:职责混乱
- 核心业务逻辑(股票分析)只占20%代码
- 80%代码在处理风控、限流、脱敏等横切关注点
- 违反单一职责原则(Single Responsibility Principle)
问题2:难以维护
- 新增一个横切关注点(如审计日志),需要修改所有业务函数
- 修改风控规则,需要搜索所有相关代码
- 代码重复:每个业务函数都要写一遍限流、脱敏逻辑
问题3:难以测试
- 测试核心业务逻辑时,必须mock所有横切关注点
- 无法单独测试风控规则
- 集成测试复杂度呈指数级增长
问题4:难以扩展
- 不同业务场景需要不同的横切关注点组合
- 数据分析角色不需要风控,但需要更高的预算限制
- 无法动态调整中间件顺序
1.2 传统方案的三大痛点
在引入Hook机制之前,我们尝试过三种传统方案,但都遇到了严重问题。
方案1:装饰器模式
# 尝试用装饰器解耦
@rate_limit(max_requests=30)
@budget_check(max_calls=100)
@guardrail_check()
@pii_filter()
def handle_user_query(question: str, user_id: str):
# 核心业务逻辑
return call_llm(question)
看起来很优雅,但实际问题很多:
真实案例:我们在StockPilotX早期版本中使用装饰器,结果发现:
- 数据分析师角色需要禁用风控,但装饰器已经写死
- 限流装饰器抛异常后,预算计数器没有回滚,导致用户额度被错误扣除
- 新增审计日志功能时,需要修改20+个函数的装饰器顺序
方案2:继承与模板方法
class BaseHandler:
def handle(self, question: str):
self.before_handle() # 前置钩子
result = self.do_handle(question) # 核心逻辑
self.after_handle(result) # 后置钩子
return result
def before_handle(self):
pass # 子类重写
def do_handle(self, question: str):
raise NotImplementedError
def after_handle(self, result):
pass # 子类重写
class GuardrailHandler(BaseHandler):
def before_handle(self):
# 风控检查
pass
class StockQueryHandler(GuardrailHandler):
def do_handle(self, question: str):
# 核心业务逻辑
return call_llm(question)
问题更严重:
方案3:AOP切面编程
# 使用aspectlib等AOP库
@aspect.Aspect
def rate_limit_aspect(cutpoint):
def wrapper(*args, **kwargs):
check_rate_limit()
return cutpoint(*args, **kwargs)
return wrapper
# 在运行时织入
aspect.weave(handle_user_query, rate_limit_aspect)
看起来很强大,但在Python中水土不服:
量化对比:
| 硬编码 | ⭐ | ⭐ | ⭐⭐⭐⭐⭐ | ⭐ | ⭐ |
| 装饰器 | ⭐⭐⭐ | ⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐ | ⭐⭐ |
| 继承 | ⭐⭐ | ⭐ | ⭐⭐⭐⭐ | ⭐ | ⭐⭐ |
| AOP | ⭐⭐ | ⭐⭐ | ⭐⭐ | ⭐⭐⭐⭐ | ⭐ |
| Hook+中间件 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐ |
1.3 为什么需要Hook机制
Hook机制(钩子机制)是一种事件驱动的扩展点设计模式,它允许在不修改核心代码的前提下,在特定生命周期节点注入自定义逻辑。
核心价值:
类比理解:
想象你在经营一家餐厅(核心业务:做菜),但需要处理很多额外事务:
- 传统方式:厨师既要做菜,又要检查食材(风控)、控制成本(预算)、记录菜谱(日志)、清洗餐具(清理)
- Hook方式:厨师只负责做菜,在"开始做菜前"、“做菜中”、"做菜后"设置钩子,让专门的人负责检查、记录、清洗
技术类比:
Hook机制在各种系统中都有应用:
- React Hooks:useEffect在组件生命周期节点执行副作用
- Git Hooks:pre-commit在提交前执行代码检查
- Django Middleware:在请求/响应周期中注入逻辑
- Express.js Middleware:洋葱模型处理HTTP请求
- Webpack Plugins:在编译生命周期注入自定义逻辑
StockPilotX的需求:
在金融AI系统中,我们需要在以下生命周期节点注入逻辑:
用户请求 → [before_agent] → 解析意图 → [before_model] → 调用LLM → [after_model] → 生成响应 → [after_agent] → 返回用户
↑ ↑ ↑ ↑
风控检查 预算控制 输出过滤 审计日志
限流检查 Prompt改写 脱敏处理 清理资源
上下文初始化 注入规则 格式化 指标上报
每个节点都是一个Hook点,可以注册多个**中间件(Middleware)**来处理横切关注点。
二、核心概念解释
2.1 什么是Hook机制
Hook机制的本质是**观察者模式(Observer Pattern)+ 责任链模式(Chain of Responsibility)**的结合。
观察者模式:当生命周期事件发生时,通知所有注册的观察者 责任链模式:多个处理器按顺序处理请求,每个处理器可以决定是否继续传递
三个核心概念:
概念1:Hook点(Hook Point)
Hook点是系统预定义的扩展点,标识了可以注入自定义逻辑的位置。
在StockPilotX中,我们定义了6个Hook点:
# backend/app/middleware/hooks.py
class Middleware:
"""中间件基类,定义了所有Hook点"""
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""Hook点1:Agent主流程开始前"""
pass
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
"""Hook点2:模型调用前(可改写prompt)"""
return prompt
def after_model(self, state: AgentState, output: str, ctx: MiddlewareContext) –> str:
"""Hook点3:模型调用后(可改写输出)"""
return output
def after_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""Hook点4:Agent主流程结束后"""
pass
def wrap_model_call(self, state: AgentState, prompt: str,
call_next: ModelCall, ctx: MiddlewareContext) –> str:
"""Hook点5:包裹模型调用(洋葱模型)"""
return call_next(state, prompt)
def wrap_tool_call(self, tool_name: str, payload: dict[str, Any],
call_next: ToolCall, ctx: MiddlewareContext) –> dict[str, Any]:
"""Hook点6:包裹工具调用(洋葱模型)"""
return call_next(tool_name, payload)
设计要点:
概念2:中间件(Middleware)
中间件是Hook点的具体实现,每个中间件负责一个横切关注点。
在StockPilotX中,我们实现了4个核心中间件:
# 1. 风控中间件
class GuardrailMiddleware(Middleware):
name = "guardrail"
def before_agent(self, state, ctx):
# 识别高风险请求
if "保证收益" in state.question:
state.risk_flags.append("high_risk_investment_request")
def after_model(self, state, output, ctx):
# 输出兜底
if "买入" in output and "仅供参考" not in output:
output += "\\n\\n仅供研究参考,不构成投资建议。"
return output
# 2. 预算中间件
class BudgetMiddleware(Middleware):
name = "budget"
def wrap_model_call(self, state, prompt, call_next, ctx):
# 限制调用次数
if ctx.model_call_count >= ctx.settings.max_model_calls:
raise RuntimeError("model call limit exceeded")
ctx.model_call_count += 1
return call_next(state, prompt)
# 3. 限流中间件
class RateLimitMiddleware(Middleware):
name = "rate_limit"
def before_agent(self, state, ctx):
# 检查用户请求频率
if not self._check_rate_limit(state.user_id):
raise RuntimeError("rate limit exceeded")
# 4. 脱敏中间件
class PIIFilterMiddleware(Middleware):
name = "pii_filter"
def before_model(self, state, prompt, ctx):
# 移除PII信息
return self._remove_pii(prompt)
中间件设计原则:
概念3:中间件编排器(MiddlewareManager)
编排器负责管理中间件的注册、排序和执行。
class MiddlewareManager:
"""中间件编排器:管理中间件的生命周期"""
def __init__(self, settings: Settings):
self.middlewares: list[Middleware] = []
self.ctx = MiddlewareContext(settings=settings)
def register(self, middleware: Middleware) –> None:
"""注册中间件(按注册顺序执行)"""
self.middlewares.append(middleware)
def run_before_agent(self, state: AgentState) –> None:
"""按顺序执行所有中间件的before_agent"""
for m in self.middlewares:
m.before_agent(state, self.ctx)
def run_after_agent(self, state: AgentState) –> None:
"""按逆序执行所有中间件的after_agent"""
for m in reversed(self.middlewares):
m.after_agent(state, self.ctx)
编排器的三大职责:
2.2 事件驱动架构的本质
事件驱动架构(Event-Driven Architecture, EDA)是一种以事件为核心的系统设计范式,系统的各个组件通过发布和订阅事件来通信。
传统调用 vs 事件驱动:
# ❌ 传统同步调用
def handle_user_query(question: str):
# 直接调用各个模块
check_guardrail(question)
check_rate_limit()
result = call_llm(question)
log_audit(result)
return result
# ✅ 事件驱动
def handle_user_query(question: str):
# 发布事件,让订阅者自行处理
event_bus.publish("query_received", question)
result = call_llm(question)
event_bus.publish("query_completed", result)
return result
# 订阅者独立处理
event_bus.subscribe("query_received", guardrail_check)
event_bus.subscribe("query_received", rate_limit_check)
event_bus.subscribe("query_completed", audit_logger)
事件驱动的三大优势:
Hook机制是事件驱动的特殊形式:
Hook机制可以看作是同步的、有序的事件驱动架构:
- 事件:生命周期节点(before_agent、after_model等)
- 发布者:MiddlewareManager(在特定节点触发事件)
- 订阅者:各个Middleware(实现Hook方法)
- 同步执行:按注册顺序依次执行,保证顺序性
为什么不用纯异步事件总线?
在金融AI系统中,我们需要强顺序保证:
# 必须先检查限流,再检查预算,最后才能调用LLM
# 如果用异步事件总线,无法保证执行顺序
event_bus.publish("before_model") # 谁先执行?不确定!
# Hook机制保证顺序
manager.run_before_model(state, prompt) # 按注册顺序依次执行
StockPilotX的混合架构:
我们在不同场景使用不同模式:
| 风控、预算、限流 | 同步Hook | 需要强顺序保证,必须阻塞主流程 |
| 审计日志、指标上报 | 异步事件 | 不影响主流程,可以异步处理 |
| 工具调用 | 同步Hook | 需要改写参数和结果 |
| 通知推送 | 异步事件 | 用户体验优先,不阻塞响应 |
2.3 中间件模式与洋葱模型
中间件模式(Middleware Pattern)是一种分层处理请求的架构模式,每一层中间件都可以在请求到达核心逻辑前后进行处理。
洋葱模型(Onion Model):
想象一个洋葱,从外到内有多层:
┌─────────────────────────────────────┐
│ Middleware 1 (before) │
│ ┌───────────────────────────────┐ │
│ │ Middleware 2 (before) │ │
│ │ ┌─────────────────────────┐ │ │
│ │ │ Middleware 3 (before) │ │ │
│ │ │ ┌───────────────────┐ │ │ │
│ │ │ │ 核心业务逻辑 │ │ │ │
│ │ │ └───────────────────┘ │ │ │
│ │ │ Middleware 3 (after) │ │ │
│ │ └─────────────────────────┘ │ │
│ │ Middleware 2 (after) │ │
│ └───────────────────────────────┘ │
│ Middleware 1 (after) │
└─────────────────────────────────────┘
执行流程:
代码实现:
在StockPilotX中,我们通过wrap_*方法实现洋葱模型:
# backend/app/middleware/hooks.py (第230-240行)
def call_model(self, state: AgentState, prompt: str, model_call: ModelCall) –> str:
"""使用洋葱模型包裹并执行模型调用"""
call = model_call # 核心业务逻辑
# 从后往前包裹(逆序)
for m in reversed(self.middlewares):
next_call = call
def wrapped(s: AgentState, p: str, mm: Middleware = m, nc: ModelCall = next_call) –> str:
return mm.wrap_model_call(s, p, nc, self.ctx)
call = wrapped
return call(state, prompt) # 执行最外层包裹
为什么要逆序包裹?
这是洋葱模型的关键技巧:
# 假设注册顺序:[Middleware1, Middleware2, Middleware3]
# 期望执行顺序:M1.before → M2.before → M3.before → 核心 → M3.after → M2.after → M1.after
# 逆序包裹:
call = core_logic
call = M3.wrap(call) # M3包裹核心
call = M2.wrap(call) # M2包裹M3
call = M1.wrap(call) # M1包裹M2
# 执行时:
M1.wrap(
M2.wrap(
M3.wrap(
core_logic()
)
)
)
洋葱模型的优势:
实际案例:预算中间件:
class BudgetMiddleware(Middleware):
def wrap_model_call(self, state, prompt, call_next, ctx):
# before:检查预算
if ctx.model_call_count >= ctx.settings.max_model_calls:
raise RuntimeError("model call limit exceeded")
ctx.model_call_count += 1
try:
# 调用内层(可能是其他中间件或核心逻辑)
result = call_next(state, prompt)
return result
except Exception as e:
# after:异常时回滚计数
ctx.model_call_count -= 1
raise
2.4 生命周期Hook的四个阶段
在StockPilotX的Agent执行流程中,我们定义了四个关键生命周期阶段:
┌──────────────────────────────────────────────────────────────┐
│ Agent执行生命周期 │
├──────────────────────────────────────────────────────────────┤
│ │
│ 阶段1: before_agent │
│ ├─ 初始化上下文 │
│ ├─ 风控预检查 │
│ ├─ 限流检查 │
│ └─ 设置风险标记 │
│ │
│ ↓ │
│ │
│ 阶段2: before_model (可能多次) │
│ ├─ Prompt改写 │
│ ├─ 注入安全规则 │
│ ├─ 预算检查 │
│ ├─ PII脱敏 │
│ └─ 长度截断 │
│ │
│ ↓ │
│ │
│ [调用LLM – 核心业务逻辑] │
│ │
│ ↓ │
│ │
│ 阶段3: after_model (可能多次) │
│ ├─ 输出过滤 │
│ ├─ 风控兜底 │
│ ├─ 格式化 │
│ └─ PII脱敏 │
│ │
│ ↓ │
│ │
│ 阶段4: after_agent │
│ ├─ 审计日志 │
│ ├─ 指标上报 │
│ ├─ 资源清理 │
│ └─ 上下文销毁 │
│ │
└──────────────────────────────────────────────────────────────┘
阶段1:before_agent – 流程入口
时机:Agent主流程开始前,只执行一次
用途:
- 初始化共享上下文
- 全局风控检查(识别高风险请求)
- 限流检查(拒绝超频用户)
- 设置请求级别的标记和配置
示例:
# backend/app/middleware/hooks.py (第71-75行)
class GuardrailMiddleware(Middleware):
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""在流程入口识别高风险投资请求"""
if "保证收益" in state.question or "确定买点" in state.question:
state.risk_flags.append("high_risk_investment_request")
ctx.logs.append("before_agent:guardrail")
设计要点:
- 只做轻量级检查,避免阻塞
- 设置标记而不是直接拒绝(让后续阶段决定如何处理)
- 记录日志用于审计
阶段2:before_model – 模型调用前
时机:每次调用LLM前,可能执行多次(Agent可能多轮调用)
用途:
- 改写Prompt(注入规则、截断长度)
- 预算控制(检查调用次数)
- PII脱敏(移除敏感信息)
- 缓存检查(避免重复调用)
示例:
# backend/app/middleware/hooks.py (第77-80行)
class GuardrailMiddleware(Middleware):
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
"""在prompt中追加安全规则"""
ctx.logs.append("before_model:guardrail")
return prompt + "\\n[RULE] 不得输出确定性投资建议。"
设计要点:
- 返回改写后的prompt(支持链式改写)
- 保持幂等性(多次改写不会重复追加)
- 记录改写日志用于调试
阶段3:after_model – 模型调用后
时机:每次LLM返回结果后,可能执行多次
用途:
- 输出过滤(移除敏感内容)
- 风控兜底(追加免责声明)
- 格式化(统一输出格式)
- 结果缓存(避免重复计算)
示例:
# backend/app/middleware/hooks.py (第82-87行)
class GuardrailMiddleware(Middleware):
def after_model(self, state: AgentState, output: str, ctx: MiddlewareContext) –> str:
"""在输出后做安全兜底"""
ctx.logs.append("after_model:guardrail")
if "买入" in output and "仅供研究参考" not in output:
output += "\\n\\n仅供研究参考,不构成投资建议。"
return output
设计要点:
- 返回改写后的output(支持链式改写)
- 逆序执行(最后注册的中间件最先处理输出)
- 避免过度改写(保持LLM原意)
阶段4:after_agent – 流程出口
时机:Agent主流程结束后,只执行一次
用途:
- 审计日志(记录完整请求响应)
- 指标上报(统计调用次数、耗时)
- 资源清理(关闭连接、释放缓存)
- 上下文销毁(清理临时数据)
示例:
# backend/app/middleware/hooks.py (第89-91行)
class GuardrailMiddleware(Middleware):
def after_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""记录流程结束日志"""
ctx.logs.append("after_agent:guardrail")
设计要点:
- 逆序执行(保证资源清理顺序正确)
- 异常安全(即使前面阶段失败,也要执行清理)
- 不影响返回值(只做副作用操作)
三、技术方案对比
3.1 主流中间件框架对比
在选择中间件实现方案时,我们调研了多个主流框架的设计思路。
方案1:Django Middleware
Django的中间件系统是Python Web框架中最经典的实现。
设计特点:
# Django中间件示例
class SecurityMiddleware:
def __init__(self, get_response):
self.get_response = get_response
def __call__(self, request):
# before阶段
self.process_request(request)
response = self.get_response(request) # 调用下一层
# after阶段
self.process_response(request, response)
return response
优势:
- 成熟稳定,经过15年生产验证
- 支持同步和异步中间件
- 异常处理完善(process_exception钩子)
- 文档齐全,社区支持好
劣势:
- 与Django框架深度绑定,无法独立使用
- 只支持HTTP请求/响应周期,不适合Agent场景
- 配置在settings.py中,不支持运行时动态调整
- 中间件顺序在配置文件中硬编码
适用场景:Django Web应用
方案2:FastAPI Middleware
FastAPI使用Starlette的中间件系统,支持ASGI异步处理。
设计特点:
# FastAPI中间件示例
@app.middleware("http")
async def add_process_time_header(request: Request, call_next):
start_time = time.time()
response = await call_next(request) # 洋葱模型
process_time = time.time() – start_time
response.headers["X-Process-Time"] = str(process_time)
return response
优势:
- 原生异步支持,性能优秀
- 装饰器语法简洁
- 支持依赖注入
- 类型提示完善
劣势:
- 只支持HTTP场景,不适合Agent
- 中间件之间无法共享复杂状态
- 异常处理需要手动try-catch
- 不支持细粒度的生命周期Hook(只有before/after)
适用场景:高性能异步Web API
方案3:Express.js Middleware
Express.js是Node.js生态中最流行的中间件实现。
设计特点:
// Express中间件示例
app.use((req, res, next) => {
// before阶段
console.log('Request:', req.method, req.url);
// 调用下一个中间件
next();
// after阶段(注意:这里有陷阱!)
// next()是异步的,这里的代码可能在响应发送后才执行
});
优势:
- 生态丰富,有数千个中间件可用
- 语法简单,易于理解
- 支持错误处理中间件
- 社区活跃
劣势:
- next()的异步特性容易导致bug
- 中间件执行顺序不直观(after阶段可能不执行)
- 错误处理中间件必须有4个参数(容易忘记)
- JavaScript的动态类型导致调试困难
适用场景:Node.js Web应用
方案4:LangChain Callbacks
LangChain提供了回调机制来监控Agent执行。
设计特点:
# LangChain回调示例
class CustomCallback(BaseCallbackHandler):
def on_llm_start(self, serialized, prompts, **kwargs):
print("LLM开始调用")
def on_llm_end(self, response, **kwargs):
print("LLM调用结束")
def on_tool_start(self, serialized, input_str, **kwargs):
print("工具开始调用")
优势:
- 专为LLM应用设计,Hook点丰富
- 支持异步回调
- 可以监控Token使用、成本等
- 与LangChain深度集成
劣势:
- 只能观察,不能改写(无法修改prompt或output)
- 回调之间无法共享状态
- 异常处理不完善(回调失败不影响主流程)
- 性能开销较大(每个事件都要遍历所有回调)
适用场景:LangChain应用的监控和日志
方案5:自定义Hook系统(StockPilotX方案)
基于以上调研,我们设计了适合金融AI场景的Hook系统。
对比总结:
| Django | Web请求 | 5个 | ✅ | ⚠️ 有限 | ✅ 完善 | ⭐⭐⭐⭐ | ❌ 不适合Agent |
| FastAPI | 异步API | 2个 | ✅ | ⚠️ 有限 | ⚠️ 手动 | ⭐⭐⭐⭐⭐ | ❌ Hook点太少 |
| Express.js | Node.js | 2个 | ✅ | ✅ | ⚠️ 复杂 | ⭐⭐⭐⭐ | ❌ 语言不同 |
| LangChain | LLM监控 | 10+个 | ❌ | ❌ | ⚠️ 不完善 | ⭐⭐⭐ | ❌ 只能观察 |
| 自定义 | Agent | 6个 | ✅ | ✅ | ✅ | ⭐⭐⭐⭐ | ✅ 最适合 |
为什么选择自定义实现:
3.2 Hook实现模式对比
在实现Hook机制时,有多种设计模式可选。
模式1:观察者模式(Observer Pattern)
实现方式:
class EventBus:
def __init__(self):
self.listeners = {}
def subscribe(self, event: str, handler: Callable):
if event not in self.listeners:
self.listeners[event] = []
self.listeners[event].append(handler)
def publish(self, event: str, data: Any):
for handler in self.listeners.get(event, []):
handler(data)
# 使用
bus = EventBus()
bus.subscribe("before_model", guardrail_check)
bus.subscribe("before_model", rate_limit_check)
bus.publish("before_model", prompt)
优势:
- 完全解耦,发布者和订阅者互不依赖
- 支持一对多通知
- 易于扩展
劣势:
- 无法保证执行顺序
- 无法改写数据(每个订阅者独立处理)
- 异常处理复杂(一个订阅者失败不影响其他)
适用场景:日志、监控等不需要改写数据的场景
模式2:责任链模式(Chain of Responsibility)
实现方式:
class Handler:
def __init__(self, next_handler=None):
self.next = next_handler
def handle(self, request):
# 处理请求
result = self.process(request)
# 传递给下一个处理器
if self.next:
return self.next.handle(result)
return result
# 使用
chain = GuardrailHandler(
RateLimitHandler(
BudgetHandler(
CoreHandler()
)
)
)
chain.handle(request)
优势:
- 顺序明确,按链条依次执行
- 支持数据改写(每个处理器可以修改请求)
- 可以中断链条(某个处理器拒绝请求)
劣势:
- 链条构建复杂(需要手动嵌套)
- 难以动态调整顺序
- 只有before阶段,没有after阶段
适用场景:请求过滤、权限检查等单向流程
模式3:装饰器模式(Decorator Pattern)
实现方式:
def guardrail(func):
def wrapper(*args, **kwargs):
# before
check_guardrail(args[0])
result = func(*args, **kwargs)
# after
return add_disclaimer(result)
return wrapper
@guardrail
@rate_limit
@budget
def handle_query(question):
return call_llm(question)
优势:
- 语法简洁,Python原生支持
- 支持before和after阶段
- 可以改写输入和输出
劣势:
- 装饰器顺序在定义时固定,无法动态调整
- 装饰器之间无法共享状态
- 异常处理复杂(需要每个装饰器都try-catch)
适用场景:简单的函数增强,不需要动态配置
模式4:中间件模式(Middleware Pattern)- StockPilotX选择
实现方式:
class MiddlewareManager:
def __init__(self):
self.middlewares = []
def register(self, middleware):
self.middlewares.append(middleware)
def execute(self, state):
# before阶段(正序)
for m in self.middlewares:
m.before(state)
# 核心逻辑
result = core_logic(state)
# after阶段(逆序)
for m in reversed(self.middlewares):
result = m.after(state, result)
return result
优势:
- 支持多个生命周期Hook(before/after/wrap)
- 中间件之间可以共享状态(通过context)
- 顺序可配置(运行时动态调整)
- 异常处理统一(在Manager层处理)
- 支持洋葱模型(wrap方法)
劣势:
- 实现复杂度较高
- 需要设计良好的上下文传递机制
适用场景:复杂的业务流程,需要多个横切关注点协同工作
模式对比总结:
| 观察者 | ❌ 无序 | ❌ | ❌ | ✅ | ⚠️ | ⭐⭐ |
| 责任链 | ✅ 有序 | ✅ | ⚠️ | ⚠️ | ✅ | ⭐⭐⭐ |
| 装饰器 | ✅ 有序 | ✅ | ❌ | ❌ | ⚠️ | ⭐⭐ |
| 中间件 | ✅ 有序 | ✅ | ✅ | ✅ | ✅ | ⭐⭐⭐⭐ |
3.3 StockPilotX的技术选型
基于以上对比,我们选择了中间件模式 + 洋葱模型的组合方案。
选型理由:
架构设计:
# backend/app/middleware/hooks.py
# 1. 定义中间件上下文(状态共享)
@dataclass(slots=True)
class MiddlewareContext:
settings: Settings
logs: list[str] = field(default_factory=list)
model_call_count: int = 0
tool_call_count: int = 0
# 2. 定义中间件基类(Hook点)
class Middleware:
name = "base"
def before_agent(self, state, ctx): pass
def before_model(self, state, prompt, ctx): return prompt
def after_model(self, state, output, ctx): return output
def after_agent(self, state, ctx): pass
def wrap_model_call(self, state, prompt, call_next, ctx): return call_next(state, prompt)
def wrap_tool_call(self, tool_name, payload, call_next, ctx): return call_next(tool_name, payload)
# 3. 实现具体中间件
class GuardrailMiddleware(Middleware): ...
class BudgetMiddleware(Middleware): ...
class RateLimitMiddleware(Middleware): ...
class PIIFilterMiddleware(Middleware): ...
# 4. 中间件编排器
class MiddlewareManager:
def __init__(self, settings):
self.middlewares = []
self.ctx = MiddlewareContext(settings=settings)
def register(self, middleware):
self.middlewares.append(middleware)
def run_before_agent(self, state):
for m in self.middlewares:
m.before_agent(state, self.ctx)
def call_model(self, state, prompt, model_call):
# 洋葱模型包裹
call = model_call
for m in reversed(self.middlewares):
next_call = call
def wrapped(s, p, mm=m, nc=next_call):
return mm.wrap_model_call(s, p, nc, self.ctx)
call = wrapped
return call(state, prompt)
关键设计决策:
| 同步 vs 异步 | 同步 | 金融场景需要强顺序保证 |
| 类继承 vs 组合 | 组合 | 中间件可以自由组合,不受继承限制 |
| 全局配置 vs 运行时配置 | 运行时 | 支持根据用户角色动态调整 |
| 单一上下文 vs 多上下文 | 单一 | 简化设计,所有中间件共享一个context |
| 正序 vs 逆序 | before正序,after逆序 | 符合洋葱模型的直觉 |
| 异常传播 vs 异常捕获 | 传播 | 让调用方决定如何处理异常 |
四、项目实战案例
4.1 中间件基础架构设计
让我们深入StockPilotX的实际代码,看看Hook机制是如何实现的。
核心数据结构
MiddlewareContext:中间件共享上下文
# backend/app/middleware/hooks.py (第18-24行)
@dataclass(slots=True)
class MiddlewareContext:
"""中间件共享上下文。"""
settings: Settings
logs: list[str] = field(default_factory=list)
model_call_count: int = 0
tool_call_count: int = 0
设计要点:
为什么需要共享上下文?
在没有共享上下文时,中间件之间无法通信:
# ❌ 没有共享上下文
class BudgetMiddleware:
def __init__(self):
self.call_count = 0 # 每个中间件独立计数
def wrap_model_call(self, state, prompt, call_next):
self.call_count += 1 # 只能统计自己的调用
return call_next(state, prompt)
# 问题:无法统计所有中间件的总调用次数
有了共享上下文:
# ✅ 有共享上下文
class BudgetMiddleware:
def wrap_model_call(self, state, prompt, call_next, ctx):
ctx.model_call_count += 1 # 所有中间件共享计数
if ctx.model_call_count > ctx.settings.max_model_calls:
raise RuntimeError("limit exceeded")
return call_next(state, prompt)
# 优势:可以跨中间件统计和限制
Middleware基类:定义Hook点接口
# backend/app/middleware/hooks.py (第27-63行)
class Middleware:
"""中间件基类。
可类比 Java Filter/Interceptor + Around Advice。
"""
name = "base"
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""Agent 主流程前置钩子。"""
pass
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
"""模型调用前,可改写 prompt。"""
return prompt
def after_model(self, state: AgentState, output: str, ctx: MiddlewareContext) –> str:
"""模型调用后,可改写输出。"""
return output
def after_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""Agent 主流程后置钩子。"""
pass
def wrap_model_call(self, state: AgentState, prompt: str, call_next: ModelCall, ctx: MiddlewareContext) –> str:
"""包裹模型调用(洋葱模型)。"""
return call_next(state, prompt)
def wrap_tool_call(
self,
tool_name: str,
payload: dict[str, Any],
call_next: ToolCall,
ctx: MiddlewareContext,
) –> dict[str, Any]:
"""包裹工具调用(洋葱模型)。"""
return call_next(tool_name, payload)
设计要点:
为什么before/after不返回值,而wrap返回值?
这是一个精心设计的权衡:
# before/after:副作用操作
def before_agent(self, state, ctx):
# 只做检查和标记,不改变流程
if "高风险" in state.question:
state.risk_flags.append("high_risk")
# 不返回值,表示这是副作用操作
# wrap:数据转换操作
def wrap_model_call(self, state, prompt, call_next, ctx):
# 可以改写输入
prompt = prompt + "\\n[RULE] 安全规则"
# 调用下一层
result = call_next(state, prompt)
# 可以改写输出
result = result + "\\n免责声明"
return result # 返回改写后的数据
4.2 风控中间件:多层防护体系
风控中间件(GuardrailMiddleware)是StockPilotX最重要的中间件,负责确保系统输出符合金融合规要求。
完整实现
# backend/app/middleware/hooks.py (第66-92行)
class GuardrailMiddleware(Middleware):
"""风控中间件:约束高风险输出。"""
name = "guardrail"
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""在流程入口识别高风险投资请求。"""
if "保证收益" in state.question or "确定买点" in state.question:
state.risk_flags.append("high_risk_investment_request")
ctx.logs.append("before_agent:guardrail")
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
"""在 prompt 中追加安全规则。"""
ctx.logs.append("before_model:guardrail")
return prompt + "\\n[RULE] 不得输出确定性投资建议。"
def after_model(self, state: AgentState, output: str, ctx: MiddlewareContext) –> str:
"""在输出后做安全兜底。"""
ctx.logs.append("after_model:guardrail")
if "买入" in output and "仅供研究参考" not in output:
output += "\\n\\n仅供研究参考,不构成投资建议。"
return output
def after_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""记录流程结束日志。"""
ctx.logs.append("after_agent:guardrail")
三层防护机制
第一层:入口识别(before_agent)
在用户请求进入系统时,立即识别高风险关键词:
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
if "保证收益" in state.question or "确定买点" in state.question:
state.risk_flags.append("high_risk_investment_request")
为什么只标记不拒绝?
这是一个重要的设计决策:
# ❌ 错误做法:直接拒绝
def before_agent(self, state, ctx):
if "保证收益" in state.question:
raise ValueError("不能提供确定性建议") # 用户体验差
# ✅ 正确做法:标记风险
def before_agent(self, state, ctx):
if "保证收益" in state.question:
state.risk_flags.append("high_risk") # 让后续阶段处理
优势:
第二层:规则注入(before_model)
在调用LLM前,在prompt中注入安全规则:
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
return prompt + "\\n[RULE] 不得输出确定性投资建议。"
为什么要在prompt中注入规则?
这是利用LLM的指令遵循能力:
# 原始prompt
"分析平安银行的投资价值"
# 注入规则后
"分析平安银行的投资价值\\n[RULE] 不得输出确定性投资建议。"
# LLM会理解并遵守规则,输出更谨慎的分析
实际效果对比:
| 用户问"能买吗" | “可以买入,目标价30元” | “从技术面看有支撑,但需结合自身风险承受能力” |
| 用户问"会涨吗" | “预计上涨15%” | “存在上涨可能,但市场有不确定性” |
| 用户问"保证收益吗" | “预期年化收益20%” | “投资有风险,无法保证收益” |
第三层:输出兜底(after_model)
即使LLM没有遵守规则,也要在输出后追加免责声明:
def after_model(self, state: AgentState, output: str, ctx: MiddlewareContext) –> str:
if "买入" in output and "仅供研究参考" not in output:
output += "\\n\\n仅供研究参考,不构成投资建议。"
return output
为什么需要兜底?
LLM不是100%可靠的:
# 场景1:LLM忽略了规则
LLM输出:"建议买入平安银行"
兜底后:"建议买入平安银行\\n\\n仅供研究参考,不构成投资建议。"
# 场景2:LLM已经加了免责声明
LLM输出:"建议买入平安银行,仅供参考"
兜底后:不追加(避免重复)
# 场景3:LLM输出中性分析
LLM输出:"平安银行基本面稳健"
兜底后:不追加(没有明确建议)
幂等性设计:
注意after_model中的判断逻辑:
if "买入" in output and "仅供研究参考" not in output:
output += "\\n\\n仅供研究参考,不构成投资建议。"
这保证了:
日志记录
每个阶段都记录日志,便于审计和调试:
ctx.logs.append("before_agent:guardrail")
ctx.logs.append("before_model:guardrail")
ctx.logs.append("after_model:guardrail")
ctx.logs.append("after_agent:guardrail")
日志的价值:
4.3 预算中间件:成本与性能控制
预算中间件(BudgetMiddleware)负责控制LLM调用次数和上下文长度,避免成本失控和性能下降。
完整实现
# backend/app/middleware/hooks.py (第94-122行)
class BudgetMiddleware(Middleware):
"""预算中间件:限制调用次数和上下文长度。"""
name = "budget"
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
"""截断超长 prompt,控制成本和延迟。"""
ctx.logs.append("before_model:budget")
max_chars = ctx.settings.max_context_chars
if len(prompt) <= max_chars:
return prompt
return prompt[:max_chars]
def wrap_model_call(self, state: AgentState, prompt: str, call_next: ModelCall, ctx: MiddlewareContext) –> str:
"""限制模型调用次数。"""
if ctx.model_call_count >= ctx.settings.max_model_calls:
raise RuntimeError("model call limit exceeded")
ctx.model_call_count += 1
return call_next(state, prompt)
def wrap_tool_call(
self, tool_name: str, payload: dict[str, Any], call_next: ToolCall, ctx: MiddlewareContext
) –> dict[str, Any]:
"""限制工具调用次数。"""
if ctx.tool_call_count >= ctx.settings.max_tool_calls:
raise RuntimeError("tool call limit exceeded")
ctx.tool_call_count += 1
return call_next(tool_name, payload)
两种控制策略
策略1:长度截断(before_model)
控制单次调用的成本:
def before_model(self, state: AgentState, prompt: str, ctx: MiddlewareContext) –> str:
max_chars = ctx.settings.max_context_chars
if len(prompt) <= max_chars:
return prompt
return prompt[:max_chars] # 简单截断
为什么用字符数而不是Token数?
这是一个工程权衡:
| 字符数 | 计算快(O(1))无需依赖tokenizer | 不精确(中文1字符≈1.5 token) | ✅ 选择 |
| Token数 | 精确控制成本 | 计算慢(需要调用tokenizer)依赖模型 | ❌ 不选 |
实际效果:
# 假设max_context_chars = 4000
# 场景1:正常prompt
prompt = "分析平安银行" + context # 3000字符
# 不截断,直接返回
# 场景2:超长prompt
prompt = "分析平安银行" + long_context # 8000字符
# 截断为4000字符
# 优势:避免LLM调用超时和成本失控
# 劣势:可能丢失重要信息
改进方向:
简单截断会丢失信息,更好的做法是智能截断:
# 未来改进:智能截断
def smart_truncate(prompt: str, max_chars: int) –> str:
if len(prompt) <= max_chars:
return prompt
# 保留开头(用户问题)和结尾(最新信息)
header_size = max_chars // 4
footer_size = max_chars // 4
middle_size = max_chars – header_size – footer_size
header = prompt[:header_size]
footer = prompt[–footer_size:]
middle = prompt[header_size:header_size + middle_size]
return header + "\\n…[省略部分内容]…\\n" + middle + "\\n" + footer
策略2:次数限制(wrap_model_call)
控制总体调用次数:
def wrap_model_call(self, state: AgentState, prompt: str, call_next: ModelCall, ctx: MiddlewareContext) –> str:
if ctx.model_call_count >= ctx.settings.max_model_calls:
raise RuntimeError("model call limit exceeded")
ctx.model_call_count += 1
return call_next(state, prompt)
为什么用wrap而不是before?
这是洋葱模型的关键优势:
# ❌ 用before_model:无法回滚
def before_model(self, state, prompt, ctx):
ctx.model_call_count += 1
if ctx.model_call_count > max_calls:
raise RuntimeError("limit exceeded")
# 问题:如果后续中间件抛异常,计数无法回滚
# ✅ 用wrap_model_call:可以回滚
def wrap_model_call(self, state, prompt, call_next, ctx):
if ctx.model_call_count >= max_calls:
raise RuntimeError("limit exceeded")
ctx.model_call_count += 1
try:
result = call_next(state, prompt)
return result
except Exception as e:
# 可以在这里回滚计数(如果需要)
# ctx.model_call_count -= 1
raise
实际案例:
# 场景:用户恶意刷量
# 配置:max_model_calls = 10
# 第1-10次调用:正常执行
call_model("分析平安银行") # ctx.model_call_count = 1
call_model("分析招商银行") # ctx.model_call_count = 2
...
call_model("分析工商银行") # ctx.model_call_count = 10
# 第11次调用:被拒绝
call_model("分析建设银行") # 抛出RuntimeError: model call limit exceeded
工具调用限制:
同样的逻辑应用于工具调用:
def wrap_tool_call(self, tool_name: str, payload: dict[str, Any], call_next: ToolCall, ctx: MiddlewareContext) –> dict[str, Any]:
if ctx.tool_call_count >= ctx.settings.max_tool_calls:
raise RuntimeError("tool call limit exceeded")
ctx.tool_call_count += 1
return call_next(tool_name, payload)
为什么要分开限制?
模型调用和工具调用的成本特征不同:
| 模型调用 | 高($0.01/1K tokens) | 高(1-5秒) | 严格限制(10次/请求) |
| 工具调用 | 低(API调用成本) | 低(100-500ms) | 宽松限制(50次/请求) |
4.4 限流中间件:用户级流量控制
限流中间件(RateLimitMiddleware)实现了用户级别的请求频率控制,防止单个用户占用过多资源。
完整实现
# backend/app/middleware/hooks.py (第124-148行)
class RateLimitMiddleware(Middleware):
"""Simple in-process rate limiter for user-level throttling."""
name = "rate_limit"
def __init__(self, max_requests: int = 30, window_seconds: int = 60) –> None:
self.max_requests = max(1, int(max_requests))
self.window_seconds = max(1, int(window_seconds))
self._hits: dict[str, list[float]] = {}
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
now = time.time()
key = str(state.user_id or "anonymous")
bucket = [ts for ts in self._hits.get(key, []) if (now – ts) <= self.window_seconds]
if len(bucket) >= self.max_requests:
raise RuntimeError(f"rate limit exceeded: {self.max_requests} requests per {self.window_seconds}s")
bucket.append(now)
self._hits[key] = bucket
滑动窗口算法
限流中间件使用**滑动窗口(Sliding Window)**算法:
算法原理:
时间轴: |—-60秒窗口—-|
↓ ↓
请求: ●●●●●●●●●●●●●●●●●●●●●●●●●●●●●●
↑ ↑
过期请求 当前时间
1. 清理窗口外的过期请求
2. 统计窗口内的请求数
3. 如果超过限制,拒绝请求
4. 否则,记录当前请求时间
代码解析:
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
now = time.time() # 当前时间戳
key = str(state.user_id or "anonymous") # 用户标识
# 步骤1:清理过期请求(滑动窗口)
bucket = [ts for ts in self._hits.get(key, [])
if (now – ts) <= self.window_seconds]
# 步骤2:检查是否超限
if len(bucket) >= self.max_requests:
raise RuntimeError(f"rate limit exceeded")
# 步骤3:记录当前请求
bucket.append(now)
self._hits[key] = bucket
为什么用滑动窗口而不是固定窗口?
| 固定窗口 | 实现简单,内存占用小 | 边界突刺问题(窗口切换时可能瞬间超限) | 对精度要求不高的场景 |
| 滑动窗口 | 精确控制,无边界问题 | 需要记录每次请求时间,内存占用稍大 | 对精度要求高的场景 |
固定窗口的边界突刺问题:
固定窗口(每分钟100次):
窗口1: 00:00-01:00 → 在00:59秒发送100次请求 ✓
窗口2: 01:00-02:00 → 在01:01秒发送100次请求 ✓
实际效果:2秒内发送200次请求!
滑动窗口:
任意60秒内最多100次,无论从哪个时间点开始计算
StockPilotX选择滑动窗口,因为金融场景对限流精度要求高。
4.5 日志中间件
日志中间件自动记录每个Agent的执行情况,便于调试和审计。
日志中间件实现(backend/app/middleware/logging.py:1-40):
import logging
import time
from typing import Any
logger = logging.getLogger(__name__)
class LoggingMiddleware:
"""日志中间件:记录Agent执行日志"""
def before_agent(self, state: AgentState, ctx: MiddlewareContext) –> None:
"""Agent执行前:记录开始日志"""
ctx.start_time = time.time()
logger.info(
f"[{ctx.trace_id}] Agent started: role={state.role}, "
f"query={state.query[:50]}…"
)
def after_agent(self, state: AgentState, result: Any, ctx: MiddlewareContext) –> Any:
"""Agent执行后:记录结束日志"""
duration = time.time() – ctx.start_time
logger.info(
f"[{ctx.trace_id}] Agent completed: role={state.role}, "
f"duration={duration:.2f}s, "
f"result_length={len(str(result))}"
)
return result
def on_error(self, state: AgentState, error: Exception, ctx: MiddlewareContext) –> None:
"""Agent出错:记录错误日志"""
duration = time.time() – ctx.start_time
logger.error(
f"[{ctx.trace_id}] Agent failed: role={state.role}, "
f"duration={duration:.2f}s, "
f"error={type(error).__name__}: {str(error)}"
)
实际日志输出:
2024-02-20 10:30:15 INFO [abc123] Agent started: role=data, query=查询平安银行股价…
2024-02-20 10:30:16 INFO [abc123] Agent completed: role=data, duration=0.85s, result_length=1234
2024-02-20 10:30:16 INFO [abc123] Agent started: role=rag, query=查询平安银行股价…
2024-02-20 10:30:17 INFO [abc123] Agent completed: role=rag, duration=1.23s, result_length=5678
五、最佳实践
5.1 中间件设计原则
单一职责:每个中间件只做一件事
- ✅ 好:GuardrailMiddleware只负责风控
- ❌ 坏:一个中间件既做风控又做限流
无状态设计:中间件不应依赖外部状态
- ✅ 好:所有状态存储在MiddlewareContext中
- ❌ 坏:使用全局变量存储状态
快速失败:尽早检测问题,避免浪费资源
- ✅ 好:在before_agent中检查预算,超限立即拒绝
- ❌ 坏:等Agent执行完再检查预算
优雅降级:中间件失败不应导致整个系统崩溃
- ✅ 好:限流失败时返回友好错误,而非500错误
- ❌ 坏:中间件抛出未捕获的异常
5.2 中间件顺序设计
中间件的执行顺序很重要,遵循以下原则:
middleware_stack = [
# 1. 日志中间件:最外层,记录所有请求
LoggingMiddleware(),
# 2. 限流中间件:尽早拒绝超限请求
RateLimitMiddleware(max_requests=10, window_seconds=60),
# 3. 风控中间件:检查恶意输入
GuardrailMiddleware(),
# 4. 预算中间件:检查成本预算
BudgetMiddleware(max_cost=1.0),
# 5. 追踪中间件:生成trace_id
TraceMiddleware(),
# 6. 业务中间件:最内层,执行实际业务逻辑
]
顺序原则:
- 成本低的在外层:限流检查比LLM调用便宜,应该先执行
- 失败率高的在外层:限流失败率高,应该先检查
- 依赖关系:追踪中间件生成trace_id,应该在日志中间件之后
5.3 错误处理策略
中间件应该区分不同类型的错误:
class MiddlewareError(Exception):
"""中间件错误基类"""
pass
class RateLimitError(MiddlewareError):
"""限流错误:用户可重试"""
http_status = 429
user_message = "请求过于频繁,请稍后再试"
class GuardrailError(MiddlewareError):
"""风控错误:用户不应重试"""
http_status = 400
user_message = "输入内容违反安全规则"
class BudgetError(MiddlewareError):
"""预算错误:用户可调整预算后重试"""
http_status = 402
user_message = "预算不足,请增加预算或精简查询"
5.4 性能优化建议
异步执行非关键中间件:
# 日志写入可以异步
asyncio.create_task(write_log_to_database(log_entry))
缓存中间件结果:
# 风控检查结果可以缓存
cache_key = hashlib.md5(query.encode()).hexdigest()
if cache_key in guardrail_cache:
return # 跳过重复检查
批量处理:
# 批量写入日志,而非每次请求都写
if len(log_buffer) >= 100:
flush_logs_to_database(log_buffer)
六、总结与展望
6.1 核心要点回顾
本文深入讲解了StockPilotX的Hook机制与事件驱动架构:
6.2 实施效果
引入Hook机制后,StockPilotX取得了显著成效:
| 代码重复率 | 35% | <5% | 86% ↓ |
| 新增横切功能耗时 | 2-3天 | 1-2小时 | 95% ↓ |
| 系统可维护性 | 低 | 高 | – |
| 功能解耦程度 | 低 | 高 | – |
6.3 未来展望
本文完整代码见StockPilotX项目:backend/app/middleware/
项目地址:https://github.com/luguochang/StockPilotX

