练习 20:并行扇出与上限

练习 19 的 sub_agent 一次只能处理一个子任务。如果模型判断一件事可以 拆成三个互不依赖的子任务,同一轮里连续发起三次 sub_agent 调用,现在 的代码是 for 循环里一个个 execute——而每一次 sub_agent 调用自己就是 一整场多轮的 LLM 对话,三个排队跑,总共要等的时间是三份时间加起来。这一章把 "这一批调用能不能并发跑"这个判断写成代码,用一个有容量上限的 channel 当信号量,防止判断成"能"之后,扇出多大都不设限,把本机资源或 provider 的并发配额一次性打满。

敲进去

在练习 19 的代码上继续写。先加一个能不能并发的判据:

// maxParallelSubAgents 限制一轮里最多同时跑几个 sub_agent。每一个都会
// 起自己的一整套 provider 连接和多轮对话,扇出多大都不设上限,就是在
// 拿本地资源和 provider 的并发限额去赌。这个数字没有理论最优解,纯粹是
// 部署环境的取舍——本书的玩具 harness 是单机跑,4 只是"够看出并发效果,
// 又不至于把本机或 provider 打满"的一个保守选择。
const maxParallelSubAgents = 4

// canFanOut 判断这一轮工具调用能不能并发跑。判据故意收得很紧:调用数
// 大于一个,且全部是 sub_agent——不是"只读工具都能并发"这种更通用的
// 规则。原因是这本书至今没有给任何工具标注过"只读",bash/write_file/
// edit_file 都会改共享状态(cwd、trash、registry 的 hasRead 记账),
// 混在一起并发执行没人能担保顺序和结果;sub_agent 不一样——它起的是
// 一整个独立的子 agent,自己的 history、自己的 childReg,跟父 agent
// 的注册表没有任何写共享(sub_agent 的参数里没有 path 字段,注册表的
// hasRead 记账不会被它触碰)。这份安全保证只覆盖 history/registry 这层
// 状态——子 agent 内部如果自己调用了需要人工确认的 bash 命令,confirm()
// 读的是同一个共享 os.Stdin,多个 goroutine 同时问会互相冲撞,见这一章
// "常见问题"里的实测记录,这份判据没有、也不打算解决这个问题。
func canFanOut(calls []toolCall) bool {
	if len(calls) < 2 {
		return false
	}
	for _, tc := range calls {
		if tc.Function.Name != "sub_agent" {
			return false
		}
	}
	return true
}

再加实际分发这一批调用的函数——能并发就开 goroutine,不能就照原来的 串行路径:

// dispatchToolCalls 跑完一轮里的全部工具调用,按原始顺序整理成待追加
// 的 tool 消息。canFanOut 为真时用一个容量 maxParallelSubAgents 的
// channel 当信号量:每个 goroutine 先占一个坑位再执行,执行完释放,
// 坑位不够的调用在 channel 上排队——这就是"goroutine + channel 限流",
// 不需要另起一个调度器或线程池。结果按 index 写回一个和 calls 等长的
// 切片,不依赖 map 的遍历顺序,保证 tool 消息和原始 tool_calls 一一
// 对应。不能并发的这一轮(只有一个调用,或混了 bash/write 这类工具)
// 走原来那条串行路径,行为跟练习 19 完全一样。
func dispatchToolCalls(reg *registry, round int, calls []toolCall) []message {
	results := make([]string, len(calls))
	if canFanOut(calls) {
		var wg sync.WaitGroup
		sem := make(chan struct{}, maxParallelSubAgents)
		var logMu sync.Mutex
		for i, tc := range calls {
			wg.Add(1)
			sem <- struct{}{} // 占坑位;坑位不够就阻塞在这一行排队
			go func(i int, tc toolCall) {
				defer wg.Done()
				defer func() { <-sem }() // 让出坑位给下一个排队的调用
				logMu.Lock()
				fmt.Fprintf(os.Stderr, "[round %d] %s(%s)\n", round, tc.Function.Name, tc.Function.Arguments)
				logMu.Unlock()
				results[i] = reg.execute(tc.Function.Name, tc.Function.Arguments)
			}(i, tc)
		}
		wg.Wait()
	} else {
		for i, tc := range calls {
			fmt.Fprintf(os.Stderr, "[round %d] %s(%s)\n", round, tc.Function.Name, tc.Function.Arguments)
			results[i] = reg.execute(tc.Function.Name, tc.Function.Arguments)
		}
	}
	out := make([]message, len(calls))
	for i, tc := range calls {
		out[i] = message{Role: "tool", ToolCallID: tc.ID, Content: results[i]}
	}
	return out
}

main() 里原来那段"一个个 execute、一个个 append"的 for 循环,换成一行:

if canFanOut(msg.ToolCalls) {
	fmt.Fprintf(os.Stderr, "[round %d 并发扇出:%d 个 sub_agent,上限 %d 个坑位]\n",
		round, len(msg.ToolCalls), maxParallelSubAgents)
}
sess.History = append(sess.History, dispatchToolCalls(reg, round, msg.ToolCalls)...)

别忘了在 import 里加一行 "sync"

跑起来

go build -o ex20 .

造三个互相独立的项目目录,一份提到废弃接口,两份没提:

mkdir -p projA projB projC
cat > projA/README.md << 'EOF'
# Project A
这是一个稳定维护的工具库,所有 API 都在积极使用中,没有计划废弃任何接口。
EOF
cat > projB/README.md << 'EOF'
# Project B
注意:`legacy_parse()` 函数已经 deprecated,请迁移到 `parse_v2()`,
下个大版本会彻底移除旧接口。
EOF
cat > projC/README.md << 'EOF'
# Project C
这个项目还在早期阶段,欢迎贡献,暂无废弃计划。
EOF

实验一:三个独立检查,扇出跑一次。

./ex20 "我这里有三个互相独立的项目目录:projA、projB、projC,各自有一份 README.md。请分别为每一个目录派一个独立的 sub agent 去检查它的 README.md 里有没有提到 'deprecated'(废弃)相关的说明,如果有就摘要是哪个接口废弃了。这三个检查任务彼此没有依赖,请一次性同时发起这三个 sub_agent 调用,不要等上一个做完再发起下一个。"

同一个任务,用上一章的 ex19 跑一遍作对照(同一份代码,没有并发分发):

time ./ex19 "……同一段任务原文……"
time ./ex20 "……同一段任务原文……"

实验二:扇出批里混进一个需要人工确认的 bash 命令。 把任务原文换成 明确要求子 agent 用 bashgrep(而不是 read_file)去查:

./ex20 "我这里有三个互相独立的项目目录:projA、projB、projC,各自有一份 README.md。请分别为每一个目录派一个独立的 sub agent,要求每个 sub agent 都必须用 bash 的 grep 命令(例如 grep -i deprecated projX/README.md)去检查有没有提到 'deprecated',不要用 read_file。这三个任务彼此没有依赖,请一次性同时发起这三个 sub_agent 调用。"

你应该看到什么

实验一,DeepSeek 一次性发起三个调用:

[round 1 并发扇出:3 个 sub_agent,上限 4 个坑位]
[round 1] sub_agent({"description": "检查 projC README 中 deprecated 说明", ...})
[子 agent "检查 projC README 中 deprecated 说明" 开始,独立的一份 history,父对话它一个字都看不到]
[round 1] sub_agent({"description": "检查 projA README 中 deprecated 说明", ...})
[子 agent "检查 projA README 中 deprecated 说明" 开始,独立的一份 history,父对话它一个字都看不到]
[round 1] sub_agent({"description": "检查 projB README 中 deprecated 说明", ...})
[子 agent "检查 projB README 中 deprecated 说明" 开始,独立的一份 history,父对话它一个字都看不到]
[子 agent "检查 projB README 中 deprecated 说明" 结束:内部消耗约 3441 tokens,……]
[子 agent "检查 projC README 中 deprecated 说明" 结束:内部消耗约 3464 tokens,……]
[子 agent "检查 projA README 中 deprecated 说明" 结束:内部消耗约 3743 tokens,……]

三个 sub agent 已并行完成,结果汇总如下:
| projA | 没有计划废弃任何接口 |
| projB | legacy_parse() 已废弃,迁移到 parse_v2() |
| projC | 暂无废弃计划 |

三条"开始"打印挨在一起,跟三条"结束"完全不按发起顺序回来——projB 最先结束、projA 最后,跟发起时的 C→A→B 顺序对不上,这正是并发跑的 证据:谁先谁后由各自那次 API 请求的实际耗时决定,不是代码里的调用顺序。

对照真实耗时:ex19(串行)跑这个任务约 19.5 秒,ex20(并发扇出) 约 16.2 秒。DeepSeek 这一侧有实打实的加速,但没有到"三倍变一倍"那么 夸张——三个子 agent 各自只需要一轮"读文件 + 回答",单个子 agent 本身 的耗时已经不长,扇出省下的是"排队等前一个"的那部分时间,不是把三份 计算量压缩成一份。

实验二,混进 bash grep 之后:

⚠️  模型想执行: grep -i deprecated projC/README.md
允许吗?(y/N) 
⚠️  模型想执行: grep -i deprecated projB/README.md; echo "exit_code=$?"
允许吗?(y/N) 
⚠️  模型想执行: grep -i deprecated projA/README.md; echo "exit code: $?"
允许吗?(y/N) 
⚠️  模型想执行: grep -i deprecated projB/README.md
允许吗?(y/N) 
……(还有更多,三个 sub agent 各自重试了两三次)

三个 sub_agent 都已返回,但结果一致:它们都没能完成检查。
……命令被系统以"权限拒绝——用户没有批准这条命令"拦下……

三个 sub_agentexecute 走的是和练习 19 一模一样的 bashTool, 一撞上 classifyBash 归类为 ask 档的命令,就调用 confirm()——而这次 是三个 goroutine 同时调用它,三条"允许吗?(y/N)"提示交错打印在一起, 没有任何东西告诉你哪一条对应哪个子 agent 的哪条命令。

本机 Ollama,实验一同样正确扇出、正确识别出 projB,但耗时反过来: ex19(串行)约 24.9 秒,ex20(并发扇出)约 38.0 秒——并发比串行 更慢。

发生了什么

"能不能并发"是一个需要显式写出来的判断,不是"起了 goroutine 就自动 安全"。 canFanOut 故意把判据收得很窄:不是"这批调用里没有明显冲突 就并发",而是"必须全部是 sub_agent,多一个都不行"。原因写在 练习 19 就交代过:sub_agent 起的是一整个独立的子 agent,有自己的 childReg、自己的 childHistory,跟父 agent 的注册表之间除了共享 只读的工具定义,没有任何写共享——hasRead 这本账,sub_agent 的参数 里压根没有 path 字段,永远碰不到。而 bash/write_file/edit_file 都会改这个进程共享的状态(工作目录、trash/ 备份、hasRead 记账), 这本书至今也没有给任何工具标注过"只读",没有这层标注就没法证明"这几个 调用放在一起并发跑是安全的",判据只能收紧到唯一一种已知安全的情况。

信号量不是防止"跑错",是防止"跑爆"。 三个子任务同时发起,代码逻辑 上完全没问题——真正的风险是模型某一轮里发起了三十个 sub_agent 调用, 每一个都要开一整套 HTTP 连接和多轮对话,一次性全冲上去,本地资源和 provider 的并发配额都扛不住。maxParallelSubAgents 容量的 channel 就是 拿来挡这件事的:坑位占满,多出来的调用在 sem <- struct{}{} 这一行 排队,等前面的 goroutine 跑完 <-sem 让出坑位,不需要另写一个任务 队列或线程池,sync.WaitGroup + 一个 buffered channel 就够。

实验一的加速幅度不算夸张,这本身就是一个诚实的数据点。 并发扇出 省下来的是"排队等前一个子 agent 跑完"这段时间,不是把三份计算压缩成 一份——每个子 agent 内部该读的文件、该发的请求,一个都不会少。子任务 本身越轻(这一章的例子只需要一轮"读文件+回答"),扇出能省下的绝对时间 就越有限;子任务越重(多轮、多次工具调用),并发的收益才会越明显。

本机 Ollama 那次"并发反而更慢",不是这一章代码写错了,是硬件层面的 真相。 DeepSeek 是云端服务,背后有能同时处理多个请求的算力;本机 Ollama 只有一份模型实例,物理上同一时刻只能做一次推理——三个 goroutine 同时把请求发过去,Ollama 内部还是得排队一个个算,额外多出来的只有 并发调度、连接管理这些开销,没有任何计算真的被并行掉,账面上自然是 负数。"并发扇出"这件事在代码层面永远成立,但它能不能换来真实的加速, 取决于运行这些请求的后端本身撑不撑得住并发——这是部署环境的事,不是 canFanOut 这几行代码能决定的。

实验二暴露了这一章故意没解决的一个洞。 confirm()os.Stdin 读一行,练习 9 写它的时候,从来没设想过会有第二个 goroutine 同时调用它——现在三个 sub agent 各自的 bashTool 都可能在同一时刻 撞上 ask 档,三次 confirm() 并发执行,三条提示打印交错在一起, 终端上没有任何标记告诉你哪一条对应哪个子任务。canFanOut 的安全论证 只覆盖 history/registry 这层状态,从没延伸到"人工确认"这个共享的 终端 I/O 通道——这本书至今的人工确认机制,就是为单线程场景设计的, 搬到并发场景下直接失效,这一章没有修它,只是让这个洞第一次露出来。

常见问题

  • 为什么不干脆把所有工具都标成能并发:因为这本书从没做过"标注只读 工具"这件事——没有这层信息,就没法证明一批混杂的调用放在一起跑是 安全的。canFanOut 收紧到只认 sub_agent,是"能证明安全的最小范围", 不是"理论上能做到的最大范围"。
  • maxParallelSubAgents 为什么是 4,不是更大或者干脆不设上限: 这个数字没有标准答案,纯粹是部署环境的取舍——本机资源、provider 的并发限额、真实任务的扇出规模,都会改变这个数字该多大。不设上限 等于把这个决定丢给"模型这一轮碰巧发起了几个调用",那不是一个决定, 是听天由命。
  • 并发扇出会不会让父 agent 收到的 tool 消息顺序乱掉:不会。 dispatchToolCallsresults[i] 按下标写回,每个 goroutine 只碰 自己那个下标,wg.Wait() 之后再按原始顺序拼出 []message——谁先 执行完不影响最终顺序,乱的只是"结束"日志打印的先后,不是最终追加进 sess.History 的顺序。
  • 如果子 agent 内部触发了需要人工确认的 bash 命令会怎样:会触发 实验二那种交错打印——三个 goroutine 同时调用同一个 confirm(), 提示互相冲撞,没有任何标记区分谁是谁。这一章没有解决这个问题, 留在这里当一个明确的坑:并发扇出目前只对"全自动、不需要人工确认" 的子任务安全。

加分练习

  1. maxParallelSubAgents 改成 1,重新跑实验一,确认耗时退化成跟 ex19 差不多——这是在验证"信号量容量决定了并发度"这句话,不是 靠肉眼猜的。
  2. confirm() 加一把互斥锁(同一时刻只允许一个 goroutine 进入这个 函数),重新跑实验二,看提示是不是不再交错——注意这只解决"打印 交错",没解决"三个子任务各自在等谁批准"这个更深的问题,想清楚 为什么,会发现单靠一把锁只是让问题从"看不清"变成"看得清但仍然 串行"。
  3. 让父 agent 一次发起 8 个 sub_agent 调用(比如让它检查 8 个不同的 目录),观察 [round %d 并发扇出] 那行日志和实际的"开始"打印数量, 确认同一时刻正在跑的子 agent 数不会超过 maxParallelSubAgents
  4. time 分别测三个子任务从"轻"(读一个文件回答一句话,这一章的 例子)到"重"(让每个子 agent 自己再多跑几轮,比如先列目录再读文件 再总结)的并发收益变化,验证"发生了什么"里那句"子任务越重,扇出 收益越明显"是不是站得住。