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这个设计有三个关键决策:
- 分发器本身不带 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_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_error(JobGroupParams专有,默认True):某个 worker 以 error 状态响应时是否取消整组;name/timeout:可选的任务名路由与超时。
- 注意
@tool_options(cancel_on_interruption=False):用户在工具执行期间打断时不会取消工具——因为任务本来就应该继续后台运行。 - 工具最后通过
params.result_callback返回job_id,让 LLM 知道"任务已启动、结果会在屏幕上出现",但提示词明确告诉它不要等待结果。
主管线与 worker 注册
主语音管线是一条标准的Pipeline:transport.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=...):发送最终响应,携带status(JobStatus枚举: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_started | create_job_group_and_request_job完成、任务组登记后 | job_id、workers(worker 名列表)、label、cancellable |
job_update | 任一 worker 调用send_job_update,经分发器的on_job_update钩子转发 | job_id、worker_name、data(进度内容) |
job_completed | 任一 worker 调用send_job_response(on_job_response钩子)或以结束流的方式收尾(on_job_stream_end钩子),经_send_job_completed发布 | job_id、worker_name、status、response |
group_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 直接合成一条status=cancelled的job_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 记录label、cancellable和workers: 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_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(override=True))。本示例需要三个变量:
| 变量 | 用途 |
|---|---|
OPENAI_API_KEY | 主 LLM(OpenAILLMService) |
DEEPGRAM_API_KEY | STT(DeepgramSTTService) |
CARTESIA_API_KEY | TTS(CartesiaTTSService,默认语音可用CARTESIA_VOICE_ID覆盖) |
两个终端:
终端 1 — bot:
cd examples/multi-worker/ui-worker/async-tasks uv run bot.pybot 启动在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),仅供参考