
tRPC 订阅与 WebSockets 实战指南端到端类型安全的实时通信从入门到 JSON-RPC 协议剖析【免费下载链接】trpc♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc本文以 tRPC 官方文档 Subscriptions / WebSockets 为主线完整讲解如何在 tRPC 中定义实时订阅subscription、搭建 WebSocket 服务端、配置客户端传输链路并深入剖析其基于 JSON-RPC 2.0 的线上消息协议。通过本仓库中 服务端适配器源码、消息信封类型 以及 standalone-server 示例 的交叉印证读完你将能够独立实现从后端事件源到前端 React 组件的完整实时推送链路。为什么实时场景要引入 SubscriptionstRPC 的query与mutation天然是一次性请求-响应模型适合处理“拉取”型数据但当服务端需要主动向客户端推送数据如新消息提醒、聊天输入状态、随机数流时就需要长连接。tRPC 的subscriptionprocedure 通过 WebSockets 传输配合官方提供的observable抽象可以在保持端到端类型安全的前提下把事件流推送变成与query/mutation一致的调用体验——客户端订阅时类型自动推断、取消订阅时代码自动清理。在服务端定义一个订阅 procedure完整示例基于 EventEmitter 的实时公告牌官方文档给出的经典场景是“添加帖子后实时通知所有订阅者”。事件源可以是内存EventEmitter也可以换成 Redis、Kafka 等任何可发布订阅的基础设施import { EventEmitter } from events; import { initTRPC } from trpc/server; import { observable } from trpc/server/observable; import { z } from zod; // create a global event emitter (could be replaced by redis, etc) const ee new EventEmitter(); const t initTRPC.create(); export const appRouter t.router({ onAdd: t.procedure.subscription(() { // return an observable with a callback which is triggered immediately return observablePost((emit) { const onAdd (data: Post) { // emit data to client emit.next(data); }; // trigger onAdd() when add is triggered in our event emitter ee.on(add, onAdd); // unsubscribe function when client disconnects or stops subscribing return () { ee.off(add, onAdd); }; }); }), add: t.procedure .input( z.object({ id: z.string().uuid().optional(), text: z.string().min(1), }), ) .mutation(async (opts) { const post { ...opts.input }; /* [..] add to db */ ee.emit(add, post); return post; }), });代码要点拆解observable来自trpc/server/observabletRPC 维护了一个精简的 Observable 实现它遵循“订阅时执行、退订时清理”的契约。订阅回调接收一个emit对象其中emit.next(data)向客户端推送数据回调返回的清理函数会在客户端断开或主动停止订阅时执行这里用ee.off(add, onAdd)移除事件监听。生产环境事件源替换内存EventEmitter只在单进程内有效。多实例部署时应将ee替换为 Redis Pub/Sub 等跨进程广播方案但 procedure 的写法不变——你只需在收到外部事件时调用emit.next(...)。addmutation 负责触发事件将新帖子写入数据库后ee.emit(add, post)于是每个正在订阅onAdd的客户端都会实时收到这条数据。从源码看packages/server/src/observable/observable.ts 中的observable()会为你创建的订阅函数补充unsubscribe语义一旦next/error/complete触发后自动执行清理逻辑teardownRef并保证多次退订幂等类型inferObservableValue、isObservable等辅助工具也由 index.ts 统一导出。创建 WebSocket 服务端安装依赖yarn add ws使用applyWSSHandler挂载路由import { applyWSSHandler } from trpc/server/adapters/ws; import ws from ws; import { appRouter } from ./routers/app; import { createContext } from ./trpc; const wss new ws.Server({ port: 3001, }); const handler applyWSSHandler({ wss, router: appRouter, createContext }); 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(); });applyWSSHandler接收三要素wssws模块的WebSocketServer、appRoutertRPC 路由、createContext与 HTTP 适配器共用的上下文工厂。接线后 WebSocket 端口上的所有消息会自动路由到对应 procedure。连接生命周期观测connection/close事件用来打印当前在线连接数。优雅停机是关键一步——收到SIGTERM时先调用handler.broadcastReconnectNotification()通知所有客户端“服务器即将重启请重连”再关闭wss从而避免停机瞬间静默丢消息。从本仓库源码看packages/server/src/adapters/ws.ts 中applyWSSHandler会在每个连接上执行getWSConnectionHandler其内部用Mapid, AbortControllerclientSubscriptions追踪每个活跃订阅收到subscription.stop时通过clientSubscriptions.get(id)?.abort()中止对应流连接关闭时则批量 abort 所有订阅并清理见 ws.ts。此外该源码还支持可选的keepAlive心跳默认关闭、prefix路径过滤与自定义experimental_encoder说明适配器在主干上持续演进。不重复造 HTTP 服务器standalone 一键双协议仓库中的 examples/standalone-server/src/server.ts 展示了同时暴露 HTTP 与 WebSocket 的最小写法——让ws复用同一个 HTTP serversubscription之外的所有调用都走 HTTP// http server const server createHTTPServer({ router: appRouter, createContext }); // ws server const wss new WebSocketServer({ server }); applyWSSHandlerAppRouter({ wss, router: appRouter, createContext }); server.listen(2022);设置 TRPCClient 使用 WebSockets仅订阅走 WebSocketwsLinkimport { createTRPCProxyClient, createWSClient, wsLink } from trpc/client; import type { AppRouter } from ../path/to/server/trpc; // create persistent WebSocket connection const wsClient createWSClient({ url: ws://localhost:3001, }); // configure TRPCClient to use WebSockets transport const client createTRPCProxyClientAppRouter({ links: [ wsLink({ client: wsClient, }), ], });推荐splitLink让 query/mutation 走 HTTP、订阅走 WebSocket官方文档建议用 Links 概览 中的链接机制做“分流”。因为 WebSocket 长连接成本更高最佳实践是让一次性请求仍走 HTTP仅把 subscription 交给 WebSocket。仓库示例 examples/standalone-server/src/client.ts 展示了这种“双终止链接”结构const wsClient createWSClient({ url: ws://localhost:2022 }); const trpc createTRPCClientAppRouter({ links: [ // call subscriptions through websockets and the rest over http splitLink({ condition(op) { return op.type subscription; }, true: wsLink({ client: wsClient }), false: httpLink({ url: http://localhost:2022 }), }), ], });这里splitLink依据op.type subscription判断订阅类操作分发给wsLink其余走httpLink。完整选项与更多用法可参考 splitLink 文档。调用订阅时与 Promise 风格不同你需要传入onData回调并显式unsubscribelet count 0; await new Promisevoid((resolve) { const subscription trpc.post.randomNumber.subscribe(undefined, { onData(data) { // ^ note that data here is inferred console.log(received, data); count; if (count 3) { subscription.unsubscribe(); // stop after 3 pulls resolve(); } }, onError(err) { console.error(error, err); }, }); }); await wsClient.close();对应服务端的randomNumber订阅定义在 standalone-server/src/server.ts它每 200msemit.next()一个随机数客户端累计收到 3 次后主动unsubscribe()此时服务端回调返回的clearInterval清理函数被触发。在 React 中使用订阅文档指出 React 场景可参考全栈示例。本仓库中的 examples/next-prisma-websockets-starter 是可直接运行的参考对应文档所指向的 Next.js Prisma WebSocket starter。在该示例的 index.tsx 中订阅与 React Hooks 无缝集成// subscribe to new posts and add trpc.post.onAdd.useSubscription(undefined, { onData(post) { addMessages([post]); }, onError(err) { console.error(Subscription error:, err); // we might have missed a message - invalidate cache utils.post.infinite.invalidate(); }, });要点onData中把推送内容写入本地状态onError里做兜底——例如连接中断错过消息时调用utils.post.infinite.invalidate()让查询缓存重新拉取从而“WebSocket 断了也不丢数据”。这样的订阅在组件卸载时由 hook 自动完成取消与服务端清理函数协同形成完整的资源释放闭环。WebSockets RPC 协议规范线上消息是 JSON-RPC 2.0 风格的 JSON 文本帧其类型骨架定义在仓库的 envelopes.ts 中对应的错误码定义见 codes.ts。query/mutation请求{ id: number | string; jsonrpc?: 2.0; // optional method: query | mutation; params: { path: string; input?: unknown; // -- pass input of procedure, serialized by transformer }; }响应{ id: number | string; jsonrpc?: 2.0; // only defined if included in request result: { type: data; // always data for mutation / queries data: TOutput; // output from procedure } }subscription/subscription.stop开启订阅请求{ id: number | string; jsonrpc?: 2.0; method: subscription; params: { path: string; input?: unknown; // -- pass input of procedure, serialized by transformer }; }取消订阅向服务端发送subscription.stop其中id必须是创建订阅时所用的 id{ id: number | string; // -- id of your created subscription jsonrpc?: 2.0; method: subscription.stop; }订阅响应形状除错误外服务端会陆续推送以下消息{ id: number | string; jsonrpc?: 2.0; result: ( | { type: data; data: TData; // subscription emitted data } | { type: started; // subscription started } | { type: stopped; // subscription stopped } ) }源码印证TRPCResultMessage在 envelopes.ts 中正是{ type: started } | { type: stopped } | TRPCResult的联合类型而 ws.ts 的处理器在收到method subscription.stop时直接查表并abort()对应订阅。这解释了规范中的约定subscription.stop的id必须指向已建立的订阅服务端随后会推进迭代器终止并向该id发送{ type: stopped }帧。错误处理错误响应遵循 JSON-RPC 2.0 的 error object 结构code / message / data具体格式与自定义请参考 Error Formatting。本仓库 codes.ts 定义了 tRPC 的整组 JSON-RPC 错误码——其设计规则是复用 JSON-RPC 标准保留区间-32000 ~ -32099并将 HTTP 4XX 状态码的最后两位映射到该区间末尾例如语义键数值对应 HTTPPARSE_ERROR-32700400JSON 解析失败BAD_REQUEST-32600400INTERNAL_SERVER_ERROR-32603500UNAUTHORIZED-32001401FORBIDDEN-32003403NOT_FOUND-32004404TIMEOUT-32008408CONFLICT-32009409TOO_MANY_REQUESTS-32029429CLIENT_CLOSED_REQUEST-32099499错误响应中的code字段携带的即为上表数值。客户端可通过TRPCClientError的data.code拿到语义化键值用于分支判断。与订阅直接相关的是服务端判断某订阅未返回observable/AsyncGenerator时会回INTERNAL_SERVER_ERROR重复使用同一id发起订阅会收到BAD_REQUEST见 ws.ts。服务端到客户端的通知重连信号在需要重启或重新部署服务前服务端应广播重连通知客户端收到后会重新建立连接从而让连接落在新实例上。文档给出的帧形状为{ id: null, type: reconnect }它由wssHandler.broadcastReconnectNotification()触发这也是上文SIGTERM优雅停机中关键的一步。在本仓库当前源码中该通知以TRPCReconnectNotification类型定义于 envelopes.ts并作为{ id: null, method: reconnect }由 ws.ts 的broadcastReconnectNotification序列化后广播给所有处于 OPEN 状态的客户端。更完整的落地示例与延伸阅读最小可用全栈示例examples/standalone-serverHTTP WebSocket 同端口、原生TRPCClient核心见 server.ts 与 client.ts。Next.js 全栈参考examples/next-prisma-websockets-starter覆盖聊天式实时列表、whoIsTyping输入状态等真实交互。服务端协议类型定义envelopes.ts 与 codes.ts。适配器实现与测试ws.ts 及其配套测试 packages/tests/server/websockets.test.ts可帮助你理解协议的各种边界行为。建议按“先跑通 standalone 示例再对照协议文档阅读 ws 适配器源码”的顺序学习把事件源从内存EventEmitter替换为 Redis 等外部总线后这套订阅骨架即可直接迁移到生产环境。【免费下载链接】trpc♀️ Move Fast and Break Nothing. End-to-end typesafe APIs made easy.项目地址: https://gitcode.com/GitHub_Trending/tr/trpc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考