首页
/ tRPC 实时订阅在 WebSockets 与 SSE 之间怎么选?官方建议与参考项目

tRPC 实时订阅在 WebSockets 与 SSE 之间怎么选?官方建议与参考项目

2026-09-09 21:08:21作者:晏闻田Solitary

tRPC 中把实时更新从服务端推给客户端用的是订阅(subscription),传输层有两条路:WebSockets 和 Server-sent Events(SSE)。二者的服务端订阅过程写法完全一致,差异集中在客户端 link 的选择、服务端额外的搭建步骤和鉴权方式上。本文基于 SubscriptionsWebSocketshttpSubscriptionLink 三篇官方文档,以及本仓库内的三个参考示例,说明如何选型,并把两条路径都跑起来。

官方建议是什么

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-polyfillEventSource 做 ponyfill,通过 eventSourceOptions 回调注入 headers(如 authorization);
  • connectionParams:会作为 URL 的 connectionParams 查询参数序列化发送,文档明确表示因此更推荐上面其他方式。

SSE 路径的两个已知限制,文档里写得很明确:

  1. 运行环境不支持原生 EventSource 时需要 polyfill。React Native 场景文档单独给了兼容说明:推荐基于 React Native 网络库的 rn-eventsource-reborn(基于 XMLHttpRequest 的 polyfill 在应用切后台后无法重连),再用 web-streams-polyfill 补 Streams API、@azure/core-asynciterator-polyfillAsyncIterator
  2. 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.WebSocketws 包的 WebSocket 实现。
  • 断线恢复:用 tracked(id, data) 发事件,断开后客户端自动重连并携带最后已收到的 ID;服务端从 opts.input?.lastEventId 读取补发缺失数据。
  • 服务端主动停止:生成器内 return,订阅停止、客户端断开。
  • 错误分支:5xx 触发客户端自动重连;其他错误传入客户端 onError()

选一个参考项目

Subscriptions 文档列了三个参考项目,本仓库内都有对应目录:

传输 定位 仓库内位置 关键文件
WebSockets 最小 Node.js 示例(HTTP 与 WS 挂在同一端口 2022) examples/standalone-server server.tsclient.ts
SSE 全栈实现(Next.js + Drizzle) examples/next-sse-chat trpc.tschannel.tspost.ts
WebSockets 全栈实现(Next.js + Prisma) examples/next-prisma-websockets-starter wssDevServer.tsutils/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_000ping.intervalMs: 3_000client.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 运行时。
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395