完全指南:SSE / WebSocket 实时事件流、tracked() 断线恢复与常见误区)
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()、用httpSubscriptionLinkSSE或wsLinkWebSocket在客户端消费、用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 AsyncIterableany, 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知道客户端最后收到哪条从而只补发遗漏的事件。官方维护者总结的通用判断原则是订阅优先选择 SSEServer-sent Events因为它搭建简单、不需要独立的 WebSocket 服务器只有在需要双向通信时才引入 WebSocket。先做选型SSE 还是 WebSockettRPC 官方文档www/docs/server/subscriptions.md给出了两条通道的入口通道服务端承载客户端 Link推荐场景SSE推荐普通 HTTP 适配器即可如 standalone / Node 的createHTTPServer、Next.js 等httpSubscriptionLink绝大多数单向服务端推送给客户端的订阅WebSocket需要ws服务端 applyWSSHandlerwsLinkcreateWSClient需要双向通信客户端也可向服务端发消息或需要 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.signalAbortSignal后请求被中止时事件监听会自动取消这是订阅清理的基石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时不会真正注入 pingclient.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第一条消息data为JSON.stringify(clientOptions)用于把reconnectAfterInactivityMs等参数告知客户端ping心跳事件data为空串return流正常结束对应生成器return客户端收到后主动close()EventSource 并结束读取serialized-error生成器内抛错时把序列化后的错误作为事件下发。当你yield的是tracked(id, data)信封时服务端会把它拆成id与data两个字段写入帧普通值则只写data字段。客户端SSE用splitLinkhttpSubscriptionLink消费订阅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 createTRPCClientAppRouter({ 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 \| Promisestring连接地址可传函数以在重连前动态计算最新 URLconnectionParams以对象或函数形式给出序列化到 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.tsexport type TrackedEnvelopeTData [TrackedId, TData, typeof trackedSymbol]; export function trackedTData(id: string, data: TData): TrackedEnvelopeTData { 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] trackedSymbolSSE 服务端遇到 tracked 信封时会写出id: id帧行从而触发 EventSource 客户端记录Last-Event-IDWebSocket 通道则由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 Patternsimport 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.mdconst 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 isAsyncIterableTValue(value: unknown): value is AsyncIterableTValue { 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让客户端主动重连到新实例实现零感知滚动重启。客户端createWSClientwsLink// client import { createTRPCClient, createWSClient, httpBatchLink, splitLink, wsLink, } from trpc/client; import type { AppRouter } from ./server; const wsClient createWSClient({ url: ws://localhost:3001 }); const trpc createTRPCClientAppRouter({ links: [ splitLink({ condition: (op) op.type subscription, true: wsLink({ client: wsClient }), false: httpBatchLink({ url: http://localhost:3000 }), }), ], });createWSClient的完整可配置项类型见 wsLink.md值得留意选项默认说明url必填WS 地址也支持() string \| PromisestringconnectionParams—建立连接后的第一条消息即连接参数服务端在createContext的opts.info.connectionParams读取WebSocket原生WS 实现 ponyfillretryDelayMs指数退避exponentialBackoff自定义重连延迟策略函数lazyenabled: false惰性模式空闲closeMs毫秒后自动断开 WSkeepAliveenabled: 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认证参数通过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 进程同时承载 HTTPquery/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 3SSE 需要自定义请求头却用了原生 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 3SSE 已够用却选择 WebSocketSSEhttpSubscriptionLink被官方维护者推荐用于大多数订阅场景WebSocket 引入连接管理、重连、心跳与独立服务器进程等额外复杂度。只有在需要双向通信或 WebSocket 专属能力时才使用wsLink。MEDIUM 4WebSocket 重连时输入参数陈旧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/applyWSSHandlerwsLink并记得SIGTERM时broadcastReconnectNotification()。仓库内可继续深挖的资料官方订阅主文档www/docs/server/subscriptions.md含 stopping、error handling、output validation 全量示例SSE 客户端链接文档www/docs/client/links/httpSubscriptionLink.mdWebSocket 服务端/协议文档www/docs/server/websockets.mdWebSocket 客户端链接文档www/docs/client/links/wsLink.mdSSE 传输层源码packages/server/src/unstable-core-do-not-import/stream/sse.tstracked/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.tsHTTP 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),仅供参考