parallel 与 parallel map
两种并发原语,都是运行时确定执行的结构:并发的是 agent 调用, 不是控制流。
parallel:分支并发
let parts = parallel {
condensed = agent(ops.main) { task "..." expect Condensed timeout 5m }
actions = agent(ops.main) { task "..." expect ActionItem[] timeout 5m }
}
// parts.condensed : Condensed
// parts.actions : ActionItem[]
- 各分支并发求值,整体形成 barrier:任一分支失败,整个 parallel 失败
- 结果类型:以分支名为字段的对象
{ 分支名: 分支类型 }——所以分支结果 可以直接按字段访问,也能整体赋给同构的具名类型(结构化可赋值) - 分支名唯一(
duplicate-branch);分支之间互相不可见
适用:同一阶段里的几个独立产出(摘要 + 行动项;风格 + 事实 + 合规), 它们逻辑上并列,谁也不依赖谁。
parallel map:批量并发
let verdicts = parallel map input.items as item limit 6 {
agent(moderator) {
task "Review one piece of content against policy"
input { item: item }
expect ModerationVerdict
timeout 90s
}
}
// verdicts : ModerationVerdict[]
- source 必须是数组(
parallel-map-source-not-array) - 循环变量(上例
item)只在 map 体内可见 limit N是本 map 的并发上限(正整数);缺省无界, 但仍受 workflow / host 并发限制约束- 结果顺序与输入顺序一致,与完成顺序无关——下游永远拿稳定序
- 类型:
T[] → U[](U 为 body 类型,通常就是 expect 类型)
适用:逐条处理批量输入(邮件、简历、发票、翻译分段、竞品列表)。
map 体内是"一个表达式"
body 是单个表达式(通常是 agent(...)),不是语句块——不能在里面写
if/require。这不是缺陷,是设计:
- 分支逻辑应该写进 agent 的 task 与 expect(让模型按输入行事),
或者在上游用
if分流成两个 map - 结构化后处理(唯一性、计数、范围)放在 map 之后的 stage 里做
stage classify -> ClassifiedEmail[] {
let items = parallel map input.emails as item limit 4 {
agent(assistant) { task "..." expect ClassifiedEmail }
}
require unique(items[*].id) else fail "ids must stay unique" // map 之后
return items
}
并发的最终上限
三层取最小:
最终并发 = min(host.maxConcurrency, workflow.limits.concurrency, map 的 limit)
宿主侧用 semaphore 强制。parallel map 的 limit 让你能给"重活"(如
联网调研)单独设一个更紧的闸,而不必放松整个 workflow 的预算。
预算联动
并发结构同样受 limits 约束:
agent_runs:所有分支/map 项的调用共同计数,超限抛LimitExceededErrorduration:墙钟上限,超限WorkflowTimeoutError
给批量 map 估预算时按条目数上界算:[1..50] 的输入 × 每条 1 次调用,
agent_runs 至少要给到 50(再加非 map 的调用)。
模式:分阶段 map
对"逐条判断 → 按判断逐条处理"的流水,用两个 stage 串起来:
stage audit -> ArticleVerdict[] { // 判断:并发读
parallel map input.articles as item limit 4 {
agent(curator) { tools [read_file, grep] expect ArticleVerdict }
}
}
stage rewrite after audit -> RewriteOutcome[] { // 处理:并发写,更小的并发
parallel map audit as item limit 2 {
agent(curator) {
tools [read_file, write_file]
write input.write_scope
expect RewriteOutcome
}
}
}
两个 map 之间可以用 require 做断言(如 unique(outcomes[*].id)),
这是 map 内嵌不进去的确定性检查。