异步压缩:让用户感知不到上下文整理
专栏信息
《从零到一构建跨平台 AI 助手:WeClaw 实战指南》专栏
本文是模块八第 8 篇,讲解非阻塞式后台摘要生成的架构设计。
作者与项目
作者简介:翁勇刚 WENG YONGGANG 新概念龙虾-WeClaw 开发团队负责人,一群专注于跨平台 AI 应用的实践者 理念:"再复杂的技术,也能用代码讲清楚"
- 项目地址:https://github.com/wyg5208/weclaw.git
- 官网地址:https://weclaw.link
- 作者 CSDN:https://blog.csdn.net/yweng18
摘要
本文结构概览: 本文从同步压缩导致用户可感知延迟的问题出发,设计异步后台压缩架构,重点讲解快照 hash 保护机制、Task 生命周期管理,以及 qasync 事件循环约束下的特殊处理。
背景:上下文压缩需要调用 LLM 生成摘要,这个过程通常需要 2-10 秒。如果在主对话循环中同步执行,用户会明显感受到"AI 卡住了"。
核心问题:如何在不阻塞主对话的前提下完成上下文压缩,并确保压缩结果不会覆盖已被用户修改的上下文?
解决方案:异步后台 Task + 快照 hash 校验 + 会话清理联动
关键成果:
- 用户感知零延迟:压缩在后台完成,主循环不阻塞
- 数据安全:快照 hash 防止过期结果回写
- 资源管理:会话清理时自动 cancel 未完成的 Task
适合读者:异步编程经验丰富的 Python 开发者,关注 Agent 性能优化
阅读时长:约 10 分钟
关键词:异步压缩、后台Task、快照Hash、非阻塞、asyncio
一、同步压缩的用户体验问题
1.1 用户感知
用户: 帮我分析这个项目的代码结构
AI: [调用 read_file, 分析代码...] ← 正常速度
AI: 以下是分析结果...
[对话持续了 30 轮,触发了压缩]
用户: 好的,再帮我看看测试覆盖率
AI: [调用工具...] ← 正常
AI: [生成回复...] ← 正常
用户: 继续
AI: ⏳ ... (3秒无响应) ← 压缩正在进行!
AI: ⏳ ... (5秒无响应) ← 还在压缩!
AI: 好的,让我继续分析... ← 压缩完成,终于可以回复了
用户在这 5-8 秒内不知道发生了什么,以为系统卡死了。
1.2 延迟来源
| 步骤 | 耗时 | 说明 |
|---|---|---|
| 计算 token 数 | ~5ms | 可忽略 |
| 构建摘要 prompt | ~10ms | 可忽略 |
| LLM 摘要生成 | 2-10s | 主要延迟 |
| 替换消息列表 | ~5ms | 可忽略 |
核心瓶颈是 LLM 摘要生成。这不是代码能优化的——必须等待外部 API 返回。
二、异步架构设计
[图片: 同步 vs 异步压缩时序对比 | 生成方式: 文生图 PROMPT: "Two swimlane diagrams comparing synchronous vs asynchronous context compression. Top diagram (Synchronous): User sends message, waits while AI compresses context (5s delay shown as yellow block), then gets response. Bottom diagram (Asynchronous): User sends message, AI responds immediately using truncated messages, background compression runs in parallel, next user message uses the fresh summary. Timeline with three lanes: User, AI Main Loop, Background Task"]
2.1 核心思路
当前轮(不阻塞):
1. 检测到需要压缩
2. 用截断结果立即返回(保证对话不中断)
3. 后台启动 Task 生成摘要
下一轮(自动使用):
4. 后台 Task 完成,摘要写入
5. 用户发送下一条消息时,自动使用新的摘要
2.2 架构实现
class DialogManager:
"""对话管理器——支持异步后台压缩"""
def __init__(self):
self._sessions = {}
# Task 管理:session_id → asyncio.Task
self._pending_tasks: dict[str, asyncio.Task] = {}
async def get_messages(self, session_id):
"""获取消息列表(主路径——不阻塞)"""
session = self._sessions[session_id]
messages = session.messages
# 检查是否需要压缩
if self._should_compress(messages):
# 先用截断结果返回(不阻塞用户)
truncated = self._quick_truncate(messages)
# 后台启动摘要生成
self._launch_background_summary(session_id, messages)
return truncated
return messages
def _launch_background_summary(self, session_id, messages):
"""启动后台摘要生成 Task"""
# 取消之前的 Task(如果还在运行)
if session_id in self._pending_tasks:
self._pending_tasks[session_id].cancel()
# 创建新 Task
task = asyncio.create_task(
self._background_compress(session_id, messages)
)
self._pending_tasks[session_id] = task
# 异常处理回调
task.add_done_callback(
lambda t: self._on_compress_done(session_id, t)
)
三、快照 Hash 保护
3.1 竞态风险
异步压缩面临一个核心竞态问题:
时间线:
T1: 用户发送消息,后台启动压缩(基于 50 条消息的快照)
T2: 用户继续对话,消息列表变为 52 条(新增了 2 条)
T3: 后台压缩完成,基于 50 条消息生成了摘要
T4: ⚠️ 如果直接写入摘要,会覆盖 T2 新增的 2 条消息!
3.2 快照 Hash 机制
async def _background_compress(self, session_id, snapshot_messages):
"""后台压缩——带快照 hash 保护"""
# Step 1: 启动时记录快照信息
snapshot_hash = self._compute_snapshot_hash(snapshot_messages)
snapshot_count = len(snapshot_messages)
# Step 2: 执行压缩(耗时 2-10s,不阻塞主循环)
try:
summary = await self._compress_engine.compress(snapshot_messages)
except asyncio.CancelledError:
logger.info(f"Compression cancelled for {session_id}")
return
except Exception as e:
logger.error(f"Background compression failed: {e}")
return
# Step 3: 完成时校验快照是否过期
session = self._sessions.get(session_id)
if not session:
return # 会话已被删除
current_hash = self._compute_snapshot_hash(session.messages)
if current_hash != snapshot_hash:
logger.warning(
f"Snapshot expired: started with {snapshot_count} msgs, "
f"now has {len(session.messages)} msgs. Discarding summary."
)
return # 快照已过期,丢弃压缩结果
# Step 4: 快照有效,写入摘要
session.apply_summary(summary)
logger.info(f"Background compression applied for {session_id}")
def _compute_snapshot_hash(self, messages):
"""计算消息快照的 hash 值"""
# 使用消息数量和最后一条消息的内容作为 hash
# 轻量级方案,不需要对所有消息做完整 hash
if not messages:
return "empty"
last_msg = messages[-1]
content = str(last_msg.get("content", ""))[:200]
return f"{len(messages)}:{hash(content)}"
3.3 Hash 策略的选择
| 策略 | 精度 | 性能 | 适用场景 |
|---|---|---|---|
| 全量 MD5 | 最高 | 慢(遍历所有消息) | 小消息列表 |
| 数量+末条内容 | 高 | 快(O(1)) | 大多数场景 |
| 仅数量 | 低 | 最快 | 粗糙校验 |
我们选择数量+末条内容:如果消息列表发生了变化(新增或删除),末条内容或数量一定会变。
四、Task 生命周期管理
4.1 Task 状态图
┌─────────┐
│ Created │
└────┬────┘
│ asyncio.create_task()
▼
┌─────────┐
│ Running │ ← 压缩进行中
└────┬────┘
┌─────┼─────┐
│ │ │
cancel() done exception
│ │ │
▼ ▼ ▼
┌────────┐ ┌────┐ ┌──────┐
│Cancelled│ │Done│ │Failed│
└────────┘ └────┘ └──────┘
4.2 会话清理联动
def clear_messages(self, session_id):
"""清除会话消息——自动 cancel 未完成的压缩 Task"""
if session_id in self._pending_tasks:
task = self._pending_tasks.pop(session_id)
if not task.done():
task.cancel()
logger.info(f"Cancelled pending compression for {session_id}")
session = self._sessions.get(session_id)
if session:
session.messages.clear()
def delete_session(self, session_id):
"""删除会话——同样清理 Task"""
if session_id in self._pending_tasks:
task = self._pending_tasks.pop(session_id)
if not task.done():
task.cancel()
self._sessions.pop(session_id, None)
4.3 异常处理回调
def _on_compress_done(self, session_id, task):
"""Task 完成回调——处理异常和清理"""
# 清理 Task 引用
self._pending_tasks.pop(session_id, None)
if task.cancelled():
return # 正常取消,不记录
exc = task.exception()
if exc:
logger.error(
f"Background compression error for {session_id}: {exc}"
)
五、qasync 事件循环约束
5.1 特殊约束
WeClaw 使用 qasync 将 asyncio 事件循环集成到 Qt 的事件循环中。这意味着:
- 不能创建新线程:所有异步操作必须在主事件循环中
- 不能使用
run_in_executor:线程池执行器会与 Qt 事件循环冲突 - LLM API 调用必须是纯异步:使用
aiohttp而非requests
5.2 AuxiliaryClient 的可变状态
辅助模型客户端(AuxiliaryClient)包含可变状态(如冷却期时间戳)。在异步环境中,多个 Task 可能同时访问这个状态:
# 潜在竞态:两个会话同时触发压缩
# Task A: 检查 auxiliary.cooldown → 未冷却 → 调用 API
# Task B: 检查 auxiliary.cooldown → 未冷却 → 调用 API
# 两个 Task 都调用了 API,浪费资源
# 解决方案:在 AuxiliaryClient 中使用 asyncio.Lock
class AuxiliaryClient:
def __init__(self):
self._lock = asyncio.Lock()
self.cooldown_until = 0
async def summarize(self, messages):
async with self._lock:
if time.time() < self.cooldown_until:
raise CooldownError()
try:
result = await self._call_api(messages)
return result
except Exception:
self.cooldown_until = time.time() + 300
raise
六、性能对比
6.1 同步 vs 异步
| 指标 | 同步压缩 | 异步压缩 |
|---|---|---|
| 用户感知延迟 | +2-10s | +0s |
| 压缩质量 | 完整摘要 | 当前轮截断 + 下轮摘要 |
| 消息一致性 | 保证 | 快照 hash 保护 |
| 实现复杂度 | 低 | 中 |
| 资源消耗 | 阻塞主循环 | 后台并行 |
6.2 "当前轮截断"的质量影响
异步压缩意味着当前轮使用的是截断结果而非摘要。这个质量损失有多大?
实际影响很小:
- 截断只是移除了最旧的消息,最近的消息完整保留
- 下一轮(通常 5-30 秒后)自动使用摘要,质量恢复
- 用户几乎感知不到当前轮的质量差异
七、总结与展望
7.1 核心要点回顾
- 异步压缩消除用户感知延迟:当前轮截断、后台摘要、下轮自动使用
- 快照 hash 保护数据一致性:防止过期压缩结果覆盖已变更的上下文
- Task 生命周期与会话绑定:清理会话时自动 cancel 未完成的 Task
- qasync 约束需要特别注意:不能创建新线程,LLM 调用必须纯异步
下期预告:《密钥脱敏:你的 API Key 可能藏在摘要里》
- 输出端脱敏 vs 输入端脱敏的选择
- 供应商前缀正则匹配设计
- 代码块保护的坑与解法
敬请期待!
版权声明:本文为 CSDN 博主「翁勇刚」的原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接及本声明。