难度:★★★(需要第 08、09 课基础)
教材:worry_debate_game/worry_debate_game/app.py(294 行)+nodes.py(355 行)
预计时间:讲解 70 分钟 + 练习 50 分钟
学完你能:把"玩具级流式接口"升级成"能多轮、防并发、失败可重试、还能演剧情"的会话式接口
第 09 课的流式是"发一次请求,收一串事件"。
但真实的一局辩论要持续好几轮,于是冒出一堆新问题:
这一课逐个解决。
cd worry_debate_game
start.bat
观察错误提示后还能不能重试同一轮
app.py
├── Session 类 第 61-67 行 一局游戏的状态 + 锁
├── sessions / _sessions_lock 第 70-71 行 所有会话 + 保护它的锁
├── _purge_sessions() 第 74-82 行 过期清理 + 超量淘汰
├── _play_round() 第 123-161 行 一轮演出(天使→打断→恶魔→裁判→结束)
├── /api/start-stream 第 184-243 行 开局(含快照与回滚)
└── /api/next-stream 第 246-294 行 下一轮(含并发锁与回滚)
nodes.py
├── resolve_api_key() 第 32-36 行 Key 的优先级
├── DemoLLM 第 94-120 行 无 Key 也能离线演示
├── _transcript() 第 176-189 行 把历史拼成上下文
├── _demon_messages() 第 207-225 行 含"打断后第一句要嚣张"的指令
├── _stream_text() 第 240-258 行 流式 + 截断(cut_at)
└── evaluate_round() 第 312-349 行 裁判打分 + 容错
Session 对象 —— 一局游戏装进一个类(第 61-67 行)class Session:
def __init__(self, state: GameState, allow_interrupt: bool, judge: bool):
self.state = state # 游戏数据(烦恼、历史台词、倾向值…)
self.allow_interrupt = allow_interrupt # 这一局允不允许打断
self.judge = judge # 要不要裁判打分
self.touched = time.time() # 最后一次活动时间(用来清理)
self.lock = threading.Lock() # 这一局的"独占锁"
对比第 08 课:那时会话就是裸的 dict;现在包成了一个类,因为要挂配置(打断/裁判)
和运行期资源(锁、活动时间)。状态多了,就值得用对象封装。
SESSION_TTL = 2 * 60 * 60 # 2 小时无活动即清理
MAX_SESSIONS = 500
def _purge_sessions():
now = time.time()
with _sessions_lock:
expired = [sid for sid, s in sessions.items() if now - s.touched > SESSION_TTL]
for sid in expired:
sessions.pop(sid, None)
if len(sessions) > MAX_SESSIONS:
for sid, _ in sorted(sessions.items(), key=lambda kv: kv[1].touched)[: len(sessions) - MAX_SESSIONS]:
sessions.pop(sid, None)
内存字典当存储的老问题:只进不出会爆内存。这里用两级策略:
now - touched > TTL)touched 从旧到新删到 500(sorted(...)[:超出数量] 就是把最旧的挑出来)
每次开局前调用一次(第 187 行 _purge_sessions()),
不用定时器也能保持内存不涨——这种"顺手清理"是轻量服务常用的小技巧。
真·生产的做法是把会话放 Redis(自动过期、多进程共享)。
但能说清"内存 + TTL + LRU 的适用边界",本身就是工程判断力。
@app.post("/api/next-stream")
def next_round_stream(req: NextRequest, x_api_key: Optional[str] = Header(None)):
session = sessions.get(req.session_id)
def generate():
if session is None:
yield sse("error", {"message": "会话不存在或已过期,请重新开始", "retryable": False})
return
if not session.lock.acquire(blocking=False): # ← 抢锁,抢不到立刻返回
yield sse("error", {"message": "上一轮还在进行中", "retryable": True})
return
...
finally:
session.lock.release() # ← 无论如何都要还锁
acquire(blocking=False):尝试拿锁,拿不到不等待,直接告诉用户"上一轮还在跑"finally: release():保证锁一定被释放,否则会话永久卡死retryable: True 告诉前端"这个错误等一会儿重试就行"为什么不用 blocking=True 排队?因为流式生成可能几十秒,
让第二个请求干等没意义,不如立刻反馈。前端也更好做提示。
先看开局(第 225-240 行):
snapshot = copy.deepcopy(state) # 开跑前拍一张"快照"
snapshot["round"] = 0
completed = False
with session.lock:
try:
yield from _play_round(session, session_id)
completed = True
except Exception as e:
state.clear()
state.update(snapshot) # 回滚:把状态恢复到快照
completed = True
yield sse("error", {...})
finally:
if not completed: # 客户端中途断开(生成器被关闭)
state.clear()
state.update(snapshot) # 同样回滚
session.touched = time.time()
三个关键点:
copy.deepcopy:深拷贝——快照和原状态完全独立,改原状态不会影响快照(浅拷贝会让嵌套的列表/字典还是同一份)
state.clear() + state.update(snapshot):注意不能用 state = snapshot(那只是换了变量指向,外面的会话对象还指着旧字典)
finally 里的 completed 判断:处理"用户关页面导致生成器被提前关闭"的情况——这时 except 不会执行,但状态已经被改了一半,所以也要回滚
为什么流式接口必须回滚?
HTTP 流一旦开始发送(200 + 数据已出门),就没法再改状态码。
失败只能以 error 事件的形式出现在流里。这时如果状态被改了半截,
用户重试就会基于错误数据——回滚让"重试"变得安全。
nodes.py 第 240-258 行)def _stream_text(llm, messages, cut_at=None):
stream = llm.stream(messages)
total = 0
try:
for chunk in stream:
content = chunk.content
if not content:
continue
if cut_at is not None and total + len(content) >= cut_at:
head = content[: max(0, cut_at - total)] # 截到指定长度
yield head, True # 标记"被截断了"
return
total += len(content)
yield content, False
finally:
close = getattr(stream, "close", None)
if close:
close() # 提前关掉模型连接,不再继续生成(省钱!)
打断的实现思路很巧:不需要"另一个进程去砍断",
只要在流式循环里累计长度,到点就停止读取并关闭流。
再看随机触发(app.py 第 128-138 行):
cut_at = None
if session.allow_interrupt:
chance = 0.3 if round_num == 1 else 0.5
if random.random() < chance:
cut_at = random.randint(38, 70)
yield sse("angel_start", {"round": round_num})
for chunk in stream_angel(state, cut_at=cut_at):
yield sse("angel", {"text": chunk})
if state.get("interrupted") and state["interrupted"][-1]:
yield sse("interrupt", {"by": "demon"})
yield sse("angel_end", {})
interrupt 事件,前端据此播放"恶魔插嘴"的动效而恶魔的 prompt 会配合剧情(_demon_messages 第 220-221 行):
if interrupted:
user += "\n注意:天使的话还没说完就被你粗暴地打断了。你的第一句必须是一句短促、嚣张的打断(例如"停停停!""够了!"),然后再展开攻击。\n"
这是本课最值得抄走的设计:状态里记下"发生了打断",
再把这件事翻译成 prompt 指令,让下一次生成延续剧情——
技术事件变成了叙事元素。
nodes.py 第 302-349 行)_JSON_RE = re.compile(r"\{.*\}", re.S)
def _clamp_int(value, lo, hi, default):
try:
return max(lo, min(hi, int(round(float(value)))))
except (TypeError, ValueError):
return default
def evaluate_round(state):
...
angel_score, demon_score, intensity, comment = 5, 5, 5, "势均力敌,胜负未分。"
try:
response = _llm_for(state, temperature=0.3).invoke([...])
match = _JSON_RE.search(response.content or "")
data = json.loads(match.group(0)) if match else {}
angel_score = _clamp_int(data.get("angel"), 0, 10, 5)
...
except Exception:
pass # 出错就用默认分,游戏继续
shift = (angel_score - demon_score) * 4
new_tendency = max(-100, min(100, state.get("tendency", 0) + shift))
和第 10 课的 evaluate_tendency 相比,这里更成熟:
\{.*\} 先"抠出" JSON:模型就算前后加了废话,也能捞到花括号里的部分_clamp_int 一个函数搞定四件事:转数字、四舍五入、限制范围、失败给默认值nodes.py 第 32-36 行)def resolve_api_key(explicit: Optional[str] = None) -> str:
key = (explicit or "").strip() or _api_key_override or os.getenv("DEEPSEEK_API_KEY")
if not key:
raise ValueError("请先在右上角设置中填写 DeepSeek API Key")
return key
优先级:请求头带的 X-API-Key > 服务端 set-key 设置的 > 环境变量。
前端每次请求带自己的 Key(client.ts 第 61 行),
服务端不必保存任何人的 Key——这也是"把 Key 交给用户自己保管"的一种产品设计。
| 事件 | 第 09 课 | 本课新增 |
|---|---|---|
| 会话 | 无(用内存 dict) | session(返回 session_id) |
| 边界 | start/end | 同左 + interrupt |
| 判定 | round_done | judge / judge_start |
| 结束 | round_done 里带 done | 单独的 game_over(带 winner) |
| 错误 | error | error + retryable 标记 |
协议是会长大的。设计时预留"边界事件"(start/end)和"可扩展字段",
后面加功能就不用改前端解析逻辑。
| 方案 | 特点 | 适用 |
|---|---|---|
| 内存 dict + TTL/LRU(你的) | 零依赖、单进程、重启即丢 | 单机原型 |
| Redis | 过期自动、多进程共享、可持久化 | 生产环境 |
| 数据库 | 强持久、可审计 | 需要留档 |
next-stream 用锁 + 快照逼近这一点)commit/rollback 是数据库版;这里是内存版的同一思想——出错时回到干净状态
blocking=False 就是这种态度)练习 1(热身):关掉打断做对比
把前端开局的"允许打断"关掉(或在 StartRequest 里把 allow_interrupt 默认改成 False),
跑一局,在 Network 里对比事件序列——少了哪个事件?
练习 2(必做):调整打断参数
把 app.py 第 130-132 行的概率(0.3/0.5)和截断点(38-70)改大,
让打断更频繁、更早,观察恶魔的"第一句嚣张"是不是更常见了。
练习 3(必做):亲手制造一次失败并重试
故意把 API Key 改错(设置里),开始一局 → 触发 error → 观察前端是否提示可重试;
用正确的 Key 再点"下一轮",确认这局还能继续(说明回滚生效了)。
如果没回滚,重试时会发生什么?说说看。
练习 4(必做):并发实验
连续快速点击"下一轮"两次,观察第二次收到的错误消息。
然后在 app.py 第 255 行把 blocking=False 改成 blocking=True,
再试一次,感受两种策略的区别(一个立刻报"进行中",一个默默等待)。
练习 5(思考题):
为什么回滚用 state.clear(); state.update(snapshot) 而不是 state = snapshot?
session.state 和函数里的 state 是什么关系?
练习 6(选做,进阶):
给 _play_round 加一个新事件 mood:当裁判 intensity ≥ 8 时,
额外 yield sse("mood", {"level": "high"})。前端在 ServerEvent 里加类型并处理它
(可以只是 console.log)。做完你就把"后端状态 → 前端表现"打通了一遍。
Session 类封装:数据 + 配置 + 锁 + 活动时间lock.acquire(blocking=False):抢不到立刻反馈,finally 保证释放clear/update → 重试安全;客户端断开也要回滚(finally 里的 completed 判断)
close() 关连接;事件化(interrupt)+ prompt 化(让下一句接戏)| 术语 | 人话解释 |
|---|---|
| 会话(session) | 一段有状态、跨多次请求的交互 |
| TTL | 生存时间,超时自动失效 |
| LRU | 最久未使用优先淘汰 |
| 快照(snapshot) | 某一刻状态的完整副本,用于回滚 |
| 深拷贝 / 浅拷贝 | 副本完全独立 / 只复制了外层 |
| 回滚(rollback) | 出错时恢复到之前的状态 |
| 可重入 | 同一个操作还没结束又被调用 |
| 幂等 | 做一次和做多次结果相同 |
| 可辨识联合 | 用某个字段区分多种类型的 TypeScript 写法 |
第 15 课:结构化输出实战:AI 总结 + 信源提取 + 关联图谱。
回到「查资料」,看它的新功能:把多轮对话一键"总结成卡片"(/api/summarize)、
从回答里自动提取信源链接、以及怎么用一张邻接表高效查出卡片之间的关联。
这课讲的是"把 AI 的自由文本变成结构化数据"的完整套路。