问题背景
原题:Evaluator-optimizer 循环上线后偶尔跑满 CPU 烧掉大量 token。缺了什么?
根因:缺少三层刹车机制:
完整复现:有缺陷的 Evaluator-optimizer 循环
1. 缺陷版本(在线路上烧掉大量 token)
from typing import TypedDict, Literal
from langgraph.graph import StateGraph, START, END
from langchain_openai import ChatOpenAI
# ⚠️ 缺陷版本:没有轮数计数器,评估器输出自然语言
class State(TypedDict):
draft: str
feedback: str
passed: bool # 用自然语言解析,不稳定
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.7)
def generate(state: State):
"""生成节点"""
prompt = "写一句关于人工智能的营销文案"
if state.get("feedback"):
prompt = f"根据反馈修改文案。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
return {"draft": llm.invoke(prompt).content}
def evaluate(state: State):
"""评估节点——用自然语言判断,容易出错"""
response = llm.invoke(f"评审这版文案是否合格:{state['draft']}\\n请回答'合格'或'不合格',并给出改进建议")
content = response.content
# ❌ 用自然语言解析,极其不稳定
if "合格" in content:
return {"passed": True, "feedback": content}
else:
return {"passed": False, "feedback": content}
def should_continue(state: State) -> Literal["generate", "__end__"]:
"""条件判断——只有 passed 判断,没有轮数限制"""
# ❌ 没有轮数计数器,永远不强制退出
if state["passed"]:
return "__end__"
return "generate"
# 构建图
builder = StateGraph(State)
builder.add_node("generate", generate)
builder.add_node("evaluate", evaluate)
builder.add_edge(START, "generate")
builder.add_edge("generate", "evaluate")
builder.add_conditional_edges("evaluate", should_continue)
graph = builder.compile()
# 执行——可能烧掉大量 token!
try:
result = graph.invoke({
"draft": "",
"feedback": "",
"passed": False
})
print(f"最终结果: {result['draft']}")
except Exception as e:
print(f"异常: {e}")
# 可能触发 recursion_limit 默认值 25,但已经烧了很多 token
2. 问题复现:为什么会出现死循环
场景一:评估器输出歧义
# 评估器输出:"这个文案还行,但可以更好。虽然不算特别不合格,但建议修改…"
# 解析结果:同时包含"合格"和"不合格",逻辑混乱
# 可能的解析结果:
# – 包含"合格" → passed = True → 退出
# – 包含"不合格" → passed = False → 继续循环
# – 同时包含 → 取决于代码逻辑,可能永远不退出
场景二:模型理解偏差
# 生成器输出:"人工智能是未来"
# 评估器认为:"太简单了,不合格"
# 生成器修改后:"人工智能将改变世界"
# 评估器认为:"还是太泛,不合格"
# 生成器再修改:"AI 技术正在重塑各个行业"
# 评估器认为:"有点进步,但还是不够具体"
# … 无限循环下去
实际烧钱演示
# 模拟一次生产事故
import time
def simulate_bad_loop():
"""模拟有缺陷的循环——烧掉大量 token"""
print("模拟缺陷循环开始…")
rounds = 0
total_tokens = 0
state = {"draft": "", "feedback": "", "passed": False}
while not state["passed"]:
rounds += 1
# 每次调用大约消耗 500 tokens
tokens_this_round = 500
total_tokens += tokens_this_round
print(f"第 {rounds} 轮,消耗 tokens: {tokens_this_round}, 累计: {total_tokens}")
# 模拟生成节点
state["draft"] = f"第 {rounds} 版文案"
# 模拟评估节点——自然语言解析失败
import random
if random.random() < 0.7: # 70% 概率认为不合格
state["passed"] = False
state["feedback"] = "还需要改进"
else:
state["passed"] = True # 偶尔才通过
if rounds >= 10: # 假设没有刹车,跑到 10 轮
print(f"⚠️ 已跑 {rounds} 轮,消耗 {total_tokens} tokens,仍未通过!")
break
print(f"总轮数: {rounds}, 总消耗 tokens: {total_tokens}")
print(f"按 GPT-4o-mini 价格估算: ${total_tokens * 0.00015 / 1000:.4f}")
simulate_bad_loop()
# 输出:
# 模拟缺陷循环开始…
# 第 1 轮,消耗 tokens: 500, 累计: 500
# 第 2 轮,消耗 tokens: 500, 累计: 1000
# …
# 第 10 轮,消耗 tokens: 500, 累计: 5000
# ⚠️ 已跑 10 轮,消耗 5000 tokens,仍未通过!
# 总轮数: 10, 总消耗 tokens: 5000
修复方案:三层刹车机制
方案一:基础修复(加轮数计数器 + 结构化输出)
from typing import TypedDict, Literal
from pydantic import BaseModel, Field
from langgraph.graph import StateGraph, START, END
from langchain_openai import ChatOpenAI
# ✅ 修复1:用 Pydantic 模型定义结构化输出
class Review(BaseModel):
"""评审结果的结构化输出"""
passed: bool = Field(description="文案是否合格")
feedback: str = Field(description="改进建议,如果合格则写'无需修改'")
score: int = Field(description="评分 1-10", ge=1, le=10)
# ✅ 修复2:状态中加入轮数计数器
class State(TypedDict):
draft: str
feedback: str
passed: bool
rounds: int # 轮数计数器
MAX_ROUNDS = 3 # 最大轮数
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.7)
def generate(state: State):
"""生成节点"""
prompt = "写一句关于人工智能的营销文案,要求简洁有力"
if state.get("feedback"):
prompt = f"根据反馈修改文案。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
return {
"draft": llm.invoke(prompt).content,
"rounds": state.get("rounds", 0) + 1 # 轮数递增
}
def evaluate(state: State):
"""评估节点——用结构化输出"""
# ✅ 用 with_structured_output 保证结构化
review = llm.with_structured_output(Review).invoke(
f"评审这版文案是否合格:{state['draft']}"
)
return {
"passed": review.passed,
"feedback": review.feedback
}
def should_continue(state: State) -> Literal["generate", "__end__"]:
"""条件判断——三层刹车"""
# 第一层:主动刹车——轮数到顶
if state["rounds"] >= MAX_ROUNDS:
print(f"⚠️ 轮数达到上限 {MAX_ROUNDS},强制退出")
return "__end__"
# 第二层:评估通过
if state["passed"]:
return "__end__"
# 继续循环
return "generate"
# 构建图
builder = StateGraph(State)
builder.add_node("generate", generate)
builder.add_node("evaluate", evaluate)
builder.add_edge(START, "generate")
builder.add_edge("generate", "evaluate")
builder.add_conditional_edges("evaluate", should_continue)
graph = builder.compile()
# 执行
result = graph.invoke({
"draft": "",
"feedback": "",
"passed": False,
"rounds": 0
}, config={"recursion_limit": 20}) # 第三层:被动刹车
print(f"最终结果: {result['draft']}")
print(f"总轮数: {result['rounds']}")
方案二:完整修复(三层刹车 + 日志 + 降级策略)
from typing import TypedDict, Literal
from pydantic import BaseModel, Field
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from langchain_openai import ChatOpenAI
import logging
# 配置日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# 结构化输出
class Review(BaseModel):
passed: bool = Field(description="文案是否合格")
feedback: str = Field(description="改进建议")
score: int = Field(description="评分 1-10", ge=1, le=10)
class State(TypedDict):
draft: str
feedback: str
passed: bool
rounds: int
final_score: int # 记录最终评分
MAX_ROUNDS = 3 # 最大轮数
RECURSION_LIMIT = 25 # 递归限制
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.7)
def generate(state: State):
"""生成节点"""
current_round = state.get("rounds", 0) + 1
logger.info(f"第 {current_round} 轮生成开始")
prompt = "写一句关于人工智能的营销文案,要求简洁有力"
if state.get("feedback"):
# 根据轮数调整 prompt
if current_round >= MAX_ROUNDS – 1:
# 最后一轮:使用更宽松的指令
prompt = f"这是最后一次修改机会。请尽量优化文案,但不要过度修改。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
else:
prompt = f"根据反馈修改文案。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
return {
"draft": llm.invoke(prompt).content,
"rounds": current_round
}
def evaluate(state: State):
"""评估节点"""
logger.info(f"第 {state['rounds']} 轮评估开始")
try:
# 结构化输出
review = llm.with_structured_output(Review).invoke(
f"评审这版文案是否合格:{state['draft']}"
)
logger.info(f"评估结果: passed={review.passed}, score={review.score}")
return {
"passed": review.passed,
"feedback": review.feedback,
"final_score": review.score
}
except Exception as e:
# 降级策略:评估失败时默认通过
logger.error(f"评估失败: {e},默认通过")
return {
"passed": True,
"feedback": "评估失败,自动通过",
"final_score": 5
}
def should_continue(state: State) -> Literal["generate", "__end__"]:
"""条件判断——三层刹车"""
current_round = state.get("rounds", 0)
# 第一层刹车:主动刹车——轮数到顶
if current_round >= MAX_ROUNDS:
logger.warning(f"轮数达到上限 {MAX_ROUNDS},强制退出")
return "__end__"
# 第二层刹车:评估通过
if state["passed"]:
logger.info(f"文案合格,退出循环")
return "__end__"
# 继续循环
logger.info(f"继续第 {current_round + 1} 轮")
return "generate"
# 构建图
builder = StateGraph(State)
builder.add_node("generate", generate)
builder.add_node("evaluate", evaluate)
builder.add_edge(START, "generate")
builder.add_edge("generate", "evaluate")
builder.add_conditional_edges("evaluate", should_continue)
# 使用 checkpointer 支持中断恢复
graph = builder.compile(checkpointer=InMemorySaver())
# 执行
config = {
"configurable": {"thread_id": "eval-optimizer-001"},
"recursion_limit": RECURSION_LIMIT # 第三层刹车:被动刹车
}
result = graph.invoke({
"draft": "",
"feedback": "",
"passed": False,
"rounds": 0,
"final_score": 0
}, config)
print(f"\\n最终结果:")
print(f"文案: {result['draft']}")
print(f"总轮数: {result['rounds']}")
print(f"最终评分: {result['final_score']}")
print(f"是否通过: {result['passed']}")
方案三:生产级实现(带监控和告警)
from typing import TypedDict, Literal
from pydantic import BaseModel, Field
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from langchain_openai import ChatOpenAI
import time
import json
# 监控指标
class Metrics:
def __init__(self):
self.total_rounds = 0
self.total_tokens = 0
self.start_time = None
self.end_time = None
def start(self):
self.start_time = time.time()
def end(self):
self.end_time = time.time()
def add_round(self, tokens):
self.total_rounds += 1
self.total_tokens += tokens
def report(self):
duration = self.end_time – self.start_time
return {
"total_rounds": self.total_rounds,
"total_tokens": self.total_tokens,
"duration_seconds": duration,
"estimated_cost": self.total_tokens * 0.00015 / 1000 # GPT-4o-mini 价格
}
metrics = Metrics()
# 结构化输出
class Review(BaseModel):
passed: bool = Field(description="文案是否合格")
feedback: str = Field(description="改进建议")
score: int = Field(description="评分 1-10", ge=1, le=10)
class State(TypedDict):
draft: str
feedback: str
passed: bool
rounds: int
final_score: int
token_usage: list # 记录每轮 token 消耗
MAX_ROUNDS = 3
RECURSION_LIMIT = 25
TOKEN_LIMIT = 5000 # 总 token 上限
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.7)
def generate(state: State):
"""生成节点"""
current_round = state.get("rounds", 0) + 1
# 检查 token 上限
total_tokens = sum(state.get("token_usage", []))
if total_tokens >= TOKEN_LIMIT:
logger.warning(f"总 token 消耗 {total_tokens} 达到上限 {TOKEN_LIMIT},强制退出")
return {
"draft": state.get("draft", "生成失败"),
"rounds": current_round,
"passed": True # 强制退出
}
prompt = "写一句关于人工智能的营销文案,要求简洁有力"
if state.get("feedback"):
if current_round >= MAX_ROUNDS – 1:
prompt = f"这是最后一次修改机会。请尽量优化文案。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
else:
prompt = f"根据反馈修改文案。\\n反馈:{state['feedback']}\\n原文:{state['draft']}"
response = llm.invoke(prompt)
# 估算 token 消耗
estimated_tokens = len(prompt) + len(response.content)
metrics.add_round(estimated_tokens)
return {
"draft": response.content,
"rounds": current_round,
"token_usage": state.get("token_usage", []) + [estimated_tokens]
}
def evaluate(state: State):
"""评估节点"""
try:
review = llm.with_structured_output(Review).invoke(
f"评审这版文案是否合格:{state['draft']}"
)
return {
"passed": review.passed,
"feedback": review.feedback,
"final_score": review.score
}
except Exception as e:
logger.error(f"评估失败: {e}")
return {
"passed": True,
"feedback": "评估失败,自动通过",
"final_score": 5
}
def should_continue(state: State) -> Literal["generate", "__end__"]:
"""条件判断——四层刹车"""
current_round = state.get("rounds", 0)
total_tokens = sum(state.get("token_usage", []))
# 第一层:轮数上限
if current_round >= MAX_ROUNDS:
logger.warning(f"轮数达到上限 {MAX_ROUNDS}")
return "__end__"
# 第二层:token 上限
if total_tokens >= TOKEN_LIMIT:
logger.warning(f"token 消耗达到上限 {TOKEN_LIMIT}")
return "__end__"
# 第三层:评估通过
if state["passed"]:
return "__end__"
return "generate"
# 构建图
builder = StateGraph(State)
builder.add_node("generate", generate)
builder.add_node("evaluate", evaluate)
builder.add_edge(START, "generate")
builder.add_edge("generate", "evaluate")
builder.add_conditional_edges("evaluate", should_continue)
graph = builder.compile(checkpointer=InMemorySaver())
# 执行并监控
metrics.start()
config = {
"configurable": {"thread_id": "eval-prod-001"},
"recursion_limit": RECURSION_LIMIT
}
try:
result = graph.invoke({
"draft": "",
"feedback": "",
"passed": False,
"rounds": 0,
"final_score": 0,
"token_usage": []
}, config)
metrics.end()
report = metrics.report()
print(f"\\n=== 执行报告 ===")
print(f"总轮数: {report['total_rounds']}")
print(f"总 token 消耗: {report['total_tokens']}")
print(f"耗时: {report['duration_seconds']:.2f} 秒")
print(f"预估成本: ${report['estimated_cost']:.6f}")
print(f"\\n最终结果:")
print(f"文案: {result['draft']}")
print(f"最终评分: {result['final_score']}")
# 告警检查
if report['total_rounds'] >= MAX_ROUNDS:
print(f"⚠️ 告警:达到最大轮数 {MAX_ROUNDS},请检查评估器效果")
if report['total_tokens'] > 3000:
print(f"⚠️ 告警:token 消耗较高 ({report['total_tokens']}),建议优化 prompt")
except Exception as e:
metrics.end()
print(f"❌ 执行异常: {e}")
print(f"异常前消耗: {metrics.total_tokens} tokens")
三种刹车机制对比
| 第一层:主动刹车 | 轮数计数器 | state["rounds"] >= MAX_ROUNDS | 精确控制最大轮数,可预测成本 |
| 第二层:结构化输出 | with_structured_output | Review.passed: bool | 避免自然语言解析失败导致的死循环 |
| 第三层:被动刹车 | recursion_limit | config={"recursion_limit": 25} | 兜底,防止意外情况 |
生产环境最佳实践
1. 配置化参数
class EvalOptimizerConfig:
"""评估-优化器配置"""
max_rounds: int = 3 # 最大轮数
max_tokens: int = 5000 # 最大 token 消耗
recursion_limit: int = 25 # 递归限制
fallback_on_failure: bool = True # 评估失败时是否通过
log_level: str = "INFO" # 日志级别
alert_threshold: float = 0.8 # 告警阈值(达到 max_rounds 的百分比)
2. 监控告警
def check_alert(metrics: dict, config: EvalOptimizerConfig):
"""检查是否需要告警"""
alerts = []
# 轮数告警
if metrics["total_rounds"] >= config.max_rounds * config.alert_threshold:
alerts.append({
"level": "WARNING",
"message": f"轮数接近上限: {metrics['total_rounds']}/{config.max_rounds}"
})
# token 告警
if metrics["total_tokens"] >= config.max_tokens * config.alert_threshold:
alerts.append({
"level": "WARNING",
"message": f"token 消耗接近上限: {metrics['total_tokens']}/{config.max_tokens}"
})
# 时长告警
if metrics["duration_seconds"] > 30:
alerts.append({
"level": "WARNING",
"message": f"执行时间过长: {metrics['duration_seconds']:.2f} 秒"
})
return alerts
3. 单元测试
def test_evaluator_optimizer():
"""测试评估-优化器循环"""
# 测试场景1:正常通过
state = {"draft": "好文案", "rounds": 1, "passed": True}
assert should_continue(state) == "__end__"
# 测试场景2:达到轮数上限
state = {"draft": "差文案", "rounds": 3, "passed": False}
assert should_continue(state) == "__end__"
# 测试场景3:继续循环
state = {"draft": "一般文案", "rounds": 1, "passed": False}
assert should_continue(state) == "generate"
# 测试场景4:结构化输出
review = Review(passed=True, feedback="优秀", score=9)
assert review.passed == True
assert review.score >= 1 and review.score <= 10
print("所有测试通过!")
总结
Evaluator-optimizer 循环的三层刹车是生产环境的生命线:
| 轮数计数器 | 精确控制成本,可预测 | 无限制循环,烧光预算 |
| 结构化输出 | 避免解析失败 | 死循环,CPU 100% |
| recursion_limit | 兜底保护 | 系统崩溃,服务不可用 |
关键修复代码:
# 1. 状态中加入轮数计数器
class State(TypedDict):
rounds: int # 轮数计数器
# 2. 主动刹车
def should_continue(state: State):
if state["rounds"] >= MAX_ROUNDS:
return "__end__"
# 3. 结构化输出
class Review(BaseModel):
passed: bool = Field(description="是否合格")
# 4. 被动刹车
config = {"recursion_limit": 25}
一句话总结:没有刹车的循环是定时炸弹——轮数计数器是油门踏板,recursion_limit 是刹车片,结构化输出是方向盘,三者缺一不可。






