AI PRO·Graph Day 6 生产环境部署

作者:

从开发到生产——监控、容错、扩展、成本控制


一、引言

图在本地跑通了,不代表能在生产环境稳定运行。

本章讲四个生产问题:监控、容错、扩展、成本。


二、监控:你得能看见图在干什么

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跑稳,再考虑建图。

评论

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注