tRPC 订阅(Subscriptions)完全指南:SSE / WebSocket 实时事件流、tracked() 断线恢复与常见误区
2026/9/9 13:27:12 网站建设 项目流程

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 是当前与未来的唯一推荐写法

订阅的典型工作流是:

  1. 客户端发起订阅请求并维持一条持久连接;
  2. 服务端事件源(数据库、EventEmitter、消息队列、轮询)产生数据时逐条yield
  3. 客户端在onData回调中消费每条事件;
  4. 一旦断线,客户端自动重连;若使用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服务端 +applyWSSHandlerwsLink+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.enabledbooleanfalse是否由服务端周期性发送 SSE 注释型 ping 消息,用于保活、防止代理/客户端超时断连
ping.intervalMsnumber1000ping 发送间隔(毫秒)。源码中当该值为Infinity<= 0时不会真正注入 ping
client.reconnectAfterInactivityMsnumberundefined(关闭)客户端在指定毫秒内未收到任何消息(含 ping)则判定连接失效并主动重连
maxDurationMsnumberundefined单条 SSE 连接的最大存活时长,到时服务端结束流
emitAndEndImmediatelybooleanfalse发送首条数据后立即结束请求;仅用于不支持流式响应的 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:第一条消息,dataJSON.stringify(clientOptions),用于把reconnectAfterInactivityMs等参数告知客户端;
  • ping:心跳事件(data为空串);
  • return:流正常结束(对应生成器return),客户端收到后主动close()EventSource 并结束读取;
  • serialized-error:生成器内抛错时,把序列化后的错误作为事件下发。

当你yield的是tracked(id, data)信封时,服务端会把它拆成iddata两个字段写入帧;普通值则只写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查询参数中,服务端可在createContextopts.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发回;
  • 同一文件也导出了被标记@deprecatedsse(event)旧辅助函数——它内部就是调用tracked(event.id, event.data),新代码直接使用tracked即可。

使用tracked后,断线恢复的完整闭环是:

  1. 客户端首次订阅时传入初始lastEventId(如null,或一个已知位置);
  2. 服务端每条事件yield tracked(post.id, post)
  3. 客户端持续消费并更新本地"最后收到的 ID";
  4. 网络断开 → 自动重连 → 把最后的 ID 作为lastEventId传给服务端;
  5. 服务端在.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

服务端需要一个wsWebSocketServer,配合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建立连接后的第一条消息即连接参数,服务端在createContextopts.info.connectionParams读取
WebSocket原生WS 实现 ponyfill
retryDelayMs指数退避(exponentialBackoff)自定义重连延迟策略函数
lazyenabled: false惰性模式:空闲closeMs毫秒后自动断开 WS
keepAliveenabled: falseintervalMs默认5_000pongTimeoutMs默认1_000客户端主动 ping、无 pong 则断开
experimental_encoderjsonEncoder自定义线上编码(如二进制格式)

tracked()在 WebSocket 通道同样生效:wsLink会在重连时自动发送客户端最后已知的 ID 作为lastEventId输入(websockets.md 的 "Automatic tracking of id" 一节与 subscriptions.md 均如此说明)。

双向连接如何做认证

由于 WebSocket 连接建立时浏览器不会自动附带 Cookie 之外的 header,认证参数通过createWSClientconnectionParams以首条消息传递;服务端侧从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 订阅在生产环境稳如磐石,记住这条主线即可:

  1. 默认走 SSE:服务端initTRPC.create({ sse: {...} })配好pingclient.reconnectAfterInactivityMs,用async function*.subscription()
  2. 一切事件都yield tracked(id, data),把断线恢复交给lastEventId闭环;
  3. 先监听、后补历史,避免补数据窗口吞事件;
  4. 需要清理副作用用try...finally,需要退出用return
  5. 只有当双向通信成为硬需求时,才引入ws/applyWSSHandler+wsLink,并记得SIGTERMbroadcastReconnectNotification()

仓库内可继续深挖的资料:

  • 官方订阅主文档: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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询