tRPC 实时订阅在 WebSockets 与 SSE 之间怎么选?官方建议与参考项目
tRPC 中把实时更新从服务端推给客户端用的是订阅(subscription),传输层有两条路:WebSockets 和 Server-sent Events(SSE)。二者的服务端订阅过程写法完全一致,差异集中在客户端 link 的选择、服务端额外的搭建步骤和鉴权方式上。本文基于 Subscriptions、WebSockets 和 httpSubscriptionLink 三篇官方文档,以及本仓库内的三个参考示例,说明如何选型,并把两条路径都跑起来。
官方建议是什么
Subscriptions 文档给出了直接结论:不确定用哪个时,官方推荐 SSE,原因是它更容易搭建,不需要单独搭一个 WebSocket 服务器。
反过来,选 WebSockets 的依据是 WebSockets 文档开头的能力说明:WebSockets 可以用于与服务端的全部或部分通信,而不只是订阅。也就是说,当你希望 query / mutation 也走 WebSocket、而不是只让订阅走持久连接时,才需要引入 ws 服务器并改用 wsLink。文档同时提示:也可以只用 Links 把订阅路由到 WebSocket,query 和 mutation 仍走 HTTP。
服务端订阅:两条路径共用
无论选哪种传输,服务端都是同一个写法:subscription 过程里用一个异步生成器 yield 数据。Subscriptions 文档的基本示例:
import EventEmitter, { on } from 'node:events';
import { initTRPC } from '@trpc/server';
const t = initTRPC.create();
type Post = { id: string; title: string };
const ee = new EventEmitter();
export const appRouter = t.router({
onPostAdd: t.procedure.subscription(async function* (opts) {
// listen for new events
for await (const [data] of on(ee, 'add', {
// Passing the AbortSignal from the request automatically cancels the event emitter when the request is aborted
signal: opts.signal,
})) {
const post = data as Post;
yield post;
}
}),
});
文档推荐使用 tracked() 辅助函数给每个事件带上 id:
yield tracked(post.id, post);
这样客户端断开后会自动重连,并回传它最后收到的 ID:SSE 走 EventSource 规范,通过 .input() 里的 lastEventId 到达服务端;WebSockets 则由 wsLink 自动发送并更新。文档提醒:如果基于 lastEventId 从数据库补数据、且不能丢事件,要先注册事件监听再查询(仓库里 next-sse-chat 示例就是这么做的),避免补发批次期间新事件被漏掉。
几个与传输无关的行为:
opts.signal是请求的AbortSignal,客户端断开时会被 abort;- 想从服务端停掉订阅,在生成器里
return即可,客户端只需调用.unsubscribe(); - 需要清理副作用时用
try...finally,tRPC 会在订阅以任何原因停止时调用生成器的.return(); - 错误处理:抛出 5xx 错误时客户端会按
tracked()记录的最后一个事件 ID 自动重连;其他错误会取消订阅并传入onError()回调。
路径一:SSE(httpSubscriptionLink)
httpSubscriptionLink 是一个 terminating link,用 SSE 承载订阅(见 httpSubscriptionLink 文档)。
服务端不需要新服务器,只需要在 initTRPC.create 时配置 sse 选项:
import { initTRPC } from '@trpc/server';
export const t = initTRPC.create({
sse: {
// Maximum duration of a single SSE connection in milliseconds
// maxDurationMs: 60_000,
ping: {
// Enable periodic ping messages to keep connection alive
enabled: true,
// Send ping message every 2s
intervalMs: 2_000,
},
// client: {
// reconnectAfterInactivityMs: 3_000
// }
},
});
其中 client.reconnectAfterInactivityMs 是客户端侧的无活动超时:如果在该时间内没收到任何消息(包括 ping),连接会被标记为 "connecting" 并自动尝试重连,适合与服务端 ping 配合使用。
客户端必须用 splitLink 显式声明"订阅走 SSE":
import {
createTRPCClient,
httpBatchLink,
httpSubscriptionLink,
loggerLink,
splitLink,
} from '@trpc/client';
import type { AppRouter } from './server';
const trpcClient = createTRPCClient<AppRouter>({
/**
* @see https://trpc.io/docs/v11/client/links
*/
links: [
// adds pretty logs to your console in development and logs errors in production
loggerLink(),
splitLink({
// uses the httpSubscriptionLink for subscriptions
condition: (op) => op.type === 'subscription',
true: httpSubscriptionLink({
url: `/api/trpc`, // 文档示例用相对路径,指向你的 tRPC 端点
}),
false: httpBatchLink({
url: `/api/trpc`,
}),
}),
],
});
鉴权按环境分四种方式(均来自同一篇文档):
- Web 应用、同域:cookie 会随请求自动发送,无需额外处理;
- Web 应用、跨域:在
eventSourceOptions里返回withCredentials: true; - 非 Web 环境 / 需要自定义 headers:用
event-source-polyfill对EventSource做 ponyfill,通过eventSourceOptions回调注入headers(如authorization); connectionParams:会作为 URL 的connectionParams查询参数序列化发送,文档明确表示因此更推荐上面其他方式。
SSE 路径的两个已知限制,文档里写得很明确:
- 运行环境不支持原生
EventSource时需要 polyfill。React Native 场景文档单独给了兼容说明:推荐基于 React Native 网络库的rn-eventsource-reborn(基于XMLHttpRequest的 polyfill 在应用切后台后无法重连),再用web-streams-polyfill补 Streams API、@azure/core-asynciterator-polyfill补AsyncIterator。 EventSource不允许对活跃的eventSourceOptions()/url()重新求值。鉴权过期后要用retryLink配合httpSubscriptionLink重建连接拿最新凭证;注意文档的警告:重启连接会从头重建EventSource,之前 tracked 的事件会丢失。
另外,服务端 SSE 选项里的 emitAndEndImmediately(数据发完立即结束请求)是为不支持流式响应的 serverless 运行时准备的,选这类部署环境时留意。
路径二:WebSockets(wsLink)
WebSockets 需要自己搭一个 WebSocket 服务器。WebSockets 文档的做法是安装 ws 并用 applyWSSHandler 挂载 tRPC 路由。文档给出的安装命令是:
yarn add ws
import { applyWSSHandler } from '@trpc/server/adapters/ws';
import { WebSocketServer } from 'ws';
import { appRouter } from './routers/app';
import { createContext } from './trpc';
const wss = new WebSocketServer({
port: 3001,
});
const handler = applyWSSHandler({
wss,
router: appRouter,
createContext,
// Enable heartbeat messages to keep connection open (disabled by default)
keepAlive: {
enabled: true,
// server ping message interval in milliseconds
pingMs: 30000,
// connection is terminated if pong message is not received in this many milliseconds
pongWaitMs: 5000,
},
});
wss.on('connection', (ws) => {
console.log(`++ Connection (${wss.clients.size})`);
ws.once('close', () => {
console.log(`-- Connection (${wss.clients.size})`);
});
});
console.log('WebSocket Server listening on ws://localhost:3001');
process.on('SIGTERM', () => {
console.log('SIGTERM');
handler.broadcastReconnectNotification();
wss.close();
});
几个要点:keepAlive 默认关闭,pingMs 是服务端 ping 间隔,pongWaitMs 内收不到 pong 会断开连接;SIGTERM 时先调 handler.broadcastReconnectNotification() 广播 { id: null, type: 'reconnect' } 通知客户端重连,再关闭服务器。
客户端用 createWSClient + wsLink:
import { createTRPCClient, createWSClient, wsLink } from '@trpc/client';
import type { AppRouter } from './server';
// create persistent WebSocket connection
const wsClient = createWSClient({
url: `ws://localhost:3001`,
});
// configure TRPCClient to use WebSockets transport
const client = createTRPCClient<AppRouter>({
links: [
wsLink({
client: wsClient,
}),
],
});
文档提示:如果只想让订阅走 WebSocket、query/mutation 仍走 HTTP,用 Links 拆分即可(splitLink 的写法与 SSE 路径相同)。
鉴权:给 createWSClient 定义 connectionParams,它会在连接建立时作为第一条消息发出,服务端在 createContext 里通过 opts.info.connectionParams 读取。做 Web 应用时可以跳过这节,因为 cookie 会随请求发送。协议细节(subscription / subscription.stop 消息、data/started/stopped 响应类型)见 WebSockets 文档的 RPC Specification 一节。
验证订阅是否工作
- WebSocket 服务器侧:启动后终端打印
WebSocket Server listening on ws://localhost:3001;每个客户端连接/断开时打印++ Connection (N)/-- Connection (N),这是 文档示例中给出的日志,可用来确认握手和连接数。 - 客户端收数:调用
trpc.xxx.subscribe(undefined, { onData, onError }),onData收到服务端yield的值即代表链路通了。仓库里最简单的完整验证脚本是 standalone-server 示例的 client.ts:服务端 server.ts 定义了一个randomNumber订阅(setInterval每 200ms 发一个随机数),客户端计数到 3 次后subscription.unsubscribe()并关闭连接。注意该示例跑在 Node 环境,客户端需要给globalThis.WebSocket赋ws包的WebSocket实现。 - 断线恢复:用
tracked(id, data)发事件,断开后客户端自动重连并携带最后已收到的 ID;服务端从opts.input?.lastEventId读取补发缺失数据。 - 服务端主动停止:生成器内
return,订阅停止、客户端断开。 - 错误分支:5xx 触发客户端自动重连;其他错误传入客户端
onError()。
选一个参考项目
Subscriptions 文档列了三个参考项目,本仓库内都有对应目录:
| 传输 | 定位 | 仓库内位置 | 关键文件 |
|---|---|---|---|
| WebSockets | 最小 Node.js 示例(HTTP 与 WS 挂在同一端口 2022) | examples/standalone-server | server.ts、client.ts |
| SSE | 全栈实现(Next.js + Drizzle) | examples/next-sse-chat | trpc.ts、channel.ts、post.ts |
| WebSockets | 全栈实现(Next.js + Prisma) | examples/next-prisma-websockets-starter | wssDevServer.ts、utils/trpc.ts |
两个全栈示例的启动命令(摘自各自 README):
# examples/next-sse-chat
pnpm i
cp .env.example .env
pnpm dev
# examples/next-prisma-websockets-starter
pnpm i
pnpm dev # starts next.js + WebSocket server
# 或 pnpm dx:另会启动本地 Postgres、跑迁移并写入种子数据
值得对照看的配置:SSE 示例的 trpc.ts 给了一组生产向的 sse 参数——maxDurationMs: 5 * 60 * 1_000、ping.intervalMs: 3_000、client.reconnectAfterInactivityMs: 5_000;WS 全栈示例的 utils/trpc.ts 则展示了 SSR 用 httpBatchLink、浏览器端用 wsLink 的 link 切换写法。
限制与边界
httpSubscriptionLink依赖EventSource、Streams API 和AsyncIterator,React Native 三者都不原生支持,必须按文档列出的包做 ponyfill 后才能继续走 Setup 流程。- SSE 的
connectionParams会被序列化进 URL,不要往里放敏感 token。 retryLink重建连接时会从头创建EventSource,之前 tracked 的事件会丢失,这是文档明确的代价。- WebSocket 的
keepAlive默认禁用;文档示例参数是pingMs: 30000/pongWaitMs: 5000。 - 服务端 SSE 的
maxDurationMs默认不限时(undefined),emitAndEndImmediately只适用于不支持流式响应的 serverless 运行时。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00