并发控制系统

📚 学习路径:本文档是架构文档的第 9 部分。建议按顺序阅读:

前言

并发控制系统是 OpenClaw 的核心组件之一,负责管理消息的并发处理和队列调度。本文档将深入解析并发控制系统的设计、队列模式和调度策略。

通过阅读本文档,你将能够:

  • 理解并发控制的作用和设计理念
  • 掌握队列模式和调度策略
  • 了解两层并发控制机制
  • 理解防抖和合并机制

一、并发控制概述

1.1 为什么需要并发控制?

解决方案

并发控制解决的问题

消息竞争

资源耗尽

状态混乱

成本控制

队列系统

并发限制

会话隔离

资源管理

问题说明解决方案
消息竞争多个消息同时处理导致状态混乱队列系统
资源耗尽并发过多导致资源耗尽并发限制
状态混乱并发修改导致状态不一致会话隔离
成本控制并发 API 调用导致成本过高资源管理

1.2 并发控制的目标

实现方式

并发控制的目标

用户体验

资源控制

状态一致性

成本优化

队列模式

并发限制

会话隔离

智能调度

目标说明实现方式
用户体验避免消息乱序,支持实时打断队列模式
资源控制防止系统过载,控制 API 调用并发限制
状态一致性避免并发修改会话状态会话隔离
成本优化减少 API 调用,优化资源使用智能调度

二、两层并发控制

2.1 架构设计

处理层

全局队列
并发度: 4

会话队列

消息输入

消息1

消息2

消息3

消息4

会话A队列

会话B队列

会话C队列

Worker 1

Worker 2

Worker 3

Worker 4

Agent实例1

Agent实例2

Agent实例3

Agent实例4

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 并发控制流程

Agent Worker 全局队列 会话队列 消息 Agent Worker 全局队列 会话队列 消息 加入会话队列 检查会话并发 请求全局 Worker 检查全局并发 分配 Worker 执行 Agent 返回结果 释放 Worker 完成通知 处理完成

三、队列模式

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 调度流程

steer

followup

collect

steer-backlog

interrupt

消息到达

会话队列是否空闲?

立即处理

加入队列

队列模式?

立即插入当前回合

排队等待

等待合并

立即插入+保留

中断当前

防抖等待

队列是否已满?

继续等待

执行丢弃策略

处理消息

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队列调度器⭐⭐⭐⭐

七、下一步

恭喜你完成了并发控制系统的学习!接下来建议:


通过理解并发控制系统的工作原理,你已经掌握了 OpenClaw 的消息调度机制!

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐