tRPC 服务端 Subscriptions 实战指南:事件流基础、tracked() 断线恢复与输出校验
tRPC 服务端 Subscriptions 实战指南事件流基础、tracked() 断线恢复与输出校验【免费下载链接】trpc♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc导读本文是 tRPC 服务端实时订阅Subscription机制的完整技术指南。Subscription 是客户端与服务器之间的一种实时事件流通道适用于需要把服务器端事件主动推送给客户端的场景如聊天消息、股票行情、后台任务进度。读完本文你将掌握如何基于异步生成器Async Generator编写订阅 procedure、如何选择 WebSockets 与 SSE 两种传输方案、如何使用tracked()实现断线自动重连与事件恢复、如何优雅停止订阅与清理副作用以及如何为订阅输出做运行时校验。本文主体内容整理自 subscriptions.md并结合本仓库的 tracked() 实现源码、SSE 编解码实现 以及 next-sse-chat 全栈示例 进行源码级印证与扩展。Subscriptions 是什么以及何时使用Subscription 是 tRPC 中「查询 / 变更 / 订阅」三类 procedure 之一。它与 query、mutation 的关键差异在于它是一条持续的实时事件流——客户端建立并维护一个持久连接服务器可以在任意时刻把事件推送给客户端当连接意外断开时客户端还会自动尝试重连并借助tracked()事件的携带 ID 优雅地从断点恢复。在服务端实现层面一个订阅 procedure 的 resolver 是一个异步生成器函数async generator它会根据过程实现产出一个AsyncIterable。这一点在 procedureBuilder.ts 的subscription()方法上有明确体现subscription$Output extends AsyncIterableany, void, any(...)。同时该文件中还保留了使用 observable 的旧式订阅定义并在 JSDoc 中标注为deprecated明确建议“请使用异步生成器async generator替代 observable”。WebSockets 还是 Server-sent EventstRPC 提供两种传输方案搭建实时订阅WebSockets需要单独启动一个 WebSocket 服务器双向全双工通信。服务端接入细节与 RPC 消息规范见 WebSockets 页面客户端连接配置见 wsLink 文档。Server-sent EventsSSE单向服务器推送基于标准 HTTP 长连接配置更轻。客户端通过httpSubscriptionLink使用见 httpSubscriptionLink 文档。官方推荐若你不确定选哪一种优先使用 SSE 做订阅——因为它搭建更简单不需要单独维护一个 WebSocket 服务器。从本仓库的 sse.ts 实现可以看到SSE 输出遵循 HTML Server-Sent Events 规范通过text/event-stream内容类型逐帧下发event:、data:、id:字段并附带Cache-Control: no-cache, no-transform、X-Accel-Buffering: no等关键响应头。补充说明tRPC 也允许用 splitLink 之类的客户端链接把 query/mutation 走 HTTP、把订阅单独走 WebSockets做到按需分流。参考项目速查类型示例形态仓库位置WebSockets极简 Node.js WebSockets 示例examples/standalone-serverSSE全栈 SSE 聊天实现含 drizzle 数据库、断线续传examples/next-sse-chatWebSockets全栈 WebSockets 实现含 Prisma 数据层examples/next-prisma-websockets-starter基础示例事件驱动的订阅 procedure最基本的订阅模式是「事件驱动」服务器内部持有一个事件源例如 Node.js 的EventEmitter当业务事件发生时向所有订阅者广播。下面的示例定义了一个onPostAdd订阅一旦有新的Post写入就会推送给客户端// target: esnext // types: node 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; } }), });这里有两个非常关键且容易踩坑的点值得展开opts.signal必须透传给事件监听。on(ee, add, { signal })来自 Node.jsevents模块的异步迭代支持把请求携带的AbortSignal传进去之后一旦订阅被终止客户端断开、连接被 abort事件迭代器会自动停止避免内存泄漏。yield post直接产出裸数据。这种写法最简单但没有任何「断线续传」能力——客户端一旦断开重连后只能收到后续新事件中间丢失的事件无法补齐。要解决这个问题就需要引入下一节的tracked()。全栈版的事件驱动实现可以参考 examples/next-sse-chat/src/server/routers/post.ts它先用mutation写入数据库并通过ee.emit(add, channelId, post)广播再用订阅 procedure 把新帖子实时推给在线用户。用 tracked() 自动追踪事件 ID推荐{#tracked}为什么需要事件 ID如果订阅过程中连接意外断开客户端会自动重连但服务器并不知道客户端「已经收到哪里了」。tracked()的作用就是为每个产出的事件附加一个id客户端收到数据时会自动记录最后一条事件的 ID重连时把该 ID 作为lastEventId重新发起订阅服务端据此补发断点之后的所有事件。在代码里tracked(id, data)产出的每个事件形如yield tracked(post.id, post);它的底层实现定义在 tracked.tsconst trackedSymbol Symbol(); type TrackedId string { __brand: TrackedId }; export type TrackedEnvelopeTData [TrackedId, TData, typeof trackedSymbol]; export function isTrackedEnvelopeTData(value: unknown): value is TrackedEnvelopeTData { return Array.isArray(value) value[2] 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, trackedSymbol]第三个位置的 Symbol 是内部约定标记用于让运行层区分「这条消息携带 ID」还是「裸数据」TrackedEnvelopeTData与isTrackedEnvelope()会随trpc/server一起导出供你在类型与运行时两侧识别这种信封结构。id不能是空字符串。因为空字符串与「未设置 id」在 SSE 语义中等价源码直接throw new Error兜底。在 sse.ts 的服务端产出逻辑中运行层通过isTrackedEnvelope(value)判断若命中就把三元组映射为 SSE 帧的id:与data:两个字段否则只产出data:字段。也就是说tracked()最终会变成标准 SSE 的id:字段天然符合浏览器EventSource对lastEventId的处理规范。客户端如何回传 lastEventIdlastEventId的传播链路按传输方案略有差异SSE这是浏览器 EventSource 规范 内置的行为——EventSource会自动记录收到的最后一个id:字段并在自动重连时通过Last-Event-ID请求头发给服务器tRPC 会把它解析后放进.input()的lastEventId字段。这一点在 sse.ts 的客户端消费端也有对应逻辑收到message事件时读取msg.lastEventId写入数据封装。WebSocketswsLink会自动发送并维护“最后已知 ID”随着浏览器收到数据不断更新。所以在订阅 procedure 的入参 schema 中你通常应该声明一个可选的lastEventId字段.input( z.object({ // lastEventId 是客户端最后收到的事件 ID // 首次订阅时它是初始化时传入的初始值 // 客户端重连时它是客户端此前收到的最后一条事件 ID lastEventId: z.string().nullish(), }).optional(), )带断线续传的完整订阅实现下面这个示例把「事件驱动」与「基于 lastEventId 的历史补偿」组合起来——它先建立事件监听避免在后续查库补发期间漏掉新事件再按lastEventId从数据库补发历史事件最后转入实时事件流// types: node import EventEmitter, { on } from node:events; import { initTRPC, tracked } from trpc/server; import { z } from zod; class IterableEventEmitter extends EventEmitter { toIterable(eventName: string, opts?: { signal?: AbortSignal }) { return on(this, eventName, opts); } } type Post { id: string; title: string }; const t initTRPC.create(); const publicProcedure t.procedure; const router t.router; const ee new IterableEventEmitter(); export const subRouter router({ onPostAdd: publicProcedure .input( z.object({ lastEventId: z.string().nullish(), }).optional(), ) .subscription(async function* (opts) { // 先订阅事件源保证补发历史数据期间不会漏掉新事件 const iterable ee.toIterable(add, { signal: opts.signal, }); if (opts.input?.lastEventId) { // [...] 查询 lastEventId 之后的帖子并补发 // const items await db.post.findMany({ ... }) // for (const item of items) { // yield tracked(item.id, item); // } } // 持续监听新事件 for await (const [data] of iterable) { const post data as Post; // 给每条事件打上 id保证客户端任意时刻断开都能从这里续传 yield tracked(post.id, post); } }), });:::tip顺序很关键如果“捕捉全部事件”对业务至关重要例如聊天消息一条都不能丢务必先建立事件监听再去数据库按lastEventId取历史数据。否则在“查库 逐条 yield 历史批次”期间新到达的事件会被忽略。本仓库的 next-sse-chat 示例 就是这个顺序的完整参考实现。 :::实际的健壮实现远比注释复杂——因为可能存在「数据库查询结果」与「事件流」之间的竞态导致重复/漏发。在 examples/next-sse-chat/src/server/routers/post.ts 里可以看到它的处理方式用一个maybeYield生成器做统一出口内部用lastMessageCreatedAt游标过滤掉「其他频道的事件」和「比已发送事件更旧的事件」先yield*数据库补发的历史批次再进入事件流的for await从而既去重又不错过断点。定期拉取数据再推送轮询模式事件驱动适合「有实时事件源」的场景如果你的数据源不支持推送比如只有数据库轮询则可以用下面的「拉取」配方订阅 procedure 在while循环里周期性查询数据库查出新增数据后推送给客户端。type Post { id: string; title: string; createdAt: Date }; declare const db: { post: { findMany(opts: { where?: { createdAt?: { gt: Date } }; orderBy?: { createdAt: string } }): PromisePost[]; }; }; declare function sleep(ms: number): Promisevoid; // ---cut--- import { initTRPC, tracked } from trpc/server; import { z } from zod; const t initTRPC.create(); export const publicProcedure t.procedure; export const router t.router; export const subRouter router({ onPostAdd: publicProcedure .input( z.object({ // lastEventId客户端最后收到的事件 ID // 这里用帖子的 createdAt 作为事件 ID lastEventId: z.coerce.date().nullish(), }), ) .subscription(async function* (opts) { // opts.signal 会在客户端断开时被 abort let lastEventId opts.input?.lastEventId ?? null; // 用 while 循环轮询条件是连接尚未被 abort while (!opts.signal!.aborted) { const posts await db.post.findMany({ // 若已有 lastEventId只查它之后创建的帖子 where: lastEventId ? { createdAt: { gt: lastEventId } } : undefined, orderBy: { createdAt: asc }, }); for (const post of posts) { // tracked() 给每条事件附上 id客户端断线重连后可从最后一条继续 yield tracked(post.createdAt.toJSON(), post); lastEventId post.createdAt; } // 睡一会儿再查避免对数据库造成压力 await sleep(1_000); } }), });这个配方有三点实践要点事件 ID 可以用业务时间戳这里把createdAt当作 ID配合z.coerce.date()解析配合createdAt: { gt: lastEventId }实现增量查询循环退出条件就是opts.signal!.aborted客户端断开后服务器自动停止轮询不会空转yield tracked(...)之后再更新lastEventId保证断线续传的游标始终指向已成功产出的事件。从服务端主动停止订阅订阅不一定永远运行。若服务器端需要主动结束某条订阅只需在生成器函数里return即可import { initTRPC } from trpc/server; import { z } from zod; const t initTRPC.create(); const publicProcedure t.procedure; const router t.router; // ---cut--- export const subRouter router({ onPostAdd: publicProcedure .input( z.object({ lastEventId: z.coerce.number().min(0).optional(), }), ) .subscription(async function* (opts) { let index opts.input.lastEventId ?? 0; while (!opts.signal!.aborted) { const idx index; if (idx 100) { // 生成器 return 后订阅结束客户端随之断开 return; } await new Promise((resolve) setTimeout(resolve, 10)); } }), });在客户端一侧对应的停止方式就是调用.unsubscribe()取消订阅。从 WebSockets 传输协议看取消对应客户端发一条subscription.stop消息消息结构见 WebSockets 页面而从本仓库的 procedureBuilder.ts 可以看到subscription本质上是把 resolver 包装成{ ..._def, type: subscription }的过程运行时对异步迭代器执行终止语义。清理订阅的副作用try...finally订阅持有的事件监听、定时器、游标等资源必须被正确释放。tRPC 在订阅因任何原因停止时都会调用生成器实例的.return()因此你完全可以用标准的try...finally模式完成清理// types: node import EventEmitter, { on } from events; import { initTRPC } from trpc/server; type Post { id: string; title: string }; const t initTRPC.create(); const publicProcedure t.procedure; const router t.router; const ee new EventEmitter(); export const subRouter router({ onPostAdd: publicProcedure.subscription(async function* (opts) { let timeout: ReturnTypetypeof setTimeout | undefined; try { for await (const [data] of on(ee, add, { signal: opts.signal, })) { timeout setTimeout(() console.log(Pretend like this is useful)); const post data as Post; yield post; } } finally { if (timeout) clearTimeout(timeout); } }), });这段代码中无论循环是自然结束、被return终止、还是抛错退出finally块都会执行并清除定时器。结合上一节的opts.signal透传机制Nodeon()会随信号中止自动退订就构成了“连接生命周期 资源生命周期”的双重保障。关于Generator.prototype.return()与try...finally的组合语义可参阅 MDN 文档。错误处理与自动重连订阅生成器中的异常传播规则如下生成器函数内抛出错误会传导到服务端的onError()钩子若抛出的错误属于5xx客户端会自动基于tracked()记录的最后事件 ID 尝试重连这正是前面建议所有订阅都使用tracked()的原因之一若是其他错误订阅会被取消并把错误传播给客户端的onError()回调。补充一个实现层面的事实在 sse.ts 中服务端会把用户代码抛出的非 abort 错误转换为getTRPCErrorFromUnknown()错误并通过一个名为serialized-error的专用 SSE 事件帧带formatError序列化下发给客户端而AbortErrorabort 导致的取消会被静默忽略不视为业务错误。此外 SSE 生产者还支持ping.enabled心跳帧默认关闭intervalMs默认 1000ms与client.reconnectAfterInactivityMs空闲重连选项二者由运行层校验 ping 间隔不得大于客户端重连间隔以免引发无谓重连。这套实现细节同样可以在 sse.test.ts 中找到对应测试佐证。订阅输出的运行时校验zAsyncIterable与 query/mutation 不同订阅过程产出的是异步迭代器Zod 的普通 schema 无法直接校验每次yield的值——必须遍历迭代器来逐项校验。官方给出的做法是封装一个针对异步迭代器的 Zod schema 辅助函数。下面是一个针对 Zod v4 的zAsyncIterable实现本仓库测试目录中也有同名可运行版本 packages/tests/server/zAsyncIterable.ts// target: esnext // lib: esnext import type { TrackedEnvelope } from trpc/server; import { isTrackedEnvelope, tracked } from trpc/server; import { z } from zod; function isAsyncIterableTValue, TReturn unknown( value: unknown, ): value is AsyncIterableTValue, TReturn { return !!value typeof value object Symbol.asyncIterator in value; } const trackedEnvelopeSchema z.customTrackedEnvelopeunknown(isTrackedEnvelope); /** * 专门用于校验异步迭代器的 Zod schema 辅助函数。它保证 * 1. 被校验的值确实是一个异步迭代器 * 2. 迭代器每次 yield 的值都符合指定的类型 * 3. 迭代器返回值若有也符合指定的类型。 */ export function zAsyncIterable TYieldIn, TYieldOut, TReturnIn void, TReturnOut void, Tracked extends boolean false, (opts: { /** * 校验异步生成器 yield 出来的值 */ yield: z.ZodTypeTYieldOut, TYieldIn; /** * 校验异步生成器的返回值 * remarks 对订阅不适用 */ return?: z.ZodTypeTReturnOut, TReturnIn; /** * yield 出来的值是否被 tracked() * remarks 仅对订阅适用 */ tracked?: Tracked; }) { return z .custom AsyncIterable Tracked extends true ? TrackedEnvelopeTYieldIn : TYieldIn, TReturnIn ((val) isAsyncIterable(val)) .transform(async function* (iter) { const iterator iter[Symbol.asyncIterator](); try { let next; while ((next await iterator.next()) !next.done) { if (opts.tracked) { // 先用 isTrackedEnvelope 结构识别再解构出 id 与 data 分别校验 const [id, data] trackedEnvelopeSchema.parse(next.value); yield tracked(id, await opts.yield.parseAsync(data)); continue; } yield opts.yield.parseAsync(next.value); } if (opts.return) { return await opts.return.parseAsync(next.value); } return; } finally { await iterator.return?.(); } }) as z.ZodType AsyncIterable Tracked extends true ? TrackedEnvelopeTYieldIn : TYieldIn, TReturnIn, unknown , AsyncIterable Tracked extends true ? TrackedEnvelopeTYieldOut : TYieldOut, TReturnOut, unknown ; }这个辅助函数的核心逻辑分三步先用Symbol.asyncIterator判断输入确实是异步迭代器然后用while循环配合iterator.next()逐项取数据对tracked: true的配置先用isTrackedEnvelope识别包裹结构并解构出[id, data]对data跑一遍opts.yield.parseAsync()后再用tracked(id, ...)原样重组从而保证校验前后的数据类型完全对称最后用finally块调用iterator.return?.()释放底层迭代器。有了这个 helper就可以在订阅 procedure 的.output()里校验产出的每条数据了// target: esnext // lib: esnext // types: node import { initTRPC, tracked } from trpc/server; import { z } from zod; import { zAsyncIterable } from ./zAsyncIterable; const t initTRPC.create(); export const publicProcedure t.procedure; export const router t.router; export const appRouter router({ mySubscription: publicProcedure .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)); } }), });在上例中zAsyncIterable({ yield: z.object({ count: z.number() }), tracked: true })声明了“订阅产出的是被tracked()包裹、每帧数据形如{ count: number }的异步迭代器”。一旦实现里yield出的对象结构不合法例如count不是数字校验会在数据离开服务器之前抛错并中断订阅。小结把本指南的要点串成一份服务端订阅「最佳实践清单」传输选型无特殊需求优先用 SSEhttpSubscriptionLink需要双向或已建 WebSocket 基础设施时用 wsLink WebSockets 服务端接入实现形态订阅 resolver 写异步生成器函数observable 旧式写法已标记废弃见 procedureBuilder.ts生命周期监听事件或定时器时务必透传opts.signal用try...finally兜底清理副作用服务端想结束订阅直接return断线续传所有需要可靠投递的订阅都建议yield tracked(id, data)并通过.input().lastEventId支持重连补偿补发历史数据前先挂上事件监听可参考 post.ts 的完整去重实现错误与校验5xx 错误会触发自动重连其他错误进入onError()对订阅输出做校验需要像zAsyncIterable一样遍历迭代器逐项校验且要能识别tracked()的信封结构。如果你希望看到上述机制在一套完整可运行应用里的实际配合直接阅读 examples/next-sse-chatSSE 断线续传 数据库与 examples/standalone-server极简 Node WebSockets两个参考工程会是投入产出比最高的学习路径。【免费下载链接】trpc♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考