从开发到生产——监控、容错、扩展、成本控制
一、引言
图在本地跑通了,不代表能在生产环境稳定运行。
本章讲四个生产问题:监控、容错、扩展、成本。
二、监控:你得能看见图在干什么
2.1 执行追踪
每个节点的输入、输出、耗时、状态都要记录。
`python
import time
import logging
async def traced_node(name, func, state):
start = time.time()
try:
result = await func(state)
duration = time.time() – start
logging.info(f”[{name}] ✅ {duration:.2f}s”)
return result
except Exception as e:
duration = time.time() – start
logging.error(f”[{name}] ❌ {duration:.2f}s: {e}”)
raise
`
2.2 状态快照
图的每个状态转换都要可回溯。
`python
# LangGraph内置支持
app = graph.compile(checkpointer=SqliteSaver.from_conn_string(“checkpoints.db”))
# 可以从任意状态快照恢复
config = {“configurable”: {“thread_id”: “user-123”}}
result = app.invoke(initial_state, config)
`
2.3 告警规则
| 指标 | 阈值 | 动作 |
|---|
|——|——|——|
| 节点执行时间 | > 30s | 告警 |
|---|---|---|
| 失败率 | > 10% | 暂停图 |
| Token消耗 | > 预算80% | 降级模型 |
| 队列积压 | > 100 | 扩容 |
三、容错:图必须能从失败中恢复
3.1 节点级容错
`python
# 重试装饰器
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3), wait=wait_exponential(min=1, max=10))
async def extract_with_retry(content):
return await llm.invoke(content)
`
3.2 图级容错
`python
# 检查点恢复
from langgraph.checkpoint.sqlite import SqliteSaver
checkpointer = SqliteSaver.from_conn_string(“graph_checkpoints.db”)
app = graph.compile(checkpointer=checkpointer)
# 图崩溃后,从最后的检查点恢复
state = app.get_state(config)
if state.next:
app.invoke(None, config) # 继续执行
`
3.3 数据级容错
`python
# 每个节点的输出都要可重放
class IdempotentNode:
def __init__(self, name):
self.name = name
self.processed = set()
async def run(self, state):
task_id = state[“task_id”]
if task_id in self.processed:
return state # 跳过已处理的
result = await self._process(state)
self.processed.add(task_id)
return result
`
四、扩展:从单机到分布式
4.1 并行度控制
`python
# 限制并发数
import asyncio
semaphore = asyncio.Semaphore(4) # 最多4个并行节点
async def limited_node(func, state):
async with semaphore:
return await func(state)
`
4.2 队列化节点
`python
# 长耗时节点用队列
from celery import Celery
app = Celery(‘graph_tasks’, broker=’redis://localhost’)
@app.task
def long_running_task(data):
# 耗时处理
return result
# 图节点调用Celery任务
async def queue_node(state):
result = long_running_task.delay(state[“data”])
return {“task_id”: result.id, “status”: “queued”}
`
4.3 分布式状态
`python
# 用Redis存储共享状态
import redis
r = redis.Redis()
def save_state(graph_id, state):
r.set(f”graph:{graph_id}”, json.dumps(state))
def load_state(graph_id):
return json.loads(r.get(f”graph:{graph_id}”))
`
五、成本控制:别让图烧光预算
5.1 模型分层
`python
# 重复性工作用便宜模型
cheap_model = ChatOpenAI(model=”gpt-4o-mini”, temperature=0)
expensive_model = ChatOpenAI(model=”gpt-4o”, temperature=0)
async def classify_node(state):
return await cheap_model.invoke(state[“content”]) # 分类用便宜模型
async def synthesize_node(state):
return await expensive_model.invoke(state[“content”]) # 合成用强模型
`
5.2 Token预算
`python
class TokenBudget:
def __init__(self, max_tokens=100000):
self.max_tokens = max_tokens
self.used = 0
def check(self, estimated):
if self.used + estimated > self.max_tokens:
raise BudgetExceeded(f”Used {self.used}/{self.max_tokens}”)
self.used += estimated
`
5.3 缓存
`python
# 相同输入不重复调用LLM
from functools import lru_cache
@lru_cache(maxsize=1000)
async def cached_llm_call(content_hash):
return await llm.invoke(content)
`
六、生产检查清单
部署前
- [ ] 每个节点都有重试机制
- [ ] 重试耗尽有明确的失败路由
- [ ] 图有检查点,可断点恢复
- [ ] Token预算有上限
- [ ] 并行度有控制
运行时
- [ ] 每个节点执行时间有监控
- [ ] 失败率有告警
- [ ] 状态变更有日志
- [ ] 成本有实时统计
故障时
- [ ] 能从检查点恢复
- [ ] 能手动干预(Human-in-the-loop)
- [ ] 能降级运行(跳过非关键节点)
- [ ] 有回滚方案
七、总结
`
开发环境: 图能跑通就行
生产环境: 监控 + 容错 + 扩展 + 成本控制
`
从开发到生产的鸿沟,不是功能,是可靠性。
系列总结
`
Day 1: 什么是Graph Engineering——从Loop到Graph的进化
Day 2: 四个核心构件——Nodes、Edges、Shared State、Failure Routing
Day 3: 五种基础拓扑——链式、扇出、菱形、条件路由、循环
Day 4: 框架选型——LangGraph vs CrewAI vs AutoGen vs 原生Python
Day 5: 图编排实操手册——完整代码实现
Day 6: 生产环境部署——监控、容错、扩展、成本控制
`
Graph Engineering不是Loop的替代品,是Loop的组织方式。
先把一个Loop跑稳,再考虑建图。
发表回复