ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Pipecat 实战:用 BaseUIWorker 扇出异步任务组,把后台工作流式渲染到 Web 客户端

Pipecat 实战:用 BaseUIWorker 扇出异步任务组,把后台工作流式渲染到 Web 客户端 Pipecat 实战用 BaseUIWorker 扇出异步任务组把后台工作流式渲染到 Web 客户端【免费下载链接】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 的 async-tasks 示例讲解如何用BaseUIWorker作为客户端可见的任务组分发器由主语音管线的 LLM 调用research工具触发一组后台 worker 并行工作生命周期事件以ui-job-group信封流式推送到浏览器页面用户还能在任务执行途中点击 Cancel 取消整组任务。读完本文你可以掌握单 LLM 无 UIWorker的异步任务模式理解四种信封group_started/job_update/job_completed/group_completed的产生与消费链路以及取消事件的协议细节。核心模式一次 LLM 工具调用扇出到多个 peer worker示例的核心架构引自 bot.py 顶部文档如下Main worker (PipelineWorker, owns transport RTVI): transport.in → STT → user_agg → LLM → TTS → transport.out → assistant_agg └── research(query) tool └── ui_jobs.request_job_group( # found by name on the runner wikipedia, news, scholar, paramsJobGroupParams(payload..., label...)) ui_jobs (BaseUIWorker): the client-visible job-group dispatcher (no LLM) Three peer workers (BaseWorker each): WikipediaResearcher · NewsResearcher · ScholarResearcher这个设计有三个关键决策分发器本身不带 LLM。BaseUIWorker是一个纯 bus worker其继承自BaseWorker的run()只是 bus 循环可以直接实例化并注册到 runner 上工具的代码通过params.worker_runner.get_worker(ui-jobs)按名字找到它。也就是说把后台工作变成客户端可见的卡片这件事完全由协议机制承担不需要第二个 LLM 在页面上驱动 UI。fire-and-forget 派发。request_job_group只等待所有目标 worker 就绪并投递请求随后立即返回job_id不等待任务组完成。因此 LLM 可以马上说出一句口头确认Researching the Mariana Trench now.然后在任务还在跑的时候继续接受用户的下一轮对话。进度与结果自动流到客户端。分发的每一个任务组其完整生命周期开始、每个 worker 的进度、完成、整组结束都会被转发为ui-job-group信封到达客户端页面上以进行中的卡片渲染每个 worker 的状态。示例中三个 peer worker 是刻意仿真的用asyncio.sleep随机 0.4–1.5 秒不等加固定文案模拟搜索 → 找到 N 条结果 → 总结 → 完成的过程让演示聚焦于协议本身而非 AI 能力。服务端实现research 工具与任务组派发工具函数如何从 LLM 侧发起任务组bot.py 中注册的工具完整代码如下tool_options(cancel_on_interruptionFalse) async def research(params: FunctionCallParams, query: str): Start background research on a topic across three worker sources. Dispatches the workers fire-and-forget: the groups progress and results stream to the client as ui-job-group envelopes, so this tool returns immediately and the LLM speaks a short acknowledgement. logger.info(fresearch({query})) ui_jobs: BaseUIWorker params.worker_runner.get_worker(ui-jobs) job_id await ui_jobs.request_job_group( wikipedia, news, scholar, paramsJobGroupParams( payload{query: query}, labelfResearch: {query}, ), ) await params.result_callback( { status: started, job_id: job_id, note: Workers run in the background; results stream to the users screen., } )要点params.worker_runner.get_worker(ui-jobs)按名字从WorkerRunner上查找分发器——分发器与工具之间没有显式依赖注入靠 runner 上的名字注册表解耦。JobGroupParams定义在 src/pipecat/pipeline/job_context.py的字段决定了整组任务的行为结合源码可以确认的取值有payload结构化任务数据本例为{query: query}peer worker 在on_job_request里从message.payload取出querylabel人类可读的任务标题BaseUIWorker会用它给客户端的进度卡片命名本例是Research: {query}cancellable默认True是否允许外部请求方如客户端 UI取消该组。worker 自己发起的取消shutdown、超时、cancel_on_error不受此限制cancel_on_errorJobGroupParams专有默认True某个 worker 以 error 状态响应时是否取消整组name/timeout可选的任务名路由与超时。注意tool_options(cancel_on_interruptionFalse)用户在工具执行期间打断时不会取消工具——因为任务本来就应该继续后台运行。工具最后通过params.result_callback返回job_id让 LLM 知道任务已启动、结果会在屏幕上出现但提示词明确告诉它不要等待结果。主管线与 worker 注册主语音管线是一条标准的Pipelinetransport.input() → DeepgramSTT → user 聚合器 → OpenAI LLM → Cartesia TTS → transport.output() → assistant 聚合器用PipelineWorker包装开启 metrics 与 usage metrics。四个 worker 通过runner.add_workers(...)注册到同一个WorkerRunnerui_jobs BaseUIWorker(ui-jobs) worker PipelineWorker( pipeline, nameMAIN_NAME, paramsPipelineParams(enable_metricsTrue, enable_usage_metricsTrue), idle_timeout_secsrunner_args.pipeline_idle_timeout_secs, processor_unusable_policyProcessorUnusablePolicy.END, ) runner WorkerRunner(handle_sigintrunner_args.handle_sigint) await runner.add_workers( ui_jobs, WikipediaResearcher(wikipedia), NewsResearcher(news), ScholarResearcher(scholar), worker, )从源码结构看所有 worker 挂在同一个 bus 上交换消息生命周期事件、帧、job RPC——这正是 Pipecat multi-worker 模式的通用形态总览见 multi-worker 示例索引。peer worker进度更新与最终响应三个研究员 worker 共享一个基类_SimulatedResearcher(BaseWorker)核心逻辑是重写on_job_requestasync def on_job_request(self, message: BusJobRequestMessage) - None: await super().on_job_request(message) job_id message.job_id query (message.payload or {}).get(query, ) try: await asyncio.sleep(random.uniform(0.4, 1.2)) await self.send_job_update(job_id, {text: fsearching {self.source_name}…}) await asyncio.sleep(random.uniform(0.6, 1.4)) n random.randint(3, 8) await self.send_job_update(job_id, {text: ffound {n} results}) await asyncio.sleep(random.uniform(0.5, 1.5)) await self.send_job_update(job_id, {text: summarizing}) await asyncio.sleep(random.uniform(0.4, 0.9)) await self.send_job_response(job_id, response{summary: self.summarize(query)}) except asyncio.CancelledError: # The base workers cancellation hook auto-emits a CANCELLED # response; just bail. raise这里体现了 job 协议的三个 API均在 src/pipecat/workers/base_worker.py 中定义send_job_update(job_id, data)发送中间进度data是任意 dict本例只放{text: ...}send_job_response(job_id, response...)发送最终响应携带statusJobStatus枚举completed/cancelled/failed/error与响应 payloadsend_job_stream_data(job_id, data)面向渐进式输出的流式数据通道本示例没有用到属于 README 明确列出的未展示能力。随机asyncio.sleep让三个 worker 以不同速度推进正好把页面卡片逐条刷新的流式效果展示出来。被取消时 worker 捕获asyncio.CancelledError直接 re-raise基类的取消钩子会自动发出CANCELLED状态的响应。BaseUIWorker四种信封如何产生BaseUIWorker的全部增量逻辑都在 src/pipecat/workers/base_ui_worker.py 中它重写了BaseWorker的 job 生命周期钩子把内部 bus 消息翻译成面向客户端的ui-job-group信封信封 kind触发时机携带的关键字段group_startedcreate_job_group_and_request_job完成、任务组登记后job_id、workersworker 名列表、label、cancellablejob_update任一 worker 调用send_job_update经分发器的on_job_update钩子转发job_id、worker_name、data进度内容job_completed任一 worker 调用send_job_responseon_job_response钩子或以结束流的方式收尾on_job_stream_end钩子经_send_job_completed发布job_id、worker_name、status、responsegroup_completed整组任务终结全部完成、被取消或超时经_send_group_completed发布job_id几个值得注意的实现细节group_started是显式发布的BaseUIWorker.create_job_group_and_request_job先调用父类创建并派发任务组再send_bus_message(BusUIJobGroupStartedMessage(...))所以卡片上出现哪些 worker在派发瞬间就确定了。job_completed与group_completed的幂等边界on_job_update/on_job_response/on_job_stream_end在转发前都会检查message.job_id not in self._job_groups——因为被取消的组已经先行拆除其 worker 迟到的消息不能再让已关闭的卡片重新变动。取消时确定性补发终态BaseUIWorker.cancel_job_group在调用父类拆解任务组之前先捕获group对象由于各 worker 自己的CANCELLED响应要等组消失之后才会到达会重复它改为对每个尚未终态的 worker 直接合成一条statuscancelled的job_completed信封然后发布唯一的group_completed。这样客户端不会因为竞态看到重复或漏掉的终态。取消入口的权限门客户端事件走on_bus_message→_handle_cancel_job_event最终调用BaseWorker.request_cancel_job_group它只在组存在且cancellableTrue时才真正执行cancel_job_group否则记录日志并忽略。worker 自身的取消超时、cancel_on_error则直接走cancel_job_group永远不会被拒绝。信封消息本身BusUIJobGroupStartedMessage等定义在 src/pipecat/bus/ui/messages.py保留的客户端取消事件名是常量__cancel_job_group_UI_CANCEL_JOB_GROUP_BUS_EVENT_NAME。客户端消费RTVIEvent.UIJobGroup 与取消浏览器端 client/main.js 是一个纯 vanilla JS 的 Vite 应用pipecat-ai/client-jspipecat-ai/small-webrtc-transport连接http://localhost:7860/api/offer。与页面上再放一个 LLM的 UIWorker 类示例不同它只新增一件事——订阅RTVIEvent.UIJobGroupclient.on(RTVIEvent.UIJobGroup, handleJobGroupEnvelope);客户端维护一个Mapjob_id, group状态表每个 group 记录label、cancellable和workers: Mapworker_name, {status, update, response}并按信封 kind 分发处理function handleJobGroupEnvelope(env) { switch (env.kind) { case group_started: { // 建 Map、渲染带 Cancel 按钮的进行中卡片 tasksList.appendChild(renderGroupCard(group)); break; } case job_update: { // 更新对应 worker 行的进度文本env.data?.text updateWorkerRow(group, env.worker_name, { update: text }); break; } case job_completed: { // 写入 status 与 response行状态变为 completed/cancelled/... break; } case group_completed: { // 把进行中卡片提升为结果卡片移入 results 面板 resultsList.prepend(renderResultsForGroup(group)); group.cardEl.remove(); groups.delete(env.job_id); break; } } }渲染规则group_started时卡片标题取env.label即服务端的JobGroupParams.label每个 worker 一行初始状态running进度文本starting…cancellable为真时卡片头部出现 Cancel 按钮点击后调用client.cancelUIJobGroup(group.job_id, user requested)并把按钮置为禁用态group_completed时进行中的卡片被提升为结果卡片统计各 worker 终态计数completed / cancelled / failed / error并把每个completedworker 的response.summary或text否则整体 JSON搬进结果面板。取消回路的完整链路是client.cancelUIJobGroup(job_id, reason)→ 向服务端发送保留事件__cancel_job_grouppayload 含job_id与reason→ 分发器BaseUIWorker的on_bus_message捕获并翻译成cancel_job_group(job_id)→ 父类向组内每个 worker 广播BusJobCancelMessage并标记组失败 → 被取消的 worker 任务抛出CancelledError终态按前述方式确定性补发 → 页面卡片关闭。运行方式运行前先按 multi-worker 示例总说明 准备好环境在仓库根执行uv sync --all-extras然后在examples/multi-worker下复制 env.example 为.env填入密钥示例自己的.env也可以bot 启动时会load_dotenv(overrideTrue)。本示例需要三个变量变量用途OPENAI_API_KEY主 LLMOpenAILLMServiceDEEPGRAM_API_KEYSTTDeepgramSTTServiceCARTESIA_API_KEYTTSCartesiaTTSService默认语音可用CARTESIA_VOICE_ID覆盖两个终端终端 1 — botcd examples/multi-worker/ui-worker/async-tasks uv run bot.pybot 启动在http://localhost:7860。终端 2 — clientcd examples/multi-worker/ui-worker/async-tasks/client npm install # one-time npm run dev打开http://localhost:5173点击Connect。建议的交互验证由于 worker 是仿真的固定摘要 随机延迟每次研究调用大约持续几秒适合逐步验证协议行为Research the Mariana Trench.— 分发器扇出三个 peerLLM 只回一句确认页面出现一张卡片逐条显示每个 peer 的状态变化searching → found N results → summarizing → completedLook up octopus cognition.— 同样的流程第二张卡片叠加出现Research the moon, then research Mars.— 两个任务组并发运行页面上两张卡片独立推进How are you?不触发研究— 快速口头回答不产生任务组对进行中的卡片点击 Cancel— 取消路由完整走通peer 任务抛出CancelledError对应行状态回显为cancelled卡片随后关闭并进入结果面板。选型BaseUIWorker 还是 UIWorkerREADME 对这一示例相对于此前 UIWorker 示例新增了什么的定位很明确其他 UI 示例是把一个LLM 放在页面上UIWorker读取页面快照、驱动 UI 交互而本示例证明流式任务组这半个协议完全不需要它——一个BaseUIWorker分发器扇出 peer worker客户端自行渲染进度即可。bot.py 的文档也给出同样的选型准则当委托方需要读取或驱动页面内容快照、deixis、UI 命令时用UIWorker可参考 document-review 示例当页面只是后台工作的视图时用本例这种BaseUIWorker。需要说明的是UIWorker正是BaseUIWorker的子类在前者基础上叠加 LLM 驱动的页面交互能力——所以任务组生命周期转发、取消协议这些机制在UIWorker应用中同样可用二者共享同一套信封协议。本示例不覆盖的范围README 明确列出了刻意留白的部分可作为后续扩展方向真实数据源集成peer 目前是仿真的实际应用中每个 worker 接自己的数据源LLM 驱动的 peer本例 peer 是纯数据抓取型的BaseWorker但 peer 本身也可以是LLMWorker流式块输出用send_job_stream_data实现渐进式内容如逐段流出的长文worker 到 worker 的扇出嵌套任务组一个 worker 收到任务后再向下游 worker 分派。任务组机制本身的更细行为JobGroup/JobContext/JobGroupContext等待与事件迭代语义、超时处理、错误传播可以在 src/pipecat/pipeline/job_context.py 与 tests/test_job_group.py 中继续深入。【免费下载链接】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),仅供参考
返回列表