如何通过tRPC将OpenAI流式响应传输至React客户端?
问题:tRPC传输OpenAI流式响应至前端失败
问题背景
已在Next.js中实现OpenAI流式补全的本地处理(使用openai-node,开启stream: true和responseType: 'stream'),但尝试通过tRPC将流式响应传输到前端React时,客户端无法识别stream.on方法,无法接收流式数据。
核心问题
tRPC的默认mutation/query基于HTTP请求实现,会等待完整响应返回后才传递给客户端,无法直接返回Node.js的Readable流:
- 前端浏览器环境不支持Node.js的
stream对象,拿到的是序列化后的空对象 - tRPC默认机制不支持流式数据的实时推送
可行解决方案
方案1:使用tRPC订阅(Subscription)
tRPC的订阅基于WebSocket,适合实时推送流式数据,是官方推荐的流式传输方式,且能兼容你代码中的authenticatedProcedure鉴权逻辑。
后端修改(tRPC Router)
import { initOpenAI } from 'lib/ai'; import { subscriptionProcedure, router } from './trpc'; // 确保已配置订阅支持 export const analysisRouter = router({ generate: subscriptionProcedure .input(z.object({ type: z.string() // 你的输入参数 })) .subscription(async ({ ctx, input }) => { const openai = initOpenAI(); let activeRole = ''; // 返回异步迭代器,tRPC通过WebSocket推送数据 return { async *[Symbol.asyncIterator]() { const result = await openai.createChatCompletion({ messages: [ { role: 'user', content: 'hello there!' } ], model: 'gpt-3.5-turbo', temperature: 0.85, stream: true }, { responseType: 'stream' }); // 监听OpenAI流数据 result.data.on('data', (data: Buffer) => { const lines = data.toString().split('\n').filter(line => line.trim() !== ''); for (const line of lines) { const message = line.replace(/^data: /, ''); if (message === '[DONE]') return; const parsed = JSON.parse(message); if (parsed.choices[0].finish_reason === 'stop') return; activeRole = parsed.choices[0].delta.role ?? activeRole; if (parsed.choices[0].delta.content) { // 通过迭代器推送chunk this.next({ role: activeRole, content: parsed.choices[0].delta.content }); } } }); // 处理流结束 result.data.on('end', () => { this.return(); }); } }; }) });
注意:需确保tRPC服务器已配置WebSocket支持(如Next.js中使用
tRPC Next.js adapter的WebSocket配置)
前端修改(React组件)
import { trpc } from '../utils/trpc'; function ChatComponent() { // 用useSubscription监听流式数据 const { data: chunk } = trpc.analysis.generate.useSubscription({ type: 'lesson' }); // 每次收到新chunk时更新UI useEffect(() => { if (chunk) { console.log(chunk); // 此处更新聊天内容状态 } }, [chunk]); return <div>{/* 渲染实时聊天内容 */}</div>; }
方案2:使用Next.js API路由 + Server-Sent Events(SSE)
若不想配置WebSocket,可绕开tRPC,直接用Next.js API路由结合SSE实现流式传输:
后端API路由(pages/api/stream.ts)
import type { NextApiRequest, NextApiResponse } from 'next'; import { initOpenAI } from 'lib/ai'; export default async function handler(req: NextApiRequest, res: NextApiResponse) { // 设置SSE响应头 res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive', // 若需要鉴权,可从请求头获取token验证 }); const openai = initOpenAI(); let activeRole = ''; const result = await openai.createChatCompletion({ messages: [{ role: 'user', content: 'hello there!' }], model: 'gpt-3.5-turbo', temperature: 0.85, stream: true }, { responseType: 'stream' }); result.data.on('data', (data: Buffer) => { const lines = data.toString().split('\n').filter(line => line.trim() !== ''); for (const line of lines) { const message = line.replace(/^data: /, ''); if (message === '[DONE]') { res.write('event: done\ndata: {}\n\n'); res.end(); return; } const parsed = JSON.parse(message); if (parsed.choices[0].finish_reason === 'stop') return; activeRole = parsed.choices[0].delta.role ?? activeRole; if (parsed.choices[0].delta.content) { const chunk = JSON.stringify({ role: activeRole, content: parsed.choices[0].delta.content }); // 发送SSE数据 res.write(`data: ${chunk}\n\n`); } } }); result.data.on('end', () => { res.end(); }); // 处理客户端断开连接,清理资源 req.on('close', () => { result.data.destroy(); res.end(); }); }
前端React组件
useEffect(() => { const eventSource = new EventSource('/api/stream'); // 接收流式数据 eventSource.onmessage = (event) => { const chunk = JSON.parse(event.data); console.log(chunk); // 更新UI状态 }; // 监听流结束事件 eventSource.addEventListener('done', () => { eventSource.close(); }); // 组件卸载时关闭连接 return () => { eventSource.close(); }; }, []);
总结
- tRPC默认的mutation/query不支持流式传输,必须使用**订阅(WebSocket)**或绕开tRPC用SSE
- 订阅方式更贴合tRPC生态,自动兼容鉴权逻辑,适合需要会话验证的场景
- SSE方式实现更简单,无需配置WebSocket,但鉴权需手动处理(如在请求头传递token)
内容的提问来源于stack exchange,提问作者thim24
相关产品推荐
相关产品推荐

