Pipecat 实战:用 BaseUIWorker 扇出异步任务组,把后台工作流式渲染到 Web 客户端
2026/9/14 1:25:12 网站建设 项目流程

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", params=JobGroupParams(payload=..., label=...)) ui_jobs (BaseUIWorker): the client-visible job-group dispatcher (no LLM) Three peer workers (BaseWorker each): WikipediaResearcher · NewsResearcher · ScholarResearcher

这个设计有三个关键决策:

  1. 分发器本身不带 LLM。BaseUIWorker是一个纯 bus worker(其继承自BaseWorkerrun()只是 bus 循环),可以直接实例化并注册到 runner 上,工具的代码通过params.worker_runner.get_worker("ui-jobs")按名字找到它。也就是说,"把后台工作变成客户端可见的卡片"这件事完全由协议机制承担,不需要第二个 LLM 在页面上驱动 UI。
  2. fire-and-forget 派发。request_job_group只等待所有目标 worker 就绪并投递请求,随后立即返回job_id,不等待任务组完成。因此 LLM 可以马上说出一句口头确认("Researching the Mariana Trench now."),然后在任务还在跑的时候继续接受用户的下一轮对话。
  3. 进度与结果自动流到客户端。分发的每一个任务组,其完整生命周期(开始、每个 worker 的进度、完成、整组结束)都会被转发为ui-job-group信封到达客户端,页面上以"进行中的卡片"渲染每个 worker 的状态。

示例中三个 peer worker 是刻意"仿真"的:用asyncio.sleep(随机 0.4–1.5 秒不等)加固定文案模拟"搜索 → 找到 N 条结果 → 总结 → 完成"的过程,让演示聚焦于协议本身而非 AI 能力。

服务端实现:research 工具与任务组派发

工具函数:如何从 LLM 侧发起任务组

bot.py 中注册的工具完整代码如下:

@tool_options(cancel_on_interruption=False) async def research(params: FunctionCallParams, query: str): """Start background research on a topic across three worker sources. Dispatches the workers fire-and-forget: the group's 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(f"research('{query}')") ui_jobs: BaseUIWorker = params.worker_runner.get_worker("ui-jobs") job_id = await ui_jobs.request_job_group( "wikipedia", "news", "scholar", params=JobGroupParams( payload={"query": query}, label=f"Research: {query}", ), ) await params.result_callback( { "status": "started", "job_id": job_id, "note": "Workers run in the background; results stream to the user's 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取出query
    • label:人类可读的任务标题,BaseUIWorker会用它给客户端的进度卡片命名(本例是Research: {query});
    • cancellable(默认True):是否允许外部请求方(如客户端 UI)取消该组。worker 自己发起的取消(shutdown、超时、cancel_on_error)不受此限制;
    • cancel_on_errorJobGroupParams专有,默认True):某个 worker 以 error 状态响应时是否取消整组;
    • name/timeout:可选的任务名路由与超时。
  • 注意@tool_options(cancel_on_interruption=False):用户在工具执行期间打断时不会取消工具——因为任务本来就应该继续后台运行。
  • 工具最后通过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(...)注册到同一个WorkerRunner

ui_jobs = BaseUIWorker("ui-jobs") worker = PipelineWorker( pipeline, name=MAIN_NAME, params=PipelineParams(enable_metrics=True, enable_usage_metrics=True), idle_timeout_secs=runner_args.pipeline_idle_timeout_secs, processor_unusable_policy=ProcessorUnusablePolicy.END, ) runner = WorkerRunner(handle_sigint=runner_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_request

async 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": f"searching {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": f"found {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 worker's 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)与响应 payload;
  • send_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_idworkers(worker 名列表)、labelcancellable
job_update任一 worker 调用send_job_update,经分发器的on_job_update钩子转发job_idworker_namedata(进度内容)
job_completed任一 worker 调用send_job_responseon_job_response钩子)或以结束流的方式收尾(on_job_stream_end钩子),经_send_job_completed发布job_idworker_namestatusresponse
group_completed整组任务终结(全部完成、被取消或超时),经_send_group_completed发布job_id

几个值得注意的实现细节:

  • group_started是显式发布的BaseUIWorker.create_job_group_and_request_job先调用父类创建并派发任务组,再send_bus_message(BusUIJobGroupStartedMessage(...)),所以"卡片上出现哪些 worker"在派发瞬间就确定了。
  • job_completedgroup_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 直接合成一条status=cancelledjob_completed信封,然后发布唯一的group_completed。这样客户端不会因为竞态看到重复或漏掉的终态。
  • 取消入口的权限门:客户端事件走on_bus_message_handle_cancel_job_event,最终调用BaseWorker.request_cancel_job_group,它只在组存在且cancellable=True时才真正执行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-js+@pipecat-ai/small-webrtc-transport),连接http://localhost:7860/api/offer。与"页面上再放一个 LLM"的 UIWorker 类示例不同,它只新增一件事——订阅RTVIEvent.UIJobGroup

client.on(RTVIEvent.UIJobGroup, handleJobGroupEnvelope);

客户端维护一个Map<job_id, group>状态表,每个 group 记录labelcancellableworkers: Map<worker_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_group(payload 含job_idreason)→ 分发器BaseUIWorkeron_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(override=True))。本示例需要三个变量:

变量用途
OPENAI_API_KEY主 LLM(OpenAILLMService
DEEPGRAM_API_KEYSTT(DeepgramSTTService
CARTESIA_API_KEYTTS(CartesiaTTSService,默认语音可用CARTESIA_VOICE_ID覆盖)

两个终端:

终端 1 — bot:

cd examples/multi-worker/ui-worker/async-tasks uv run bot.py

bot 启动在http://localhost:7860

终端 2 — client:

cd examples/multi-worker/ui-worker/async-tasks/client npm install # one-time npm run dev

打开http://localhost:5173,点击Connect

建议的交互验证

由于 worker 是仿真的(固定摘要 + 随机延迟),每次研究调用大约持续几秒,适合逐步验证协议行为:

  • "Research the Mariana Trench."— 分发器扇出三个 peer,LLM 只回一句确认,页面出现一张卡片,逐条显示每个 peer 的状态变化(searching → found N results → summarizing → completed);
  • "Look up octopus cognition."— 同样的流程,第二张卡片叠加出现;
  • "Research the moon, then research Mars."— 两个任务组并发运行,页面上两张卡片独立推进;
  • "How are you?"(不触发研究)— 快速口头回答,不产生任务组;
  • 对进行中的卡片点击 Cancel— 取消路由完整走通:peer 任务抛出CancelledError,对应行状态回显为cancelled,卡片随后关闭并进入结果面板。

选型:BaseUIWorker 还是 UIWorker

README 对这一示例"相对于此前 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询