T3 Stack中TRPC WebSocket订阅报错:Subscriptions should use wsLink
TRPC WebSocket订阅功能报错:Subscriptions should use wsLink
问题现象
在T3 Stack中使用TRPC WebSocket,其余功能正常,但useSubscription无法工作,报错:
Stream error: TRPCClientError: Subscriptions should use wsLink
调用streamChatbot功能可正常获取结果,但无法接收chatWSRouter中的事件推送。
依赖版本:
"@trpc/client": "10.45.0", "@trpc/next": "10.45.0", "@trpc/react-query": "10.45.0", "@trpc/server": "10.45.0",
相关代码
客户端WS配置(src/utils/ws.ts)
import { createWSClient, loggerLink, wsLink } from '@trpc/client'; import { createTRPCNext } from '@trpc/next'; import getConfig from 'next/config'; import superjson from 'superjson'; import { AppRouterWS } from '~/server/api/root'; const { publicRuntimeConfig } = getConfig(); const { WS_URL } = publicRuntimeConfig; const client = createWSClient({ url: WS_URL, }); export const ws = createTRPCNext<AppRouterWS>({ ssr: false, config() { return { links: [ loggerLink({ enabled: (opts) => (process.env.NODE_ENV === 'development' && typeof window !== 'undefined') || (opts.direction === 'down' && opts.result instanceof Error), }), wsLink({ client, }), ], transformer: superjson, }; }, });
useSubscription使用代码
import { ws } from '~/utils/ws'; ws.chatWS.stream.useSubscription(undefined, { onStarted() { console.log('STARTING'); }, onData(data: { message: any; status: 'STREAMING' | 'ERROR' }) { console.log(data); setMessages((prev) => [ ...prev, { id: currentMessageId, type: 'response', content: data.message, status: data.status === 'STREAMING' ? 'streaming' : data.status, error: data.status === 'ERROR', }, ]); }, onError(error) { console.log('ASDASD'); console.error('Stream error:', error); }, });
后端路由(src/server/api/routers/websocket/chatWSRouter.ts)
import { observable } from '@trpc/server/observable'; import { z } from 'zod'; import { createTRPCRouter, publicProcedure } from '../../trpc'; import awsWsClient from '~/server/wssAwsServer'; export const chatWSRouter = createTRPCRouter({ invokeChatbot: publicProcedure .input(z.object({ prompt: z.string() })) .mutation(async ({ input }) => { if (awsWsClient.isConnected()) { awsWsClient.send({ action: 'streamChatbot', message: input.prompt }); } else { console.error('WebSocket is not connected.'); } }), invokeAgent: publicProcedure .input(z.object({ prompt: z.string(), agentAliasId: z.string() })) .mutation(async ({ input }) => {}), stream: publicProcedure.subscription(() => { return observable<{ message: string; status: 'STREAMING' | 'ERROR' }>((emit) => { const onMessage = (data: { message: string; status: 'STREAMING' | 'ERROR' }) => { console.log(data); emit.next(data); }; awsWsClient.on('stream_chat', onMessage); }); }), });
后端根路由(src/server/api/root.ts)
import { fieldRouter } from './routers/rest/fieldRouter'; import { formRouter } from './routers/rest/formRouter'; import { formTemplateRouter } from './routers/rest/formTemplate'; import { transcriptionRouter } from './routers/rest/transcriptionRouter'; import { userGroupRouter } from './routers/rest/userGroupRouter'; import { userRouter } from './routers/rest/userRouter'; import { narrativeRouter } from './routers/rest/narrativeRouter'; import { createTRPCRouter, mergeRouters } from './trpc'; import { narrativeTemplateRouter } from './routers/rest/narrativeTemplate'; import { eventRouter } from './routers/rest/eventRouter'; import { chatWSRouter } from './routers/websocket/chatWSRouter'; import { formWSRouter } from './routers/websocket/formWSRouter'; import { narrativeWSRouter } from './routers/websocket/narrativeWSRouter'; import { transcribeWSRouter } from './routers/websocket/transcribeWSRouter'; /** * This is the primary router for your server. * * All routers added in /api/routers should be manually added here. */ export const appRouter = createTRPCRouter({ formTemplate: formTemplateRouter, field: fieldRouter, form: formRouter, transcription: transcriptionRouter, user: userRouter, userGroup: userGroupRouter, narrative: narrativeRouter, narrativeTemplate: narrativeTemplateRouter, events: eventRouter, }); export const appRouterWS = createTRPCRouter({ chatWS: chatWSRouter, formWS: formWSRouter, narrativeWS: narrativeWSRouter, transcibreWS: transcribeWSRouter, }); export const combinedRouter = mergeRouters(appRouter, appRouterWS); // export type definition of API export type AppRouter = typeof appRouter; export type AppRouterWS = typeof appRouterWS; export type CombinedRouter = typeof combinedRouter;
解决方案
1. 配置后端WebSocket服务器
在Next.js的TRPC API路由(src/pages/api/trpc/[trpc].ts)中,添加WebSocket适配器配置,确保服务器同时支持HTTP和WebSocket:
import { createNextApiHandler } from '@trpc/server/adapters/next'; import { combinedRouter } from '~/server/api/root'; import { createTRPCContext } from '~/server/api/trpc'; import { applyWSSHandler } from '@trpc/server/adapters/ws'; import { WebSocketServer } from 'ws'; const handler = createNextApiHandler({ router: combinedRouter, createContext: createTRPCContext, }); // 启动WebSocket服务,端口与客户端WS_URL一致 const wss = new WebSocketServer({ port: 3001 }); applyWSSHandler({ wss, router: combinedRouter, createContext: createTRPCContext }); export default handler;
2. 优化客户端链路配置
使用splitLink区分订阅请求和普通HTTP请求,确保订阅走WebSocket链路:
import { createWSClient, loggerLink, wsLink, splitLink, httpLink } from '@trpc/client'; // ...其他导入 export const ws = createTRPCNext<AppRouterWS>({ ssr: false, config() { return { links: [ loggerLink({ enabled: (opts) => (process.env.NODE_ENV === 'development' && typeof window !== 'undefined') || (opts.direction === 'down' && opts.result instanceof Error), }), splitLink({ condition(op) { return op.type === 'subscription'; }, true: wsLink({ client }), false: httpLink({ url: '/api/trpc' }), }), ], transformer: superjson, }; }, });
3. 修复Observable订阅清理逻辑
在后端stream订阅中添加清理函数,避免内存泄漏和重复监听:
stream: publicProcedure.subscription(() => { return observable<{ message: string; status: 'STREAMING' | 'ERROR' }>((emit) => { const onMessage = (data: { message: string; status: 'STREAMING' | 'ERROR' }) => { console.log(data); emit.next(data); }; awsWsClient.on('stream_chat', onMessage); // 客户端取消订阅时清理监听 return () => { awsWsClient.off('stream_chat', onMessage); }; }); }),
4. 验证WS_URL正确性
确保客户端WS_URL格式正确,例如开发环境为ws://localhost:3001,生产环境对应部署域名,且WebSocket服务器确实在该端口运行。
内容的提问来源于stack exchange,提问作者Lane Floyd
相关产品推荐
相关产品推荐

