
Electric Agents map-reduce 模式并行分片处理与结果归约的实体设计实战【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electricmap-reduce 是 Electric Agentselectric-ax/agents-runtime中用于把输入切分、并行处理、统一归约的实体协调模式父实体将输入拆成多个 chunk为每个 chunk 独立 spawn 一个 worker 并行执行最后把全部结果收集归约。本文以 designing-entities skill 的 pattern 文档 为主体结合仓库内的官方文档、源码实现与测试用例完整讲解该模式的状态设计、handler 骨架、不变量、审查清单与反模式并给出可复制的完整代码。一、模式定义与 manager-worker、pipeline 的本质区别map-reduce 的核心定义是把输入拆成多个 chunk用相互独立的 worker 并行处理所有 chunk最后收集并归约结果。它与manager-worker的关键差异在于与 manager-worker 不同map-reduce 的 worker 数量随输入变化——每个 chunk 对应一个 worker而不是一组固定命名的角色。换句话说manager-worker 是固定专家集map-reduce 是按数据量动态伸缩的工作池。相关模式的横向对照可参考 pattern-triggers.md 中的触发词表和消歧流程以及同为 patterns 参考的 manager-worker.md 与 pipeline.md。何时适用该模式输入是一个可并行处理的集合数据行、文档、URL、搜索查询等每个 chunk 用相同的方式处理同一个systemPrompt不同的 payloadworker 数量取决于输入大小而非固定值最后有一个reduce 阶段把所有结果聚合起来。何时不该用如果 worker 是固定命名集合→ 应改用manager-worker如果 worker串行执行A 的输出是 B 的输入→ 应改用pipeline。从消歧流程看pattern-triggers.md描述中出现 parallel、all at once、chunks、divide and process、batch、fan out and collect、process in parallel、map over 等短语时优先匹配 map-reduce若并行但角色固定则归入 manager-worker若串行则归入 pipeline。二、必备状态Required state三个集合撑起整个流程map-reduce 实体需要声明三个 state collection缺一不可state: { children: { schema: z.object({ key: z.string(), // chunk-${i}-${timestamp}-${counter} url: z.string(), chunk: z.number(), // index in the chunks array }), primaryKey: key, }, status: { schema: z.object({ key: z.literal(current), value: z.enum([idle, mapping, reducing]), }), primaryKey: key, }, spawnCounter: { schema: z.object({ key: z.literal(value), count: z.number() }), primaryKey: key, }, }三个集合各自职责集合用途关键点children记录每个已 spawn 的 workerkey唯一 ID、url子实体的 entityUrl、chunk在 chunks 数组中的下标reduce 阶段靠它把结果与输入 chunk 对应起来status阶段状态机取值idle → mapping → reducing → idle防止多个 map 阶段并发交错spawnCounter单调递增的 spawn 计数防止同一 map 操作被多次唤醒时子 ID 重复其中 spawnCounter 尤其关键Date.now()单独使用在工具被快速连续调用时会碰撞spawnCounter 保证即使同一次 map 操作被重复唤醒re-wake生成的子 ID 依然全局唯一。这一点在 官方文档版 map-reduce 中被凝练为一张集合表childrenkey、URL、chunk index而 skill 版本则进一步补全了status与spawnCounter。三、Handler 骨架从 firstWake 初始化到 map_chunks 工具下面是以 patterns/map-reduce.md 为准的完整 handler 骨架并补充了内置 worker 的契约说明import { entity } from electric-ax/agents-runtime async handler(ctx, wake) { if (ctx.firstWake) { ctx.db.actions.status_insert({ row: { key: current, value: idle } }) ctx.db.actions.spawnCounter_insert({ row: { key: value, count: 0 } }) } const mapChunksTool: AgentTool { name: map_chunks, parameters: Type.Object({ chunks: Type.Array(Type.String()), task: Type.String(), }), execute: async (_id, { chunks, task }) { transition(ctx.db, MR_TRANSITIONS, mapping) const counterRow ctx.db.collections.spawnCounter?.get(value)! const startCounter counterRow.count const spawnedIds: string[] [] for (let i 0; i chunks.length; i) { const spawnNum startCounter i 1 const id chunk-${i}-${Date.now()}-${spawnNum} const child await ctx.spawn( worker, id, { systemPrompt: task, tools: chunkTools, // required for built-in worker, e.g. [web_search, fetch_url] }, { initialMessage: chunks[i], wake: { on: runFinished, includeResponse: true } } ) ctx.db.actions.children_insert({ row: { key: id, url: child.entityUrl, chunk: i } }) spawnedIds.push(id) } ctx.db.actions.spawnCounter_update({ key: value, row: { count: startCounter chunks.length } }) transition(ctx.db, MR_TRANSITIONS, reducing) const results await Promise.all( spawnedIds.map(async (id) { const row ctx.db.collections.children?.get(id)! const handle await ctx.observe(entity(row.url)) return { id, chunk: row.chunk, text: row.output ?? } }) ) transition(ctx.db, MR_TRANSITIONS, idle) return { content: [{ type: text, text: results.map(r ## Chunk ${r.chunk}\n${r.text}).join(\n\n) }], details: {}, } }, } ctx.useAgent({ /* ... */, tools: [...ctx.electricTools, mapChunksTool] }) await ctx.agent.run() }三个阶段的职责拆解初始化firstWake写入status idle和spawnCounter 0。ctx.db.actions.*是集合的写接口insert/update/deletectx.db.collections.*是读接口get/toArray这组代理 API 的完整说明见 designing-entities skill 的 Canonical material。Map 阶段spawn 循环把状态推进到mapping从spawnCounter读出起始计数然后在一个循环里为每个 chunk 依次 spawn 一个worker子实体子 ID 形如chunk-${i}-${Date.now()}-${spawnNum}其中spawnNum来自单调递增计数器spawn 参数为{ systemPrompt: task, tools: chunkTools }——内置 worker 强制要求这两个字段wake: { on: runFinished, includeResponse: true }保证每个子实体跑完后父实体被重新唤醒并收到其输出每个 spawn 的结果立即写入children集合循环结束后一次性更新 spawnCounter 并推进到reducing。Reduce 阶段Promise.all 归约遍历spawnedIds对每个子实体调用ctx.observe(entity(row.url))获取句柄从children行中读出chunk下标与output最后用Promise.all并行收集全部结果并拼接成 markdown 文本返回给 LLM。注意骨架中transition(ctx.db, MR_TRANSITIONS, ...)是对status集合状态机迁移的封装MR_TRANSITIONS 即idle → mapping → reducing → idle实际项目中可以按相同语义自行实现或直接用ctx.db.actions.status_update(...)逐阶段更新。内置 worker 的 least-privilege 契约文档明确指出真实应用 spawn 的worker拿到的是服务端内置的最小权限 worker它不会收到ctx.electricTools。其完整契约在 manager-worker.md 中有权威定义interface WorkerArgs { systemPrompt: string tools: ArrayWorkerToolName // non-empty subset of: bash | read | write | edit | web_search | fetch_url | spawn_worker }两条硬性约束systemPrompt与tools都必须传二者缺一不可tools必须是WorkerToolName 的非空子集——省略tools或传空数组会在 spawn 时直接抛错。该约束在代码中也有对应体现内置 worker 的 spawn 失败会在 setup-context.ts 处报出child type:id already exists之类的错误因此 spawn 前必须做状态查证而runFinished/includeResponse的合法 schema 定义在 entity-schema.ts。四、官方文档版注册代码一个可直接落地的实体形态网站官方文档 给出了该模式的实体注册形态展示了完整实体在 registry 中的样子export function registerMapReduce(registry: EntityRegistry) { registry.define(map-reduce, { description: Map-reduce orchestrator that splits input into chunks, processes them in parallel with worker agents, then synthesizes results, state: { children: { primaryKey: key }, }, async handler(ctx) { ctx.useAgent({ systemPrompt: MAP_REDUCE_SYSTEM_PROMPT, model: claude-sonnet-4-6, tools: [...ctx.electricTools, createMapChunksTool(ctx)], }) await ctx.agent.run() }, }) }其工作原理描述与 skill 文档完全一致agent 暴露一个map_chunks工具调用时1map 阶段为每个 chunk 同时 spawn 一个 worker全部并行运行2工具立即返回实体在每个 worker 完成时被重新唤醒3所有 worker 通过 wake 事件汇报完毕后LLM 对结果进行综合synthesize。同时它提供了一个更精简的 map 阶段参考实现强调两个要点spawn 循环内不等待子任务立即返回以及 spawn 后立即把子实体信息写入children// Map phase - spawn all workers in parallel for (let i 0; i chunks.length; i) { spawnCounter const id chunk-${i}-${Date.now()}-${spawnCounter} const child await ctx.spawn( worker, id, { systemPrompt: task, tools: [read] }, { initialMessage: chunks[i], wake: { on: runFinished, includeResponse: true }, } ) ctx.db.actions.children_insert({ row: { key: id, url: child.entityUrl, chunk: i }, }) } return { content: [ { type: text as const, text: Spawned ${chunks.length} parallel workers. You will be woken as each finishes with its output in finished_child.response., }, ], details: { chunkCount: chunks.length }, }注意这里返回给 LLM 的文本明确提示了后续 wake 的行为finished_child.response会携带每个子任务完成后的输出父实体的 LLM 依靠这些 wake 事件增量掌握进度并最终完成综合。五、不变量Invariants五条必须守住的设计铁律Spawn counter 保证跨调用 ID 唯一。Date.now()单独使用在工具被快速连续调用时会碰撞必须叠加单调计数器。单一 spawn 循环循环内不对子任务 await。所有 worker 先全部启动reduce 阶段再统一收集——循环内 await 等于把并行打回串行。每个产生结果的 spawn 都要设置wake: { on: runFinished, includeResponse: true }。这是父实体在子任务完成时收到续接唤醒的前提。reduce 使用Promise.all。所有收集并行进行与并行的 map 阶段相匹配。状态迁移强制阶段顺序。idle → mapping → reducing → idle防止 map 阶段并发交错。六、模式专属审查清单Review Checklist MR1–MR7#规则原因MR1state.children记录每个 worker 的 chunk 下标与 urlreduce 阶段需要把结果与输入 chunk 对应MR2state.status迁移idle → mapping → reducing → idle防止并发 map 阶段交错MR3子 ID 中包含 spawn counter或等价单调源仅用Date.now()在快速调用下会碰撞MR4spawn 循环只记录子实体元数据并返回reduce 只在这些 worker 全部汇报后的后续runFinished唤醒中执行顺序 await 会破坏并行MR5reduce 对所有子实体使用Promise.all并行收集与并行 map 匹配MR6每个产生结果的 spawn 都设置wake: { on: runFinished, includeResponse: true }父实体在子任务完成时收到续接唤醒MR7每次 spawn 内置worker都传入非空tools数组否则内置 worker 在解析阶段直接抛错这组清单与 designing-entities skill 第 4 阶段Review的工作方式一致加载通用清单与模式清单后逐条机械比对设计输出✓ / ✗ / N/A报告再循环修订直至开发者明确批准。七、反模式Anti-patterns四种常见翻车写法只用随机 IDMath.random()。快速调用下碰撞频繁、难以调试。应该用 计数器 时间戳。在 map 循环里逐个 await 每个 worker。这会把本应并行的流程串行化——那就成了pipeline不是 map-reduce。固定 chunk 数量。如果 chunk 永远是 3 个固定角色请改用manager-worker。没有 spawn counter。重新唤醒时计数器被重置会重新生成相同 ID → 碰撞。八、源码与测试佐证runtime 如何支撑这个模式wake 事件的底层实现runFinished/includeResponse并非魔法而是 runtime 的一等公民entity-schema.ts 中Wake类型明确包含{ on: runFinished; includeResponse?: boolean }分支process-wake.ts 负责把wake参数归一化为runFinished字符串并透传includeResponseprocess-wake.ts 在 spawn 子实体时把wake: { on: runFinished, includeResponse: true }转译为子任务的完成唤醒配置。测试用例G1 map-reduce orderingruntime-dsl.test.ts 中有专门的G: map-reduce ordering测试组验证结果按 chunk 顺序返回这一核心不变量describe(G: map-reduce ordering, () { it(G1: map-reduce returns results in chunk order, async () { const parent await t.spawn(TYPES.g1MapReduce, map-1) await parent.send(map_chunks summarize :: alpha0|beta0|gamma0) // 期望输出按 chunk-1 / chunk-2 / chunk-3 顺序排列 // chunk-1:chunk-1::summarize:alpha | chunk-2:...:beta | chunk-3:...:gamma }) })该测试通过向父实体发送map_chunks summarize :: alpha0|beta0|gamma0验证 3 个 chunk 并行处理后reduce 结果严格保持 chunk 下标顺序。这正是children集合记录chunk下标的实战价值。实际项目中的应用形态在 agents-playground 示例 中researcher.ts 是一个贴近真实业务的研究型实体它的 system prompt 要求先判断请求是否需要拆解需要时用工具创建多个并行 specialist 子 agent 分别研究不同子问题并在每个 specialist 完成时被自动唤醒finished_child.response全部汇报后综合成带引用的最终回答。虽然它按稳定 ID 复用 specialist更接近 manager-worker 的变形但其拆解 → 并行 → 唤醒收集 → 综合的执行路径正是 map-reduce 思想的落地范本可作为理解本模式的业务化参考。九、实操建议与自检要点编写一个符合规范的 map-reduce 实体建议按以下顺序自检状态先行先声明children、status、spawnCounter三个集合不要遗漏ID 确定性spawn 循环中用chunk-${i}-${Date.now()}-${spawnNum}形式拼 ID并同步更新 spawnCounter不 await 子任务map 循环只 spawn 记录元数据收集一律放到 reducewake 配置所有产生结果的 spawn 都带上wake: { on: runFinished, includeResponse: true }worker 参数内置 worker 务必传非空tools且内容必须是WorkerToolName合法子集对照 MR1–MR7逐条过一遍清单尤其是 MR3计数器防碰撞与 MR4串行 await 陷阱。完成实体文件后按 designing-entities skill 的流程把registerXxx(registry)挂载进你的 registry 组合文件通常是entities/registry.ts或server.ts即可接入应用。【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考