OpenClaw 学习系列之十一:并发控制系统
并发控制系统
📚 学习路径:本文档是架构文档的第 9 部分。建议按顺序阅读:
- 01. 技术基础 - 了解 TypeScript 技术特性
- 02. 整体框架 - 理解 OpenClaw 架构
- 03. 消息流转 - 掌握消息生命周期
- 04. 设计原理 - 理解反共识设计
- 05. Gateway 深度解析 - 理解控制平面
- 06. Agent 运行机制 - 了解 Agent 如何工作
- 07. 会话管理 - 理解会话生命周期
- 08. 记忆系统 - 了解长期记忆管理
前言
并发控制系统是 OpenClaw 的核心组件之一,负责管理消息的并发处理和队列调度。本文档将深入解析并发控制系统的设计、队列模式和调度策略。
通过阅读本文档,你将能够:
- 理解并发控制的作用和设计理念
- 掌握队列模式和调度策略
- 了解两层并发控制机制
- 理解防抖和合并机制
一、并发控制概述
1.1 为什么需要并发控制?
| 问题 | 说明 | 解决方案 |
|---|---|---|
| 消息竞争 | 多个消息同时处理导致状态混乱 | 队列系统 |
| 资源耗尽 | 并发过多导致资源耗尽 | 并发限制 |
| 状态混乱 | 并发修改导致状态不一致 | 会话隔离 |
| 成本控制 | 并发 API 调用导致成本过高 | 资源管理 |
1.2 并发控制的目标
| 目标 | 说明 | 实现方式 |
|---|---|---|
| 用户体验 | 避免消息乱序,支持实时打断 | 队列模式 |
| 资源控制 | 防止系统过载,控制 API 调用 | 并发限制 |
| 状态一致性 | 避免并发修改会话状态 | 会话隔离 |
| 成本优化 | 减少 API 调用,优化资源使用 | 智能调度 |
二、两层并发控制
2.1 架构设计
2.2 第一层:会话级别
原则:同一会话内的消息必须串行处理
原因:
- 避免状态混乱(两个消息同时修改同一会话状态)
- 保持对话顺序(消息按到达顺序处理)
- 防止竞态条件(Race Condition)
代码位置:src/agents/pi-embedded-runner/run/lanes.ts
export function resolveSessionLane(sessionId: string): string {
return `session:${sessionId}`;
}
// 同一会话的消息进入同一个队列
const enqueueSession = (task) => enqueueCommandInLane(sessionLane, task);
2.3 第二层:全局级别
原则:系统整体并发度限制(默认 4)
原因:
- 防止资源耗尽(内存、CPU、API 调用限制)
- 控制成本(避免同时发起大量 AI 调用)
- 保证响应质量(太多并发会降低每个请求的质量)
代码位置:src/agents/pi-embedded-runner/run/lanes.ts
export function resolveGlobalLane(lane?: string): string {
return lane || `global:default`;
}
// 全局并发控制
const enqueueGlobal = (task) => enqueueCommandInLane(globalLane, task, {
concurrency: 4 // 最多4个并发
});
2.4 并发控制流程
三、队列模式
3.1 队列模式类型
代码位置:src/auto-reply/reply/queue/types.ts
export type QueueMode =
| "steer" // 转向模式:立即注入到当前 agent 回合中
| "followup" // 跟进模式:当前运行结束后,为下一个 agent 回合排队
| "collect" // 收集模式(默认):将所有排队的消息合并成单个后续回复
| "steer-backlog" // 转向+积压模式:现在转向当前回合,然后保留消息用于后续回合
| "interrupt" // 中断模式:中断当前运行
| "queue"; // 队列模式:普通队列
export type QueueSettings = {
mode: QueueMode;
debounceMs?: number; // 防抖延迟(毫秒),默认 1000ms
cap?: number; // 队列容量上限,默认 20
dropPolicy?: "old" | "new" | "summarize"; // 默认 summarize
};
3.2 Collect 模式(收集模式)
行为:将所有排队的消息合并成单个后续回复
适用场景:用户连续发送多条消息
示例:
{
"id": "e1c9d464",
"message": {
"content": [
{
"text": "[Queued messages while agent was busy]\n\n---\nQueued #1\n[Slack x +1s 2026-02-09 16:58 GMT+8] 算了 [slack message id: x channel: x]\n[message_id: x]\n\n---\nQueued #2\n[Slack x +4s 2026-02-09 16:58 GMT+8] 查一下天津的 [slack message id: x channel: x]\n[message_id: x]",
"type": "text"
}
],
"role": "user",
"timestamp": 1770627650797
},
"parentId": "7588527b",
"timestamp": "2026-02-09T09:00:50.802Z",
"type": "message"
}
代码位置:src/auto-reply/reply/queue/drain.ts
if (queue.mode === "collect") {
// 批量收集模式:合并所有等待消息为一个 Prompt
const items = queue.items.splice(0, queue.items.length);
const summary = buildQueueSummaryPrompt({ state: queue, noun: "message" });
const prompt = buildCollectPrompt({
title: "[Queued messages while agent was busy]",
items,
summary,
renderItem: (item, idx) => `---\nQueued #${idx + 1}\n${item.prompt}`.trim(),
});
await runFollowup({ prompt, run, enqueuedAt: Date.now() });
continue;
}
3.3 Steer 模式(转向模式)
行为:立即注入到当前 agent 回合中
适用场景:用户想纠正或补充正在处理的问题
示例:
用户:帮我写一段Python代码
[Agent正在生成代码...]
用户:等等,要用Java写
Agent立即收到补充消息,停止Python代码生成,改为生成Java代码
代码位置:src/agents/pi-embedded-runner/run/lanes.ts
if (shouldSteer && isStreaming) {
const handle = ACTIVE_EMBEDDED_RUNS.get(sessionId);
void handle.queueMessage(text); // 立即插入
}
// 实际:使用 pi-agent 框架的 steer 插入消息
const queueHandle: EmbeddedPiQueueHandle = {
queueMessage: async (text: string) => {
await activeSession.steer(text);
}
...
};
3.4 Followup 模式(跟进模式)
行为:当前运行结束后,为下一个 Agent 回合排队
适用场景:正常的新消息,等待当前处理完成
示例:
用户:帮我查一下天气
[Agent正在查询...]
用户:顺便查一下明天
第二条消息排队,等第一条处理完再处理
代码位置:src/auto-reply/reply/queue/drain.ts
// Followup 模式:处理下一个消息
const next = queue.items.shift();
await runFollowup(next);
3.5 Steer-Backlog 模式(转向+积压模式)
行为:现在转向当前回合,同时保留消息用于后续回合
适用场景:复杂的交互场景
3.6 Interrupt 模式(中断模式)
行为:中断当前运行
适用场景:紧急情况,需要立即停止当前任务
四、防抖与合并机制
4.1 防抖(Debounce)
作用:短时间内连续消息不立即触发处理
实现:
export async function waitForQueueDebounce(
queue: Queue,
): Promise<void> {
if (queue.debounceMs) {
await sleep(queue.debounceMs);
}
}
默认防抖时间:1000ms
4.2 合并(Merge)
作用:将多条消息合并为一条
实现:
export function buildCollectPrompt(options: {
title: string;
items: QueueItem[];
summary?: string;
renderItem: (item: QueueItem, idx: number) => string;
}): string {
const parts = [options.title];
if (options.summary) {
parts.push(options.summary);
}
parts.push(
...options.items.map((item, idx) =>
options.renderItem(item, idx)
)
);
return parts.join('\n\n');
}
4.3 队列容量控制
配置:
export type QueueSettings = {
mode: QueueMode;
debounceMs?: number; // 防抖延迟(毫秒),默认 1000ms
cap?: number; // 队列容量上限,默认 20
dropPolicy?: "old" | "new" | "summarize"; // 默认 summarize
};
丢弃策略:
| 策略 | 说明 |
|---|---|
old | 丢弃最旧的消息 |
new | 丢弃最新的消息 |
summarize | 将消息合并为摘要 |
五、队列调度策略
5.1 调度流程
5.2 优先级调度
export interface QueueItem {
id: string;
prompt: string;
priority: number; // 优先级,数字越大优先级越高
enqueuedAt: number;
run: FollowupRun;
}
export function sortQueueByPriority(items: QueueItem[]): QueueItem[] {
return items.sort((a, b) => {
// 1. 先按优先级排序
if (a.priority !== b.priority) {
return b.priority - a.priority;
}
// 2. 优先级相同,按时间排序
return a.enqueuedAt - b.enqueuedAt;
});
}
5.3 资源分配
export function allocateResources(
queue: Queue,
availableResources: Resources,
): Allocation {
// 1. 计算需要的资源
const required = estimateRequiredResources(queue);
// 2. 检查是否有足够资源
if (required.memory > availableResources.memory) {
throw new Error('Insufficient memory');
}
if (required.cpu > availableResources.cpu) {
throw new Error('Insufficient CPU');
}
// 3. 分配资源
return {
memory: required.memory,
cpu: required.cpu,
timeout: required.timeout,
};
}
六、核心代码文件索引
| 文件路径 | 功能 | 重要性 |
|---|---|---|
| src/auto-reply/reply/queue/types.ts | 队列类型定义 | ⭐⭐⭐⭐⭐ |
| src/auto-reply/reply/queue/drain.ts | 队列排空 | ⭐⭐⭐⭐⭐ |
| src/agents/pi-embedded-runner/run/lanes.ts | 并发控制 | ⭐⭐⭐⭐⭐ |
| src/auto-reply/reply/queue/manager.ts | 队列管理器 | ⭐⭐⭐⭐ |
| src/auto-reply/reply/queue/scheduler.ts | 队列调度器 | ⭐⭐⭐⭐ |
七、下一步
恭喜你完成了并发控制系统的学习!接下来建议:
- 📖 阅读 10. 安全架构 - 了解安全机制
- 🔄 查看 03. 消息流转 - 理解完整消息流程
- 🛠️ 探索 message_flow/ - 查看详细步骤文档
通过理解并发控制系统的工作原理,你已经掌握了 OpenClaw 的消息调度机制!
更多推荐



所有评论(0)