跳到主要内容

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 maplimit 让你能给"重活"(如 联网调研)单独设一个更紧的闸,而不必放松整个 workflow 的预算。

预算联动

并发结构同样受 limits 约束:

  • agent_runs:所有分支/map 项的调用共同计数,超限抛 LimitExceededError
  • duration:墙钟上限,超限 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 内嵌不进去的确定性检查。

下一步

  • verify —— 并发结果如何被判定
  • 示例:content-moderation · translation-flow · competitive-analysis (见 示例库)