tRPC 订阅(Subscriptions)完全指南:SSE / WebSocket 实时事件流、tracked() 断线恢复与常见误区
【免费下载链接】trpc🧙♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc
订阅是 tRPC 中一类"实时事件流"能力:由客户端与服务端建立并维持持久连接,服务端可随时把数据推送给客户端;连接中断后,配合tracked()事件 ID 机制,客户端会自动重连并基于lastEventId优雅地补齐遗漏事件。本篇将以 tRPC v11(本仓库库版本 11.16.0)为背景,完整讲解在服务端用 async generator 定义.subscription()、用httpSubscriptionLink(SSE)或wsLink(WebSocket)在客户端消费、用initTRPC.create()配置 SSE 心跳与客户端超时、以及用tracked(id, data)实现断线恢复的完整姿势,并逐个拆解官方维护者总结的高/中危常见误区。读完你可以直接照搬文中的可运行代码,搭建出具备自动重连能力的实时推送后端与前端。
订阅在 tRPC 中的定位与核心概念
tRPC 的订阅(subscription)是一类 procedure,它与query/mutation并列,但解析器返回的不是普通值,而是一个异步可迭代对象(AsyncIterable)。在服务端,推荐用async generator 语法(async function*)来定义:
t.procedure.subscription(async function* (opts) { // ...持续 yield 事件数据 });从类型层面可以验证这一点:在 procedureBuilder.ts 中,.subscription()的非废弃重载要求解析器产出$Output extends AsyncIterable<any, void, any>;而基于observable的重载带有@deprecated注释,明确写着"Using subscriptions with an observable is deprecated. Use an async generator instead. This feature will be removed in v12 of tRPC."。也就是说,async generator 是当前与未来的唯一推荐写法。
订阅的典型工作流是:
- 客户端发起订阅请求并维持一条持久连接;
- 服务端事件源(数据库、
EventEmitter、消息队列、轮询)产生数据时逐条yield; - 客户端在
onData回调中消费每条事件; - 一旦断线,客户端自动重连;若使用
tracked()携带事件 ID,服务端可以通过输入参数lastEventId知道"客户端最后收到哪条",从而只补发遗漏的事件。
官方维护者总结的通用判断原则是:订阅优先选择 SSE(Server-sent Events),因为它搭建简单、不需要独立的 WebSocket 服务器,只有在需要双向通信时才引入 WebSocket。
先做选型:SSE 还是 WebSocket?
tRPC 官方文档(www/docs/server/subscriptions.md)给出了两条通道的入口:
| 通道 | 服务端承载 | 客户端 Link | 推荐场景 |
|---|---|---|---|
| SSE(推荐) | 普通 HTTP 适配器即可(如 standalone / Node 的createHTTPServer、Next.js 等) | httpSubscriptionLink | 绝大多数单向"服务端推送给客户端"的订阅 |
| WebSocket | 需要ws服务端 +applyWSSHandler | wsLink+createWSClient | 需要双向通信(客户端也可向服务端发消息),或需要 WebSocket 专属能力 |
SSE 之所以被推荐,是因为它本质上是基于 HTTP 的流式响应,无需管理独立的 WebSocket 连接池、心跳、重连握手,也不需要额外跑一个 WS server 进程。接下来先按"SSE 优先"展开完整落地流程。
服务端:在initTRPC.create()中配置 SSE 并定义订阅 procedure
订阅的 SSE 参数在初始化 tRPC 实例时就一次性配置好。先看一个完整、可直接运行的 SSE 服务端示例(出自 SKILL.md):
// server.ts import EventEmitter, { on } from 'node:events'; import { initTRPC, tracked } from '@trpc/server'; import { createHTTPServer } from '@trpc/server/adapters/standalone'; import { z } from 'zod'; const t = initTRPC.create({ sse: { ping: { enabled: true, intervalMs: 2000, }, client: { reconnectAfterInactivityMs: 5000, }, }, }); type Post = { id: string; title: string }; const ee = new EventEmitter(); const appRouter = t.router({ onPostAdd: t.procedure .input(z.object({ lastEventId: z.string().nullish() }).optional()) .subscription(async function* (opts) { for await (const [data] of on(ee, 'add', { signal: opts.signal })) { const post = data as Post; yield tracked(post.id, post); } }), }); export type AppRouter = typeof appRouter; createHTTPServer({ router: appRouter, createContext() { return {}; }, }).listen(3000);要点拆解:
on(ee, 'add', { signal: opts.signal })是node:events提供的"把 EventEmitter 变成异步迭代器"的工具,监听名add的事件会逐个流经for await;- 传入
opts.signal(AbortSignal)后,请求被中止时事件监听会自动取消,这是订阅清理的基石; yield tracked(post.id, post)让每条事件携带可恢复的 ID,客户端据此自动重连;.input()定义了lastEventId输入(首次连接由客户端传入初始值,重连时自动替换为客户端最后收到的 ID)。
sse配置项的完整含义与默认值
在 sse.ts 中,服务端通过SSEStreamProducerOptions定义这些选项(该类型同时被 rootConfig.ts 用于约束initTRPC.create()的sse字段):
| 配置项 | 类型 | 默认值 | 作用 |
|---|---|---|---|
ping.enabled | boolean | false | 是否由服务端周期性发送 SSE 注释型 ping 消息,用于保活、防止代理/客户端超时断连 |
ping.intervalMs | number | 1000 | ping 发送间隔(毫秒)。源码中当该值为Infinity或<= 0时不会真正注入 ping |
client.reconnectAfterInactivityMs | number | undefined(关闭) | 客户端在指定毫秒内未收到任何消息(含 ping)则判定连接失效并主动重连 |
maxDurationMs | number | undefined | 单条 SSE 连接的最大存活时长,到时服务端结束流 |
emitAndEndImmediately | boolean | false | 发送首条数据后立即结束请求;仅用于不支持流式响应的 serverless 运行时 |
client相关 | — | {} | 会被序列化进 SSE 流的第一条消息(connected事件)下发给客户端,客户端据此设置自身行为 |
一个值得注意的工程细节:在 sse.ts 中,如果同时开启 ping 并配置了client.reconnectAfterInactivityMs,且ping.intervalMs > client.reconnectAfterInactivityMs,服务端会直接抛错:
Ping interval must be less than client reconnect interval to prevent unnecessary reconnection这从代码层面锁定了"心跳间隔必须小于客户端失活重连阈值"的规则(下一节"常见误区"还会展开边界情况)。
SSE 流的传输细节(源码佐证)
同一文件的sseHeaders常量展示了 SSE 响应应有的响应头:
export const sseHeaders = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', 'X-Accel-Buffering': 'no', Connection: 'keep-alive', } as const;而流的实际编码在 TransformStream 中逐字段拼装为event:/data:/id:/ 注释行。内部还定义了若干协议级事件名:
connected:第一条消息,data为JSON.stringify(clientOptions),用于把reconnectAfterInactivityMs等参数告知客户端;ping:心跳事件(data为空串);return:流正常结束(对应生成器return),客户端收到后主动close()EventSource 并结束读取;serialized-error:生成器内抛错时,把序列化后的错误作为事件下发。
当你yield的是tracked(id, data)信封时,服务端会把它拆成id与data两个字段写入帧;普通值则只写data字段。
客户端(SSE):用splitLink+httpSubscriptionLink消费订阅
SSE 客户端方案是"用httpSubscriptionLink处理订阅、用普通 HTTP link 处理 query/mutation"的路由组合(来自 SKILL.md 与 httpSubscriptionLink.md):
// client.ts import { createTRPCClient, httpBatchLink, httpSubscriptionLink, splitLink, } from '@trpc/client'; import type { AppRouter } from './server'; const trpc = createTRPCClient<AppRouter>({ links: [ splitLink({ condition: (op) => op.type === 'subscription', true: httpSubscriptionLink({ url: 'http://localhost:3000' }), false: httpBatchLink({ url: 'http://localhost:3000' }), }), ], }); const subscription = trpc.onPostAdd.subscribe( { lastEventId: null }, { onData(post) { console.log('New post:', post); }, onError(err) { console.error('Subscription error:', err); }, }, ); // To stop: // subscription.unsubscribe();几个关键点:
.subscribe(input, { onData, onError })的第一个参数即 procedure 的输入,lastEventId: null表示"首次连接、不需要历史";此字段会随客户端收到的带id事件被持续更新;- 返回的
subscription对象调用.unsubscribe()即可停止订阅; condition: (op) => op.type === 'subscription'是路由依据,splitLink据此把三类 operation 分发到不同终止 link。
httpSubscriptionLink在内部使用浏览器原生EventSourceAPI 建立长连接,这带来一个直接好处:EventSource 规范天然内置自动重连,一旦网络抖动或收到非 2xx 响应,客户端会自动重试;而重连时,若最后收到的事件携带id字段,浏览器会自动通过Last-Event-ID请求头/参数把该 ID 带给服务端,配合服务端输入里的lastEventId即可补齐漏掉的事件。
客户端/服务端两侧的可用配置一览
httpSubscriptionLink的选项(类型定义见 httpSubscriptionLink.md):
| 选项 | 说明 |
|---|---|
url: string \| () => string \| Promise<string> | 连接地址,可传函数以在重连前动态计算最新 URL |
connectionParams | 以对象或函数形式给出,序列化到 URL 的connectionParams查询参数中,服务端可在createContext的opts.info.connectionParams里读取 |
transformer | 数据转换器(如superjson),须与服务端一致 |
EventSource | 传入 EventSource ponyfill/polyfill(见下文自定义请求头) |
eventSourceOptions | 传给 EventSource 构造器的选项或返回这些选项的异步回调(回调可拿到当前op,用于按操作生成签名/新 token) |
服务端侧则由initTRPC.create({ sse: {...} })统一配置(上一节的表格),其中两个最重要的联调参数是:
sse.ping:服务端心跳;sse.client.reconnectAfterInactivityMs:客户端失活重连阈值——它在服务端配置,但会通过connected首帧下发给客户端生效(见 httpSubscriptionLink.md 的"Timeout Configuration"一节)。
联调建议:服务端ping.intervalMs取 2s、客户端reconnectAfterInactivityMs取 5s,这样客户端有充足余量等到下一次心跳,不会误判连接死亡。
tracked(id, data)深潜:断线恢复的核心机制
tracked是本仓库从@trpc/server导出的辅助函数,实现位于 tracked.ts:
export type TrackedEnvelope<TData> = [TrackedId, TData, typeof trackedSymbol]; export function tracked<TData>(id: string, data: TData): TrackedEnvelope<TData> { if (id === '') { throw new Error( '`id` must not be an empty string as empty string is the same as not setting the id at all', ); } return [id as TrackedId, data, trackedSymbol]; }实现要点:
tracked(id, data)返回一个三元组:[id, data, 私有符号],这个私有符号让运行时可判别该值是否"被追踪过",对应的判断函数是isTrackedEnvelope(value)(Array.isArray(value) && value[2] === trackedSymbol);- SSE 服务端遇到 tracked 信封时会写出
id: <id>帧行,从而触发 EventSource 客户端记录Last-Event-ID;WebSocket 通道则由wsLink在重连时自动把最后已知 ID 作为lastEventId发回; - 同一文件也导出了被标记
@deprecated的sse(event)旧辅助函数——它内部就是调用tracked(event.id, event.data),新代码直接使用tracked即可。
使用tracked后,断线恢复的完整闭环是:
- 客户端首次订阅时传入初始
lastEventId(如null,或一个已知位置); - 服务端每条事件
yield tracked(post.id, post); - 客户端持续消费并更新本地"最后收到的 ID";
- 网络断开 → 自动重连 → 把最后的 ID 作为
lastEventId传给服务端; - 服务端在
.input中读到opts.input?.lastEventId,从自己的存储(DB、日志、消息队列)查询并补发此 ID 之后的所有事件,再衔接实时流。
从lastEventId恢复 + 防止事件丢失的顺序
重连恢复的完整服务端模式(来自 SKILL.md 的 Core Patterns):
import EventEmitter, { on } from 'node:events'; import { initTRPC, tracked } from '@trpc/server'; import { z } from 'zod'; const t = initTRPC.create(); const ee = new EventEmitter(); const appRouter = t.router({ onPostAdd: t.procedure .input(z.object({ lastEventId: z.string().nullish() }).optional()) .subscription(async function* (opts) { const iterable = on(ee, 'add', { signal: opts.signal }); if (opts.input?.lastEventId) { // Fetch and yield events since lastEventId from your database // const missed = await db.post.findMany({ where: { id: { gt: opts.input.lastEventId } } }); // for (const post of missed) { yield tracked(post.id, post); } } for await (const [data] of iterable) { yield tracked(data.id, data); } }), });顺序至关重要:必须先on(ee, 'add', ...)建立事件监听(先iterable),再去数据库补拉lastEventId之后的历史。如果反着来——先await db.getEvents()再挂监听——补拉历史期间新产生的事件就会白白丢失(详见"常见误区"中的 HIGH 项)。
轮询式订阅:Pull 数据库新数据并下推
当事件源不支持推送、需要定时去数据库捞"上次游标之后的新数据"时,可采用官方"Pull data in a loop"配方(见 SKILL.md 与 subscriptions.md)。此时lastEventId直接使用时间游标:
import { initTRPC, tracked } from '@trpc/server'; import { z } from 'zod'; const t = initTRPC.create(); const appRouter = t.router({ onNewItems: t.procedure .input(z.object({ lastEventId: z.coerce.date().nullish() })) .subscription(async function* (opts) { let cursor = opts.input?.lastEventId ?? null; while (!opts.signal?.aborted) { const items = await db.item.findMany({ where: cursor ? { createdAt: { gt: cursor } } : undefined, orderBy: { createdAt: 'asc' }, }); for (const item of items) { yield tracked(item.createdAt.toJSON(), item); cursor = item.createdAt; } await new Promise((r) => setTimeout(r, 1000)); } }), });配方要点:
- 用
z.coerce.date()把客户端传来的时间字符串输入强转为Date,天然充当游标; while (!opts.signal?.aborted)是循环的退出条件:客户端一旦断开,opts.signal被 abort,循环即停;- 每次用
cursor过滤createdAt > cursor并升序取出,逐条yield tracked(...)后推进游标; - 末尾
await sleep(1000)防止打爆数据库,实现约 1s 一次的轻量轮询。
停止订阅与副作用清理
服务端主动结束:直接return
若需在服务端终止某条订阅,在 generator 内return即可。以官方文档示例的逻辑为例:一旦计数超过阈值就结束流,客户端会随之断开:
.subscription(async function* (opts) { let index = opts.input?.lastEventId ?? 0; while (!opts.signal!.aborted) { const idx = index++; if (idx > 100) { // With this, the subscription will stop and the client will disconnect return; } await new Promise((resolve) => setTimeout(resolve, 10)); } });客户端停止订阅则调用订阅对象上的.unsubscribe()。
清理副作用:try...finally
订阅可能持有定时器、文件句柄、外部监听器等副作用。官方文档明确说明:订阅因任何原因停止时,tRPC 都会调用 generator 实例的.return(),因此try...finally是可靠的清理时机(SKILL.md):
const appRouter = t.router({ events: t.procedure.subscription(async function* (opts) { const cleanup = registerListener(); try { for await (const [data] of on(ee, 'event', { signal: opts.signal })) { yield data; } } finally { cleanup(); } }), });同时别忘了opts.signal本身:把它传给on(..., { signal }),请求中止时 EventEmitter 的异步迭代会自动停止,多数基于监听器的清理其实已经由它兜底。
错误处理与订阅输出校验
错误语义
- 服务端 generator 内抛错,会传播到 tRPC 后端的
onError(); - 若抛出的错误属于 5xx:客户端会基于
tracked()记录的最后事件 ID 自动重连(因为这是一个可能"补齐后可恢复"的瞬态错误); - 其它错误:订阅被取消,错误进入客户端的
onError()回调(见 subscriptions.md 的 Error handling 一节)。
订阅输出校验:必须遍历 async iterable
订阅是 async iterable,.output()无法像 query/mutation 那样直接校验单值,需要深入迭代器逐条校验。仓库里有一套现成的 Zod v4 辅助工具zAsyncIterable(完整实现见 zAsyncIterable.ts,配套文档在 subscriptions.md)。其核心思路是:
function isAsyncIterable<TValue>(value: unknown): value is AsyncIterable<TValue> { return !!value && typeof value === 'object' && Symbol.asyncIterator in value; }先断言值是 async iterable,再用.transform(async function* (iter) {...})逐条解析;当启用tracked: true时,每一条都会先用trackedEnvelopeSchema拆出[id, data],分别校验数据后用tracked(id, 校验后的数据)重新组装。服务端使用示例:
export const appRouter = t.router({ mySubscription: t.procedure .input(z.object({ lastEventId: z.coerce.number().min(0).optional() })) .output( zAsyncIterable({ yield: z.object({ count: z.number() }), tracked: true, }), ) .subscription(async function* (opts) { let index = opts.input?.lastEventId ?? 0; while (true) { index++; yield tracked(String(index), { count: index }); await new Promise((resolve) => setTimeout(resolve, 1000)); } }), });WebSocket 路线:仅在需要双向通信时使用
wsLink同样是 tRPC 官方支持且完整的订阅通道。若你确定需要双向通信或 WebSocket 专属特性,可以照下面的配置落地。
服务端:applyWSSHandler
服务端需要一个ws的WebSocketServer,配合applyWSSHandler把 tRPC 路由桥接上去(来自 SKILL.md,更完整的服务器示例见 standalone-server/src/server.ts,同一 HTTP server 同时挂载 HTTP 与 WS):
// server import { applyWSSHandler } from '@trpc/server/adapters/ws'; import { WebSocketServer } from 'ws'; import { appRouter } from './router'; const wss = new WebSocketServer({ port: 3001 }); const handler = applyWSSHandler({ wss, router: appRouter, createContext() { return {}; }, keepAlive: { enabled: true, pingMs: 30000, pongWaitMs: 5000, }, }); process.on('SIGTERM', () => { handler.broadcastReconnectNotification(); wss.close(); });配置与运维要点:
keepAlive默认关闭;开启后pingMs是服务端 ping 间隔,pongWaitMs是"未收到 pong 即判定死亡并断开"的阈值(本例每 30s ping、5s 内无 pong 断开);- 收到
SIGTERM优雅下线前调用handler.broadcastReconnectNotification():它会向所有客户端广播{ id: null, type: 'reconnect' }通知(见 websockets.md 的 "Notifications from Server to Client"),让客户端主动重连到新实例,实现零感知滚动重启。
客户端:createWSClient+wsLink
// client import { createTRPCClient, createWSClient, httpBatchLink, splitLink, wsLink, } from '@trpc/client'; import type { AppRouter } from './server'; const wsClient = createWSClient({ url: 'ws://localhost:3001' }); const trpc = createTRPCClient<AppRouter>({ links: [ splitLink({ condition: (op) => op.type === 'subscription', true: wsLink({ client: wsClient }), false: httpBatchLink({ url: 'http://localhost:3000' }), }), ], });createWSClient的完整可配置项(类型见 wsLink.md)值得留意:
| 选项 | 默认 | 说明 |
|---|---|---|
url | 必填 | WS 地址,也支持() => string \| Promise<string> |
connectionParams | — | 建立连接后的第一条消息即连接参数,服务端在createContext的opts.info.connectionParams读取 |
WebSocket | 原生 | WS 实现 ponyfill |
retryDelayMs | 指数退避(exponentialBackoff) | 自定义重连延迟策略函数 |
lazy | enabled: false | 惰性模式:空闲closeMs毫秒后自动断开 WS |
keepAlive | enabled: false,intervalMs默认5_000,pongTimeoutMs默认1_000 | 客户端主动 ping、无 pong 则断开 |
experimental_encoder | jsonEncoder | 自定义线上编码(如二进制格式) |
tracked()在 WebSocket 通道同样生效:wsLink会在重连时自动发送客户端最后已知的 ID 作为lastEventId输入(websockets.md 的 "Automatic tracking of id" 一节与 subscriptions.md 均如此说明)。
双向连接如何做认证
由于 WebSocket 连接建立时浏览器不会自动附带 Cookie 之外的 header,认证参数通过createWSClient的connectionParams以首条消息传递;服务端侧从opts.info.connectionParams取出(如token)完成鉴权:
const wsClient = createWSClient({ url: 'ws://localhost:3000', connectionParams: async () => ({ token: 'supersecret' }), });作为对照,SSE 在浏览器同域场景下 Cookie 随请求自动携带;跨域则可用eventSourceOptions里的withCredentials: true。原生EventSource不支持自定义请求头,需要自定义 Header(如Authorization: Bearer ...)时必须传入基于 ponyfill 的 EventSource(详见下节)。
仓库内的可参考实现
仓库中与之配套的可运行参考包括:
- examples/next-sse-chat:全栈 SSE 订阅示例(Next.js App Router + Drizzle + docker-compose 起库),服务端含大量
tracked/lastEventId恢复逻辑; - examples/next-prisma-websockets-starter:全栈 WebSocket 订阅示例;
- examples/standalone-server:单个 Node 进程同时承载 HTTP(query/mutation)与 WebSocket(订阅)的最小服务器。
常见误区清单(官方维护者视角,按严重度排序)
下面是 SKILL.md 中沉淀的常见错误。每一条都给出错误/正确写法、原因与源码出处,可直接作为代码评审的 Checklist。
HIGH 1:用observable而不是 async generator
// Wrong import { observable } from '@trpc/server/observable'; t.procedure.subscription(({ input }) => { return observable((emit) => { emit.next(data); }); });// Correct t.procedure.subscription(async function* ({ input, signal }) { for await (const [data] of on(ee, 'event', { signal })) { yield data; } });Observable 形式的订阅已被废弃,将在 v12 移除(procedureBuilder.ts 中有明确的@deprecated注释)。顺带一提,仓库内的旧示例(如 examples/standalone-server/src/server.ts 的randomNumber)仍用 observable 演示,但那属于历史写法,新代码应一律使用 async generator。
HIGH 2:先拉历史、后挂事件监听,导致事件丢失
// Wrong t.procedure.subscription(async function* (opts) { const history = await db.getEvents(); // events may fire here and be lost yield* history; for await (const event of listener) { yield event; } });// Correct t.procedure.subscription(async function* (opts) { const iterable = on(ee, 'event', { signal: opts.signal }); // listen first const history = await db.getEvents(); for (const item of history) { yield tracked(item.id, item); } for await (const [event] of iterable) { yield tracked(event.id, event); } });若在设置监听器之前异步抓取历史数据,抓取与监听建立之间的窗口期产生的事件会永久丢失(出处:www/docs/server/subscriptions.md)。
HIGH 3:SSE 需要自定义请求头却用了原生 EventSource
// Wrong httpSubscriptionLink({ url: 'http://localhost:3000', // Native EventSource does not support custom headers });// Correct import { EventSourcePolyfill } from 'event-source-polyfill'; httpSubscriptionLink({ url: 'http://localhost:3000', EventSource: EventSourcePolyfill, eventSourceOptions: async () => ({ headers: { authorization: 'Bearer token' }, }), });原生EventSourceAPI 不支持自定义 header,必须把基于 polyfill 的实现通过EventSource选项注入httpSubscriptionLink(出处:httpSubscriptionLink.md)。鉴权话题的更多细节(connectionParams、Cookie、EventSource polyfill headers)可进一步参考仓库内 auth 相关技能说明。
MEDIUM 1:把空字符串当作 tracked 事件 ID
// Wrong yield tracked('', data);// Correct yield tracked(event.id.toString(), data);tracked()会在 ID 为空字符串时抛出异常,因为空字符串在 SSE 语义中等同于"未设置 id"(见 tracked.ts 中的抛错分支)。注意 Event ID 是字符串,数字 ID 记得.toString()。
MEDIUM 2:服务端 ping 间隔 ≥ 客户端重连阈值
// Wrong initTRPC.create({ sse: { ping: { enabled: true, intervalMs: 10000 }, client: { reconnectAfterInactivityMs: 5000 }, }, });// Correct initTRPC.create({ sse: { ping: { enabled: true, intervalMs: 2000 }, client: { reconnectAfterInactivityMs: 5000 }, }, });若服务端心跳间隔大于等于客户端失活重连阈值,客户端会在收到下一次 ping 之前就误判"连接死亡"而反复重连。sse.ts 会对"间隔严格大于阈值"的配置直接抛错兜底;即便恰好相等也建议避免(出处:sse.ts)。
MEDIUM 3:SSE 已够用却选择 WebSocket
SSE(httpSubscriptionLink)被官方维护者推荐用于大多数订阅场景;WebSocket 引入连接管理、重连、心跳与独立服务器进程等额外复杂度。只有在需要双向通信或 WebSocket 专属能力时才使用wsLink。
MEDIUM 4:WebSocket 重连时输入参数陈旧
WebSocket 重连后,各订阅会重新发送创建时的原始输入参数,框架没有"重连前重新求值输入"的钩子,可能导致客户端拿到陈旧数据(关联 issue 见 SKILL.md 引用)。缓解手段仍是使用tracked()+lastEventId:让服务端以"最后收到的 ID"为准补发数据,而不是依赖重连时的静态输入。
小结与进一步阅读
要让 tRPC 订阅在生产环境稳如磐石,记住这条主线即可:
- 默认走 SSE:服务端
initTRPC.create({ sse: {...} })配好ping与client.reconnectAfterInactivityMs,用async function*写.subscription(); - 一切事件都
yield tracked(id, data),把断线恢复交给lastEventId闭环; - 先监听、后补历史,避免补数据窗口吞事件;
- 需要清理副作用用
try...finally,需要退出用return; - 只有当双向通信成为硬需求时,才引入
ws/applyWSSHandler+wsLink,并记得SIGTERM时broadcastReconnectNotification()。
仓库内可继续深挖的资料:
- 官方订阅主文档:www/docs/server/subscriptions.md(含 stopping、error handling、output validation 全量示例)
- SSE 客户端链接文档:www/docs/client/links/httpSubscriptionLink.md
- WebSocket 服务端/协议文档:www/docs/server/websockets.md
- WebSocket 客户端链接文档:www/docs/client/links/wsLink.md
- SSE 传输层源码:packages/server/src/unstable-core-do-not-import/stream/sse.ts
tracked/isTrackedEnvelope实现:packages/server/src/unstable-core-do-not-import/stream/tracked.ts- 订阅 procedure 类型约束:packages/server/src/unstable-core-do-not-import/procedureBuilder.ts
- 订阅输出校验辅助工具:packages/tests/server/zAsyncIterable.ts
- HTTP + WebSocket 双通道最小示例:examples/standalone-server/src/server.ts
- 全栈 SSE 示例:examples/next-sse-chat、全栈 WebSocket 示例:examples/next-prisma-websockets-starter
【免费下载链接】trpc🧙♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考