Pipecat 集成

使用 Pipecat pipeline 作为 Speech Engine 背后的 LLM 大脑。

本指南介绍如何在 Speech Engine 大脑服务器中使用 Pipecat 作为 LLM pipeline。Speech Engine 负责语音流程——语音转文本、轮次管理和文本转语音;Pipecat 则通过由处理器组成的可组合 pipeline 处理文本生成(LLM 调用、RAG、函数调用、护栏、内容过滤器)。

本指南仅适用于 Python,因为 Pipecat 服务端是 Python 框架。pipeline 处理器没有 Node 等效实现;虽然有 pipecat-client-js 包,但它是与 Pipecat 服务器通信的浏览器客户端, 并不能用于通过 TypeScript 构建 pipeline。

架构

Speech Engine SDK 作为外层运行——每当用户说完话,都会触发其 on_transcript 回调。在回调中,你需要构建 Pipecat pipeline,将对话历史作为 LLMContextFrame 输入,并将 pipeline 的文本输出流式传回 Speech Engine。ElevenLabs 会将文本转为语音并播放给用户。

User speaks (audio) on_transcript(history) LLMContextFrame(history) LLMTextFrame chunks send_response(async iterator) Agent speaks (audio) Browser ElevenLabs Brain Server (engine.serve) Pipecat Pipeline

Pipecat pipeline 仅在单个轮次期间运行。新转录文本到达时,上一个 pipeline 会在下一个运行前取消——Speech Engine 的打断处理机制便以此方式传递到 pipeline 中。

何时使用此模式

当大脑不仅需要一次 LLM 调用时,Pipecat 会很有用:

  • 用于检索增强生成、函数调用或护栏的可组合处理器
  • 可在每一步检查、转换或阻止流量的基于帧的中间件
  • 可在多个智能体间共享的可复用 pipeline 片段

如果大脑只是“输入转录文本,输出 LLM 调用”,Speech Engine 快速入门会更简单。当 pipeline 本身才是重点时,再选择 Pipecat。

前提条件

  • 一个 Speech Engine。请按Speech Engine 快速入门创建。
  • Python 3.10+(pipecat-ai 的要求)。
  • 面向大脑服务器的公共 HTTPS 隧道(例如 ngrok)。

安装依赖

pip install "pipecat-ai[openai]" "elevenlabs" "python-dotenv"

pipecat-ai[openai] 会安装 OpenAI LLM 服务。如果你偏好其他提供商,可替换为相应扩展(anthropic、google 等)。

构建 Pipecat 大脑

大脑由两部分组成:将流式文本写入 asyncio.Queue 的 TextSink 处理器,以及构建单轮 pipeline 并以异步迭代器形式产出分块的 run_pipecat_brain 协程。

brain.py
import asyncio
import os
from typing import AsyncIterator
from dotenv import load_dotenv
from pipecat.frames.frames import (
Frame,
LLMContextFrame,
LLMFullResponseEndFrame,
LLMTextFrame,
)
from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.services.openai.llm import OpenAILLMService
load_dotenv()
SYSTEM_PROMPT = (
"You are a helpful voice assistant. Keep responses concise and conversational."
)
class TextSink(FrameProcessor):
"""Drain LLMTextFrame text into an asyncio.Queue."""
def __init__(self, queue: asyncio.Queue):
super().__init__()
self._queue = queue
async def process_frame(self, frame: Frame, direction: FrameDirection):
await super().process_frame(frame, direction)
if isinstance(frame, LLMTextFrame):
await self._queue.put(frame.text)
elif isinstance(frame, LLMFullResponseEndFrame):
await self._queue.put(None) # sentinel
await self.push_frame(frame, direction)
def build_messages(transcript: list[dict]) -> list[dict]:
messages = [{"role": "system", "content": SYSTEM_PROMPT}]
for turn in transcript:
role = "assistant" if turn["role"] == "agent" else turn["role"]
messages.append({"role": role, "content": turn["content"]})
return messages
async def run_pipecat_brain(transcript: list[dict]) -> AsyncIterator[str]:
"""Yield response text chunks from a one-turn Pipecat pipeline."""
llm = OpenAILLMService(
api_key=os.environ["OPENAI_API_KEY"],
model="gpt-4o-mini",
)
queue: asyncio.Queue[str | None] = asyncio.Queue()
sink = TextSink(queue)
task = PipelineTask(Pipeline([llm, sink]))
runner = PipelineRunner(handle_sigint=False)
async def drive():
context = LLMContext(build_messages(transcript))
await task.queue_frame(LLMContextFrame(context))
await task.stop_when_done()
run_task = asyncio.create_task(runner.run(task))
drive_task = asyncio.create_task(drive())
try:
while True:
chunk = await queue.get()
if chunk is None:
break
yield chunk
finally:
await task.cancel()
await asyncio.gather(run_task, drive_task, return_exceptions=True)

pipeline 仅包含 LLM 服务和 sink,不包含 STT 或 TTS 处理器,因为这些由 Speech Engine 处理。输入为 LLMContextFrame,输出为 LLMTextFrame 分块。

run_pipecat_brain 是异步生成器。每个产出的分块都会直接发送至 Speech Engine,因此无需等待完整响应准备好,智能体便可开始说话。

接入 Speech Engine 服务器

Speech Engine SDK 的 send_response 接受字符串或任意字符串异步可迭代对象,因此可以直接传入 run_pipecat_brain(transcript)。先将 Speech Engine 提供的 ConversationMessage 对象转换为普通字典,再传递给大脑。

server.py
import asyncio
import os
from dotenv import load_dotenv
from elevenlabs import AsyncElevenLabs
from brain import run_pipecat_brain
load_dotenv()
elevenlabs = AsyncElevenLabs(api_key=os.environ["ELEVENLABS_API_KEY"])
SPEECH_ENGINE_ID = os.environ["SPEECH_ENGINE_ID"]
async def on_transcript(transcript, session):
history = [{"role": m.role, "content": m.content} for m in transcript]
await session.send_response(run_pipecat_brain(history))
async def main():
engine = await elevenlabs.speech_engine.get(SPEECH_ENGINE_ID)
await engine.serve(
port=3001,
path="/ws",
debug=True,
on_transcript=on_transcript,
)
if __name__ == "__main__":
asyncio.run(main())

新转录文本到达时,Speech Engine SDK 会取消上一轮任务,这会通过 run_pipecat_brain 中的 try/finally 块取消异步生成器及底层 PipelineTask。

运行服务器

ngrok http 3001
python server.py

使用快速入门中相同的令牌端点和客户端代码,通过浏览器连接 Speech Engine。Pipecat pipeline 在服务器端运行;浏览器看到的是普通的 Speech Engine 对话。

扩展 pipeline

纯文本 Pipecat pipeline 可以包含任何作用于 LLMTextFrame 或 LLMContextFrame 的帧处理器。以下是一些常见扩展:

  • 护栏:置于 LLM 前的 FrameProcessor,用于检查 LLMContextFrame 并替换或阻止不安全的上下文。
  • 函数调用:在 OpenAILLMService 上注册工具,Pipecat 会原生处理工具调用帧。最终的助手文本仍将以 LLMTextFrame 到达。
  • 多阶段推理:串联两个 OpenAILLMService 实例,并在中间添加自定义处理器,为第二次处理重写上下文。
  • 输出过滤:置于 LLM 后的 FrameProcessor,用于检查每个 LLMTextFrame,在内容到达 TextSink 前丢弃或重写不允许的内容。

pipeline 形状保持不变——Pipeline([processor_a, llm, processor_b, sink])——run_pipecat_brain 无需修改。

生产环境注意事项

  • 取消安全性:如果在 pipeline 完全启动前调用,PipelineTask.cancel() 可能导致死锁(pipecat-ai/pipecat#4276)。上述 try/finally 模式是安全的,因为 cancel() 只会在至少一个帧已加入队列后运行。
  • 提示词注入:语音转文本输出属于用户输入。将其传入 LLM 前,请验证或规范化转录文本,尤其是下游处理器会在工具调用或数据库查询中使用文本时。
  • 大脑服务器身份验证:在 Speech Engine 上设置共享密钥,并在大脑服务器中进行检查,防止未经授权的连接访问 /ws 端点:
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • LLM 提供商:pipecat-ai[openai] 包含 OpenAILLMService。若使用 Anthropic,请安装 pipecat-ai[anthropic] 并使用 AnthropicLLMService;pipeline 的其余部分无需更改。

后续步骤