Pipecat 多 Worker 并行辩论实战:用 job_group 实现三个 LLM Worker 的并行扇出与结果聚合
Pipecat 多 Worker 并行辩论实战用 job_group 实现三个 LLM Worker 的并行扇出与结果聚合【免费下载链接】pipecatOpen Source framework for voice agents, multimodal apps, and realtime AI. Maintained by Daily and the community.项目地址: https://gitcode.com/GitHub_Trending/pi/pipecat本文围绕 Pipecat 官方示例examples/multi-worker/parallel-debate展开。这是一个语音主持人 三个并行 LLM Worker的多 worker 应用用户说出一个辩题主持人 LLM 调用debate工具通过worker.job_group(...)把任务并行分发给 advocate正方、critic反方、analyst分析师三个LLMContextWorker等待三方全部回复后聚合成一段平衡的口头总结。读完后你将掌握 Pipecat worker/job 体系的核心用法如何用job_group做结构化并发、如何让每个 worker 跨轮次保持独立的LLMContext、以及如何通过 bus 消息把子 worker 的回复传回主 worker 完成合成。整体架构主持人主 Worker 加三个角色化子 Worker示例的 READMEexamples/multi-worker/parallel-debate/README.md给出的架构如下Main worker (transport LLM debate tool) └── job_group(advocate, critic, analyst) └── DebateWorker (LLMContextWorker, one per role)主 Worker承载 transportSTT、TTS、负责主持对话的 LLM以及一个debate直接函数direct function。这个函数内部调用worker.job_group(...)做并行扇出。辩论子 Worker三个LLMContextWorker在同一个 runner 上并行运行。每个 worker 持有自己独立的LLMContext因此在多轮辩论中都能记得之前讨论过哪些辩题每个 worker 完成自己的回合后通过 assistant-aggregator 的on_assistant_turn_stopped事件把完整回复作为 job 响应发回。这一形态在 Pipecat 多 worker 模式库中被归类为Parallel fan-out一个 worker 向多个对等 worker 并行派发同一个 job并等待全部响应适合多视角汇总成一个答复的场景。关于 setup 与环境变量的说明见顶层 examples/multi-worker/README.md其中对多 worker 模式的定义是一个 worker 是挂到共享 bus 上、通过消息生命周期事件、帧传输、job RPC协作的工作单元。运行方式与环境变量从仓库根目录进入示例目录后cd examples/multi-worker uv run parallel-debate/parallel-debate.py启动后在浏览器打开http://localhost:7860/client即可与 bot 对话。若要使用 Daily 云端 transportuv run parallel-debate/parallel-debate.py --transport daily首次使用前需按 examples/multi-worker/README.md 的 Setup 章节安装依赖并配置环境uv sync --all-extras source .venv/bin/activate cd examples/multi-worker cp env.example .env本示例依赖的环境变量模板见 examples/multi-worker/env.example变量用途本示例中对应OPENAI_API_KEYLLM主 Worker 与三个 DebateWorker 的OpenAILLMServiceDEEPGRAM_API_KEYSTTDeepgramSTTServiceCARTESIA_API_KEYTTSCartesiaTTSServiceDAILY_API_KEY可选仅--transport daily时需要Daily transport主 Workertransport、聚合器与 debate 直接函数核心实现在 examples/multi-worker/parallel-debate/parallel-debate.py。run_bot中主 Worker 的组装要点服务栈DeepgramSTTService做语音识别CartesiaTTSServicevoice 固定为 Jacqueline做语音合成OpenAILLMService作为主持人 LLM其 system instruction 明确要求当用户给出辩题时调用debate工具收集三个视角再把结果合成为简洁、适合口头表达的平衡总结。上下文与聚合器LLMContext(tools[debate])把debate函数注册为主 LLM 的工具LLMContextAggregatorPair同时提供用户侧聚合带SileroVADAnalyzer做端点检测与助手侧聚合。Pipeline 顺序transport.input() → stt → aggregators.user() → llm → tts → transport.output() → aggregators.assistant()。注意 assistant 聚合器放在 pipeline 末尾这样主 Worker 自己的助手回合结束时才会触发on_assistant_turn_stopped与子 Worker 使用同名事件的机制一致。Worker 化pipeline 被包进PipelineWorker名字parallel-debate开启enable_metrics与enable_usage_metricsProcessorUnusablePolicy.END随后和三个DebateWorker一起交给WorkerRunner.add_workers(...)在同一个 runner 下并行运行——这就是单进程内、共享 in-memory bus的本地多 worker 形态。debate是主 LLM 的 direct function签名和实现值得逐行看tool_options(cancel_on_interruptionFalse, timeout_secs60) async def debate(params: FunctionCallParams, topic: str): Analyze a topic from multiple perspectives (advocate, critic, analyst). logger.info(fStarting debate on {topic}) async with params.pipeline_worker.job_group( *ROLE_PROMPTS, payload{topic: topic}, timeout30 ) as tg: pass result \n\n.join(f{r[role].upper()}: {r[text]} for r in tg.responses.values()) logger.info(Debate complete, synthesizing) await params.result_callback(result)关键点params.pipeline_worker.job_group(...)FunctionCallParams上直接拿到当前PipelineWorker把 worker 名这里就是三个角色名 advocate / critic / analyst展开传入payload{topic: topic}把辩题带给所有子 workertimeout30秒内必须全部响应。*ROLE_PROMPTS利用字典按键展开的特性把三个角色名作为 worker 名单传给job_group。上下文退出后tg.responses是一个以 worker 名为 key 的响应字典代码把三条回复按ADVOCATE / CRITIC / ANALYST前缀拼成一个字符串再经params.result_callback(result)回注给主持人 LLM——合成synthesis这一步交给主 LLM 的下一轮自然语言生成完成而不是模板拼接。tool_options(cancel_on_interruptionFalse, timeout_secs60)工具允许用户打断时不取消辩论在后台继续跑但给整个工具调用设了 60 秒的硬上限与 job_group 内部 30 秒超时形成两层保护。DebateWorker每个角色独立 LLMContext 的跨轮记忆DebateWorker继承自pipecat.workers.llm.LLMContextWorker实现位于 src/pipecat/workers/llm/llm_context_worker.py是worker 即带 LLM 上下文的 LLM 消费者这一模式的落地。示例中的实现ROLE_PROMPTS { advocate: You argue IN FAVOR of the topic. ... Be concise, just 2-3 sentences., critic: You argue AGAINST the topic. ... Be concise, just 2-3 sentences., analyst: You provide a BALANCED, NEUTRAL analysis. ... Be concise, just 2-3 sentences., } class DebateWorker(LLMContextWorker): def __init__(self, role: str): llm OpenAILLMService( api_keyos.environ[OPENAI_API_KEY], settingsOpenAILLMService.Settings(system_instructionROLE_PROMPTS[role]), ) super().__init__(role, llmllm) self._role role self._current_job_id: str | None None self.assistant_aggregator.event_handler(on_assistant_turn_stopped) async def on_assistant_turn_stopped(aggregator, message: AssistantTurnStoppedMessage): text message.content if self._current_job_id: job_id self._current_job_id self._current_job_id None await self.send_job_response(job_id, {role: self._role, text: text}) async def on_job_request(self, message: BusJobRequestMessage) - None: await super().on_job_request(message) self._current_job_id message.job_id await self.queue_frame( LLMMessagesAppendFrame( messages[{role: developer, content: fTopic: {message.payload[topic]}}], run_llmTrue, ) )它的设计要点有三worker 名即角色名。super().__init__(role, llmllm)用角色名作为 worker 名所以主 worker 里job_group(*ROLE_PROMPTS)传的字典键恰好就是 bus 上可寻址的 worker 名。每个角色拿到独立的OpenAILLMService实例同一 key不同 system prompt。job 请求注入辩题并触发 LLM。on_job_request收到BusJobRequestMessage后记录message.job_id然后向自己队列推一个LLMMessagesAppendFrame追加一条developer角色的 Topic: ... 消息并置run_llmTrue立即跑 LLM。因为追加是发生在该 worker 自己的LLMContext上历史辩题会一直留在上下文里——这就是 README 所说each worker keeps its own LLM context across rounds的实现基础。回合结束即回包。LLMContextWorker内置了 assistant 侧聚合器用self.assistant_aggregator.event_handler(on_assistant_turn_stopped)挂监听当该 worker 的助手回合完整结束流式文本聚合完毕时拿到完整文本message.content调用send_job_response(job_id, {role: ..., text: ...})把结果作为 job 响应发回发起方。_current_job_id用后即清保证并发多次 job 时响应路由不串号。值得注意的一个细节这里把完整回合文本而非流式增量作为 job 响应避免了主 worker 拿到半截句子代价是合成回复要等最慢的那个角色说完timeout30就是为这种木桶效应兜底。job_group 机制请求、等待、响应收集与取消语义job_group的类型与上下文实现位于 src/pipecat/pipeline/job_context.py。从源码结构看一次job_group调用的生命周期是进入上下文JobGroupContext为组生成共享job_id向worker_names中每个 worker 发送BusJobRequestMessage携带payload并等待所有目标 worker 注册就绪ready等待期间tg支持async for迭代中间事件——JobGroupEvent有UPDATE、STREAM_START、STREAM_DATA、STREAM_END四种类型每条事件带worker_name因此可以实时展示各角色的进度。parallel-debate 中没有迭代async with ... : pass只关心最终结果退出上下文等待全组响应若正常完成tg.responses是按 worker 名组织的dict[str, dict]若超时、worker 报 error默认cancel_on_errorTrue会取消其余成员、或块内抛异常/被取消如工具被打断则抛出JobGroupError块内异常场景下还会向剩余 worker 发送BusJobCancelMessagereason 为context exited with error。组级参数由JobGroupParams控制继承自JobParamspayload、timeout同时覆盖等待 worker ready与 job 执行两个阶段、labelBaseUIWorker用它给客户端进度卡命名、cancellable是否允许外部如 UI 请求取消、cancel_on_error默认True。这些语义在 tests/test_job_group.py 中有系统性验证与本示例直接相关的用例包括test_job_group_collects_responses两个 worker 的响应按名字收集进tg.responsestest_job_group_raises_on_timeout/test_job_group_raises_on_ready_timeoutjob 超时抛JobGroupErrortimeoutworker 未就绪超时抛 not readytest_job_group_cancels_on_block_exception块内抛异常时剩余 worker 收到取消消息test_job_group_partial_responses_on_error一个 worker 报错时已成功的部分响应仍可访问test_job_group_iterates_updates/test_job_group_iterates_stream_eventsasync for可拿到进度与流式事件。也就是说parallel-debate 里那句看似平淡的timeout30背后对应的是等 ready 等执行的完整超时预算超时会以JobGroupError形式中断工具调用而不是无限挂起。会话生命周期打招呼与收尾示例还展示了两个 transport 事件钩子构成完整会话边界transport.event_handler(on_client_connected) async def on_client_connected(transport, client): context.add_message( {role: developer, content: Greet the user and tell them you can moderate a debate on any topic. ...} ) await worker.queue_frame(LLMRunFrame()) transport.event_handler(on_client_disconnected) async def on_client_disconnected(transport, client): await runner.cancel()客户端连上后向主 Worker 的LLMContext追加一条 developer 指令并推LLMRunFrame触发 LLM 主动开口完成先自我介绍、再问用户想辩什么的开场客户端断开时runner.cancel()停掉整个 runner主 worker 与三个 DebateWorker 一并退出。入口函数bot(runner_args)通过create_transport(runner_args, transport_params)按 CLI 参数选择 eval / daily / webrtc 三种 transport 参数这也是 README 中--transport daily选项的来源if __name__ __main__下调用pipecat.runner.run.main()启动。小结与延伸parallel-debate 示例用不到 250 行代码串起了 Pipecat 多 worker 体系的几个核心概念WorkerRunner统一管理多个 worker、job_group提供带超时与取消语义的并行扇出、LLMContextWorker让子 worker 拥有独立且持久的 LLM 上下文、bus 消息BusJobRequestMessage/ job response完成跨 worker 的 RPC。若要继续扩展可以在async with ... as tg:块内改用async for event in tg:迭代JobGroupEvent向用户实时播报各角色进度把 DebateWorker 换成LLMWorker见 src/pipecat/workers/llm/llm_worker.py做无上下文的纯函数式角色或换成自定义BaseWorker接入检索、代码执行等非 LLM 能力参考同目录的 handoff、code-assistant 等示例见 examples/multi-worker/README.md 的示例表观察 fan-out 之外的其他 worker 协作形态。【免费下载链接】pipecatOpen Source framework for voice agents, multimodal apps, and realtime AI. Maintained by Daily and the community.项目地址: https://gitcode.com/GitHub_Trending/pi/pipecat创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考