You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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流:

  1. 前端浏览器环境不支持Node.js的stream对象,拿到的是序列化后的空对象
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 08:24:56