完全指南:从 WebSocket 协议到事件过滤)
Prisma 实时订阅Subscriptions完全指南从 WebSocket 协议到事件过滤【免费下载链接】prisma1 Database Tools incl. ORM, Migrations and Admin UI (Postgres, MySQL MongoDB) [deprecated]项目地址: https://gitcode.com/gh_mirrors/pr/prisma1导读GraphQL 订阅Subscriptions是 Prisma 数据层提供的实时通知机制当数据库中的数据发生创建CREATED、更新UPDATED或删除DELETED三类变更时服务端会通过专用 WebSocket 端点把变更事件实时推送给客户端。本篇指南以 Prisma 1.x 官方参考文档为主体结合本仓库 server/servers/subscriptions 模块的源码实现完整讲解订阅的触发事件、类型订阅语法、where过滤体系、关系变更的变通方案以及底层 WebSocket 消息协议与认证握手原理。读完本文你将能独立写出可复用、可精确过滤的订阅查询并理解订阅端到端的执行链路。订阅概览三种触发事件GraphQL 订阅让你能够在数据发生变更时实时收到通知。驱动订阅的事件event共有三种CREATED一个新的节点被创建UPDATED一个已有节点被更新DELETED一个已有节点被删除。下面的例子订阅了所有新建的Post节点当订阅被触发时服务端推送的 payload 中会携带该Post的description与imageUrl字段subscription newPosts { post(where: { mutation_in: [CREATED] }) { mutation node { description imageUrl } } }订阅使用一个专用的 WebSocket 端点wss://__CLUSTER__.prisma.sh/__WORKSPACE__/__SERVICE__/__STAGE__而不是普通的 HTTP 查询端点。Prisma 会为数据模型中的每个对象类型自动生成对应的类型订阅type subscription。需要注意的一个已知限制是在 relation 中连接或断开节点connect/disconnect不会触发订阅事件官方提供了“触碰字段”的变通方案见下文 Relation subscriptions 小节。你可以在单个订阅请求中组合多个触发器mutation_in: [CREATED, UPDATED, DELETED]精确控制要接收哪些事件。订阅 API 还复用了查询 API 那套丰富的过滤系统可按字段条件、AND/OR逻辑自由组合。发起订阅请求的三种方式方式一GraphQL PlaygroundGraphQL Playground 内置了对订阅的支持可以在服务内直接打开 Playground 来探索、编写并运行订阅查询非常适合调试where过滤条件。方式二Apollo Client 与 apollo-link-ws在浏览器或 Node.js 环境中最常用的做法是使用 Apollo Client 生态的apollo-link-ws库来建立订阅连接。它会自动完成下文描述的 WebSocket 握手、subscription_start与消息分发流程开发者只需提供订阅文档和变量即可。方式三原生 WebSocket 客户端协议详解Prisma 的订阅协议兼容两个子协议服务端在 WebSocket 握手阶段根据客户端声明的子协议选择版本子协议名对应版本说明graphql-subscriptionsV05本文档对应版本init/subscription_start/subscription_end等消息类型graphql-wsV07connection_init/start/stop等消息类型这一协商逻辑可以在源码 WebSocketHandler.scala 中确认它声明了v5ProtocolName graphql-subscriptions与v7ProtocolName graphql-ws两个受支持的协议握手时客户端必须携带其中之一否则连接会被拒绝UnsupportedWebSocketSubprotocolRejection。两种协议下完整的消息类型枚举定义在 SubscriptionProtocol.scala。以下以 V05 协议graphql-subscriptions为例演示从零手写一个原生 WebSocket 订阅客户端。第 1 步建立连接订阅由 WebSocket 管理。首先建立 WebSocket 连接并在握手阶段声明graphql-subscriptions子协议let webSocket new WebSocket(wss://__CLUSTER__.prisma.sh/__WORKSPACE__/__SERVICE__/__STAGE__, graphql-subscriptions);第 2 步发起握手监听open事件然后向服务端发送一条type为init的 JSON 消息webSocket.onopen (event) { const message { type: init } webSocket.send(JSON.stringify(message)) }关于握手有一个容易被忽略但很重要的细节认证。如果服务配置了secretinit消息的payload中必须携带Authorization字段token否则握手会失败。这一点在源码 SubscriptionSessionActorV05.scala 中有清晰的实现当项目没有配置 secret 时直接返回init_success当项目配置了 secret 时会调用AuthImpl.verify(project.secrets, token)校验 token——无效 token 返回init_failAuthentication token is invalid.未提供 token 则返回init_failNo Authorization field was provided in payload.。另外在握手完成之前发送任何订阅消息都会收到init_fail提示必须先发送init。第 3 步接收消息服务端可能返回多种消息通过type属性区分应用可按需响应webSocket.onmessage (event) { const data JSON.parse(event.data) switch (data.type) { case init_success: { console.log(init_success, the handshake is complete) break } case init_fail: { throw { message: init_fail returned from WebSocket server, data } } case subscription_data: { console.log(subscription data has been received, data) break } case subscription_success: { console.log(subscription_success) break } case subscription_fail: { throw { message: subscription_fail returned from WebSocket server, data } } } }对照源码 SubscriptionProtocol.scalaV05 协议的全部服务端消息类型还包括keepalive。服务端会按固定间隔向客户端发送keepalive心跳消息以维持长连接间隔由keepAliveIntervalSeconds配置控制见 WebsocketSession.scala。第 4 步订阅数据变更发送type为subscription_start的消息并在query字段中携带订阅文档const message { id: 1, type: subscription_start, query: subscription newPosts { post(filter: { mutation_in: [CREATED] }) { mutation node { description imageUrl } } } } webSocket.send(JSON.stringify(message))你会先收到一条subscription_success消息表示订阅建立成功此后每当数据发生符合条件的变化就会收到subscription_data消息。你提供的id会出现在所有subscription_data消息中因此可以通过不同的id在同一个 WebSocket 连接上多路复用multiplex多个订阅。从源码看subscription_start消息还支持variables与operationName字段SubscriptionProtocol.scala。operationName用于在查询文档中包含多个操作时指定执行哪一个——这一点在协议测试 SubscriptionsProtocolV05Spec.scala 中有专门用例验证同一条消息内同时带有subscription x {...}与mutation y {...}时通过operationName:x可以正确选中所要执行的订阅。此外若订阅文档语法非法服务端会返回subscription_fail并提示 the GraphQL Query was not valid若文档中不包含任何已知的模型名比如用了旧的createTodo顶层字段语法同样会返回subscription_fail提示 The provided query doesnt include any known model name.见 SubscriptionSessionActorV05.scala 与协议测试中的断言。第 5 步取消订阅发送type为subscription_end的消息即可停止接收某条订阅的数据const message { id: 1, type: subscription_end } webSocket.send(JSON.stringify(message))服务端收到后会从该模型的订阅表中移除对应条目之后该id不再收到任何subscription_data。测试中还展示了不带id的subscription_end是合法的被安全忽略见 SubscriptionsProtocolV05Spec.scala。类型订阅Type Subscriptions数据模型中的每个对象类型都会自动生成一个对应的类型订阅。以如下只有一个Post类型的数据模型为例type Post { id: ID! unique title: String! description: String }生成的 Prisma API 中会提供post订阅用于在Post节点被创建、更新或删除时通知客户端。订阅创建的节点CREATED订阅所有新建节点使用where对象并设置mutation_in: [CREATED]subscription { post(where: { mutation_in: [CREATED] }) { mutation node { description imageUrl author { id } } } }payload 中包含mutation此处返回CREATEDnode可查询新建节点的信息以及关联节点。订阅特定的新建节点通过where对象中的node参数可以使用与查询一致的字段过滤系统。例如只有当特定用户**关注follow**了author时才通知我新建了Postsubscription { post(where: { AND: [{ mutation_in: [CREATED] }, { node: { author: { followedBy_some: { id: cj03x3nacox6m0119755kmcm3 } } }] }) { mutation node { description imageUrl author { id } } } }node过滤器还支持与查询变量variables配合。协议测试 SubscriptionsProtocolV05Spec.scala 演示了subscription_start中携带variables: {text: some}、查询里使用$text: String!与node: {text_contains: $text}的写法服务端会按变量值过滤后推送。订阅删除的节点DELETED订阅所有被删除的节点设置mutation_in: [DELETED]即可订阅所有删除事件subscription deletePost { post(where: { mutation_in: [DELETED] }) { mutation previousValues { id } } }payload 包含mutation此处返回DELETEDpreviousValues被删除节点的标量字段旧值。注意对于CREATED订阅previousValues永远为null。订阅特定的删除事件同样可以使用node参数做字段过滤。例如只有当特定用户关注了author时才通知我某个Post被删除subscription { post(where: { mutation_in: [DELETED] node: { author: { followedBy_some: { id: cj03x3nacox6m0119755kmcm3 } } } }) { mutation previousValues { id } } }从测试用例看删除事件中node内查询的字段在推送时会被解析为null节点已不存在因此删除订阅中真正有意义的 payload 是previousValues见 SubscriptionsProtocolV05Spec.scala。订阅更新的节点UPDATED订阅所有被更新的节点设置mutation_in: [UPDATED]subscription { post(where: { mutation_in: [UPDATED] }) { mutation node { description imageUrl author { id } } updatedFields previousValues { description imageUrl } } }payload 包含mutation此处返回UPDATEDnode可查询更新后的节点及其关联节点updatedFields发生变更的字段列表previousValues节点的标量字段旧值。注意对于CREATED与DELETED订阅updatedFields永远为null对于CREATED订阅previousValues永远为null。订阅特定字段的更新where对象支持对更新字段做过滤。例如仅当Post的description字段发生变化时才通知subscription { post(where: { mutation_in: [UPDATED] updatedFields_contains: description }) { mutation node { description } updatedFields previousValues { description } } }与updatedFields_contains类似还有更多过滤条件updatedFields_contains_every: [String!]指定的字段全部被更新时才匹配updatedFields_contains_some: [String!]指定的字段中至少一个被更新时匹配。注意不能将updatedFields系列过滤条件与mutation_in: [CREATED]或mutation_in: [DELETED]组合使用会产生错误。这与服务端实现一致updatedFields仅在UPDATED事件中有意义——源码 SubscriptionResolver.scala 中只有DatabaseUpdateEvent才会携带changedFields对应updatedFields并传给查询执行器而创建与删除事件中该值为None。updatedFields_contains与mutation_in的组合过滤在测试 SubscriptionFilterSpec.scala 中也有覆盖包括对枚举类型字段旧值previousValues.status的过滤。关系订阅Relation Subscriptions如前文所述连接/断开关系节点connect / disconnect本身不会触发订阅事件。目前只能借助UPDATED订阅通过变通方式实现关系变更通知订阅关系变更的变通方案触碰touch字段给相关类型添加一个dummy: String字段然后在关系状态发生变化时更新该字段的值从而“逼出”一条UPDATED事件mutation updatePost { updatePost( where: { id: some-id } data: { dummy: dummy # do a dummy change to trigger update subscription } ) }随后再配合updatedFields_contains: dummy或组合updatedFields_contains_every/updatedFields_contains_some来精确匹配这类“触碰”更新即可模拟出关系变更通知的效果。组合订阅Combining Subscriptions可以在同一个订阅中对同一类型订阅多种变更事件。订阅所有节点的所有变更利用where对象的mutation_in参数选择要订阅的变更类型。例如同时订阅createPost、updatePost与deletePost三类变更subscription { post(where: { mutation_in: [CREATED, UPDATED, DELETED] }) { mutation node { id description } updatedFields previousValues { description imageUrl } } }订阅特定节点的所有变更用node参数选定关心的节点并可与mutation_in组合。例如只通知我特定用户关注的作者发布的文章发生创建、更新或删除subscription { post( where: { mutation_in: [CREATED, UPDATED, DELETED] } node: { author: { followedBy_some: { id: cj03x3nacox6m0119755kmcm3 } } } ) { mutation node { id description } updatedFields previousValues { description imageUrl } } }注意previousValues对于CREATED订阅永远为nullupdatedFields对于CREATED与DELETED订阅永远为null。高级订阅过滤器AND / OR 逻辑组合where参数完整支持查询 API 的过滤系统包括AND与OR逻辑操作符。例如订阅所有CREATED与DELETE事件外加所有imageUrl字段被更新的UPDATED事件subscription { post(where: { OR: [{ mutation_in: [CREATED, DELETED] }, { mutation_in: [UPDATED] updatedFields_contains: imageUrl }] }) { mutation node { id description } updatedFields previousValues { description imageUrl } } }注意将任何updatedFields过滤条件与CREATED或DELETED组合会报错previousValues对CREATED恒为nullupdatedFields对CREATED与DELETED恒为null。深入源码订阅的端到端执行链路理解订阅不仅能“会用”还能“知其所以然”。以下是结合本仓库源码梳理的订阅完整链路。1. 连接建立与协议协商客户端发起 WebSocket 握手时WebSocketHandler.scala 根据Sec-WebSocket-Protocol头选择 V05 或 V07 协议并生成一个携带sessionIdcuid的 WebsocketSession。该 Session Actor 负责解析入站 JSON 消息、按协议路由到对应的会话 Actor、维护 keepalive 心跳、以及超时默认 10 分钟无消息则断开处理。2. 握手与认证init消息到达后SubscriptionSessionActorV05 按上文所述三种情况无 secret / 有 secret 且 token 有效 / 有 secret 且 token 无效或缺失返回init_success或init_fail。认证通过的会话进入readyReceive状态才能处理后续订阅消息。3. 订阅注册与按模型分发subscription_start到达后会话 Actor 用 sangria 解析查询文档随后把CreateSubscription请求发给全局 SubscriptionsManager。这里有两件事通过SubscriptionQueryValidator校验查询并定位目标模型不包含任何已知模型名会在此处被拒绝通过QueryTransformer.getMutationTypesFromSubscription从查询中提取mutation_in指定的变更类型集合CREATED/UPDATED/DELETED。随后每个模型对应一个 SubscriptionsManagerForModel Actor它在启动时向消息总线订阅该模型的三个变更通道。4. 变更事件的广播与通道命名数据库变更事件通过消息总线message bus广播。通道命名规则定义在 MutationChannelUtil.scala格式为subscription:event:projectId:channelName其中channelName按模型生成createModel、updateModel、deleteModel例如createTodo、updateTodo、deleteTodo。这正是协议测试中向subscription:event:${project.id}:createTodo等通道发布测试事件的原因。5. 数据库事件 → GraphQL payloadSubscriptionsManagerForModel 收到消息总线事件后会按“查询文本 变量”分组对相同查询的多个订阅只执行一次过滤查询并复用结果性能优化随后逐个通知订阅者。真正的数据组装发生在 SubscriptionResolver.scala依据事件类型DatabaseCreateEvent/DatabaseUpdateEvent/DatabaseDeleteEvent定义见 DatabaseEvents.scala分别构造previousValues与updatedFields通过SubscriptionExecutor.execute用原始查询文档对变更节点执行一次 GraphQL 查询得到推送给客户端的 payload生产环境中数据库副本可能落后主库最多约 20ms因此代码中人为加入了35ms 的缓冲延迟源码注释明确要求不要移除该延迟确保读取到最新数据。最终payload 以subscription_data消息的形式、带上发起订阅时的id返回客户端客户端断开或发送subscription_end后订阅被移除见handleTerminatedSubscriber与EndSubscription的处理逻辑。结语Prisma 的订阅能力覆盖了数据层实时通知的完整闭环通过mutation_in选择事件类型、通过node字段过滤与updatedFields_*条件做精细匹配、通过AND/OR组合任意复杂逻辑、通过专用 WebSocket 协议V05graphql-subscriptions/ V07graphql-ws在一条连接上多路复用多个订阅。配合本仓库 server/servers/subscriptions 下的协议实现与测试用例你可以将文档中的每个语法点都落到可验证的源码行为上从而在生产环境中自信地使用实时订阅构建聊天、协作编辑、消息推送等实时功能。【免费下载链接】prisma1 Database Tools incl. ORM, Migrations and Admin UI (Postgres, MySQL MongoDB) [deprecated]项目地址: https://gitcode.com/gh_mirrors/pr/prisma1创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考