Nhost Constellation 订阅系统深度解析:从 WebSocket 握手到 cohort 多路复用轮询的端到端实现
【免费下载链接】nhostThe Open Source Firebase Alternative with GraphQL.项目地址: https://gitcode.com/GitHub_Trending/nh/nhost
导读
本文以 services/constellation/docs/developers/subscriptions.md 为骨架,结合 Nhost 开源仓库中 Constellation 控制器的实际源码,完整拆解一条 GraphQL 订阅的生命周期:WebSocket 握手、查询解析与路由、cohort 批量聚合、多路复用轮询(multiplexed polling)、变更检测、背压处理与最终拆除。读完本文,你将掌握 Constellation 如何用「一个 SQL 查询服务成百上千个相同订阅」的核心思想,以及_stream游标订阅与普通实时订阅在实现上的本质差异,并能在自己的接入层设计中复用这套「cohort 聚合 + 单轮询扇出」的架构模式。
前提说明:以下所有结论均以当前仓库(
services/constellation/目录)的代码与文档为准;文中给出的文件路径均为仓库根目录下的相对路径。
为什么需要 cohort(订阅聚合)
在 GraphQL 订阅场景里,最朴素的实现是为每个订阅启动一个轮询 goroutine,每次轮询执行一条独立 SQL。但真实负载中大量订阅者在问同一个问题——subscription { messages { ... } }是典型形态。若 N 个订阅各自查询,数据库就要承受 N 倍的查询压力。
Constellation 的答案是把这些订阅聚合进cohort(同类批次):同一 cohort 内所有订阅共享同一条多路复用 SQL,把 N 次查询压缩成 1 次,再由 Postgres 的UNNEST把结果按订阅者扇出(fan-out)。
代价是 cohort 成员资格要求严格一致:
- GraphQL 查询字符串(以 xxhash 摘要形式参与键值);
- 变量(
$limit: 10与$limit: 20会生成不同的 SQL,不能共享 cohort); - 角色(role,决定行级权限如何应用);
- 操作名(operationName)。
会话变量(session variables,即x-hasura-*)不参与 cohort 键——它们在多路复用查询内部按订阅者逐个绑定,这正是把 N 个订阅压缩成一个查询的关键前提(详见 connector/sql/subscription/cohort.go 中cohortKey的定义)。
端到端调用链总览
文档给出了完整的调用链图,整理为如下流程:
- 客户端打开 WebSocket →
controller/handlers.go的HandlerWebsocket入口; controller/websocket/(纯协议层):升级 HTTP、启动 readPump/writePump goroutine、解析graphql-transport-ws帧并分发到MessageHandler;controller/websocket.go的webSocketHandler(每连接状态):连接时快照controllerState,OnConnectionInit提取会话,OnSubscribe解析校验查询并路由到subscription.Handler,OnComplete/OnClose停止对应订阅;subscription.Handler(接口定义在subscription/,实现按 connector 区分):connector/sql/subscription.Handler依据 stream 还是普通实时查询分派;- 普通实时查询 →
cohortManager;subscription_stream→streamCohortManager; Driver.ExecuteMultiplexedOperation(目前仅 Postgres):UNNEST($1::text[], $2::json[])展开_subs("result_id", "result_vars"),内层查询通过 JSON path 运算符读取会话变量与游标;- 结果分发回订阅者:按订阅者计算 xxhash,负载未变化则跳过发送;
forwardUpdatesgoroutine 把next帧写入sendCh;背压采用「1 深度通道 + 最新值胜出(latest-wins)」。
1. WebSocket 协议层:纯协议、零业务逻辑
controller/websocket/是一个纯粹的协议处理器(包内架构见 controller/websocket/doc.go)。它实现graphql-transport-ws规范,负责帧的读写泵(readPump/writePump)与帧封装,其余一切通过MessageHandler接口委托给调用方:
type MessageHandler interface { OnConnectionInit(ctx, payload) // auth / session OnSubscribe(ctx, id, payload) // start a sub OnComplete(ctx, id) // stop a sub OnClose(ctx) // tear down all subs }该包不持有任何业务逻辑。ping/pong 与 connection-ack 帧由协议层自动发送;调用方只通过共享的sendCh通道发送next、error、complete帧。这种分层让协议语义与订阅调度彻底解耦:未来哪怕把协议换成 SSE 或别的传输,订阅侧完全不用动。
2. 每连接桥接层:webSocketHandler
controller/websocket.go 是协议层与订阅系统之间的桥。webSocketHandler为每个连接构造一次,持有三类关键状态:
state——连接建立时对controllerState的快照,保证该连接上所有订阅看到一致的 schema/connector 视图,即使期间发生元数据重载也不受影响;当state.done通道关闭时连接自毁(与元数据重载联动,详见 architecture.md);session——由OnConnectionInit通过middleware.ExtractSession填充,优先级为 admin secret → JWT → public role;subs——syncmap.Map[string, *subscriptionState],以订阅 ID 为键(该类型 map 来自仓库根目录 internal/lib/syncmap/syncmap.go)。每个条目记住哪个subscription.Handler拥有该订阅,确保OnComplete/OnClose总能停在正确(可能已变老)的 handler 上。
OnSubscribe完成每个订阅的预检:
operation, fragments, validatedVariables, err := parseAndValidateQuery(...) dbName := getConnectorForOperation(state, operation) subHandler := state.subHandlers[dbName] h.startSubscription(ctx, id, payload, subHandler, operation, fragments, validatedVariables, logger)从源码看,parseAndValidateQuery是 HTTP 查询路径Resolve解析步骤的「订阅孪生」:同样命中queryCache、运行 gqlparser 校验、做变量强制转换(coerce),并额外对根选择集做@skip/@include求值与片段展开归一化(见 controller/websocket.go)。路由则取最简策略:由第一个根字段所属的 connector 决定(getConnectorForOperation按state.fieldToConnector映射查找);订阅不会跨 connector 扇出——Controller.execute会拒绝包含 remote relationship 的订阅查询。
startSubscription构造subscription.Request(通过NewRequest校验必填字段),调用subHandler.Start得到<-chan subscription.Update,再启动forwardUpdatesgoroutine 把更新翻译成next/error帧写回 WebSocket。
3. Handler 接缝:三方法接口与纯数据结构
subscription包(顶层 subscription/types.go)只承载三个纯数据形状,不含任何行为:
Request——查询字符串 + 已解析的Operation+ 角色 + 变量 + 会话变量。NewRequest强制校验四个承重字段(ID、QueryString、Operation、Role)非空,缺一即返回ErrInvalidRequest;Update——{SubscriptionID, Data jsontext.Value, Error}。Data采用jsontext.Value(来自encoding/json/v2),让 connector 可以把序列化字节直接交给下游而无需二次 marshal;Update通过NewUpdateData/NewUpdateError构造,从调用点保证「Data/Error 二选一」不变量。注意Error非终止性:通道在Start关闭前会持续投递;Handler——每个 connector 必须实现的 3 方法接口:
type Handler interface { Start(ctx context.Context, req Request, logger *slog.Logger) (<-chan Update, error) Stop(ctx context.Context, subscriptionID string) Shutdown(ctx context.Context) }包注释(subscription/types.go)说明了该接缝存在的三个理由:控制器依赖稳定接口而非具体 connector 类型;未来 CDC、消息总线等新策略可直接插拔而不触碰 WebSocket 层;控制器与 connector 之间无 import cycle。此外,ErrInvalidSubscription哨兵错误用于把「客户端查询不可规划」这类由订阅形态导致的失败与驱动/运行时故障区分开,前者会被协议层原样呈现给客户端而非折叠成模糊的 internal server error。
4. SQL connector handler:stream 与实时查询的路由分派
connector/sql/subscription.Handler(connector/sql/subscription/handler.go)在两个 manager 之间路由:
isStream, cursorValues, cursorColumns, err := h.detectStreamSubscription(req) if isStream { return h.streamCohortMgr.addSubscription(ctx, req, cursorValues, cursorColumns, logger) } return h.cohortMgr.addSubscription(ctx, req, logger)stream 检测本身是廉价的:QueryBuilder.IsStreamSubscription(field)是 O(1) 的名字检查(根字段以_stream结尾),普通实时订阅无需构建任何 SQL 即可快速返回。只有 stream 路径才会付出BuildQuery的成本——而且目的是收集游标元数据(ExtractInitialCursorValues与游标列名),而非为了 SQL 本身。handler.go 顶部还定义了QueryExecutor与QueryBuilder两个接口(均带 mockgen 指令),其中ExecuteMultiplexedQueryWithCursor专门服务 stream 订阅的游标参数。
Stop对两个 manager 各调用一次removeSubscription,由于订阅在Start时已分区,另一个 manager 只是 O(1) 的索引未命中,成本可忽略。
5. cohortManager:普通实时查询的聚合轮询
cohortManager处理只依赖时间变化(底层表数据变动)而非客户端游标的订阅。包级架构图见 connector/sql/subscription/doc.go。
5.1 cohort 键:变量值参与、会话变量缺席
type cohortKey struct { queryHash string // xxhash of the GraphQL query string role string operationName string variablesHash string // xxhash of GraphQL variable values }变量值是键的一部分:$limit: 10与$limit: 20会产生不同 SQL,无法共享 cohort(newCohortKey对变量做排序键的确定性 xxhash,见 cohort.go)。会话变量不在键中——它们在多路复用查询内按订阅者绑定。
5.2 容量与溢出链
maxCohortSize为100(cohort.go 的常量定义)。findOrCreateCohort沿key、key_overflow_1、key_overflow_2……依次查找第一个有空位的 cohort,全部满员则新建一个;溢出 cohort 拥有独立的轮询 goroutine,行为与主 cohort 完全一致(createOverflowKey通过给 operationName 追加_overflow_N后缀实现区分,见 cohort_manager.go)。
5.3 轮询循环:与请求上下文解耦
pollCohort运行在context.Background()下,刻意与任何订阅者的请求上下文解耦——否则第一个订阅者断开就会连带取消所有人的轮询(源码在 cohort_manager.go 有//nolint:contextcheck注释说明)。终止路径只有两条:Handler.Shutdown,或轮询 tick 中发现的空 cohort 清理。
每个 tick 的五个步骤:
- 在 cohort 锁下快照订阅集合(
getSubscriptionsCopy,避免查询执行期间长期持锁); - 组装订阅者输入:订阅 ID 数组 + 每个变量名的
[]any值数组(buildSubscriberInputs); - 首次调用构建 SQL,之后复用缓存的
*core.SQLOperation(getOrBuildSQL)——cohort 键固定意味着 SQL 形状稳定,整个轮询执行计划只构建一次。缓存构建时使用core.SessionVarValue{Name: varName}模板标记会话变量(只带名字、不带值),让多路复用转换器能按类型识别真正的权限会话变量并重写为逐订阅者的result_vars查找,而用户提供的以x-hasura-开头的字面量仍是普通数据; - 经
QueryExecutor.ExecuteMultiplexedQuery执行(底层是 SQL 驱动的ExecuteMultiplexedOperation); - 把结果解复用回各订阅者(
distributeResults)。
5.4 变更检测:xxhash 跳过重复负载
每个 cohort 订阅维护lastHash(负载字节的 xxhash)。distributeResults计算新哈希,相同则跳过发送,且只在sendUpdate成功时才更新lastHash。首次轮询必然发送(lastHash初始为空串),保证订阅者拿到基线数据。
5.5 背压:1 深度通道 + 最新值胜出
cohortSubscription.updateCh是容量为 1 的缓冲通道。sendUpdate尝试发送;缓冲满则先排空陈旧条目再重试放入新条目(cohort.go 的 drain-and-retry 逻辑)。语义是latest-wins:慢消费者永远不会看到陈旧数据,但可能错过中间更新——对「当前状态」型订阅这是正确默认值。
sendMu互斥锁串行化sendUpdate与stop,保证关闭通道永远不会与并发发送竞争(stop幂等,重复调用是 no-op)。若轮询中 SQL 构建失败,broadcastError把同一错误发给 cohort 内所有订阅者,错误经ErrInvalidSubscription包装后由协议层原样呈现。
6. streamCohortManager:游标驱动的流式订阅
subscription_stream订阅与普通实时订阅不同:每个订阅者追踪自己的游标位置(典型为自增序列列或时间戳列)。cohort 键因此包含游标哈希,使处于同一位置的订阅者聚合在一起;随着游标推进到相同值,cohort 自然合并:
type streamCohortKey struct { queryHash string role string operationName string variablesHash string cursorHash string // hash of serialised cursor values }6.1 每轮重建(executeAndRebuild)
每次轮询后,executeAndRebuild产出新的 cohort 映射(stream_cohort_manager.go):
- 游标前进的订阅者按新游标位置重新入键(
rebuildCohortMap→reseatStreamCohort重设键,或mergeStreamCohort合并进已存在的同键 cohort——后者先搬订阅、clearSubscriptions再stop源 cohort,保证被迁移订阅者的通道保持打开); - 轮询期间新到达的订阅者进入独立的「initial-data」cohort(
startPolling/endPolling期间写入newSubscriptions,轮询结束后由processNewSubscribers按游标哈希分组,attachOrCreateCohortForCursor决定并入已有 cohort 或新建),先收到追赶(catch-up)基线负载,再在当前位置并入主 cohort。
正是这个重建机制让 cohort 合并成为涌现行为:多个 cohort 的游标推进到同一值后,下一次 tick 上它们的订阅者就落在同一个键里。
6.2 游标提取:单次解析优化
pickCursorFromResults对每个结果行只解析一次并读取游标列以推进 cohort。早期实现每轮对每个负载重复解析三次;当前单解析路径是刻意的性能优化(parseStreamResults把空结果跳过与游标提取合并到同一次 unmarshal,见 stream_cohort_manager.go)。游标值缺失或为 NULL 的列回退到上一轮游标,语义对齐 Hasura 的mergeOldAndNewCursorValues。另注意sendStreamResults中初始轮询总是发送(即使无行),让订阅者先看到基线再进入静默。
7. 多路复用 SQL:UNNEST + LATERAL + JSON path
connector/sql/graphql/queries/multiplexed/multiplexed.go把单个 GraphQL 订阅的 SQL 改写成 Hasura 风格的多路复用形态:
SELECT "_subs"."result_id", "_fld_resp"."root" AS "result" FROM UNNEST($1::text[], $2::json[]) AS "_subs"("result_id", "result_vars") LEFT OUTER JOIN LATERAL ( SELECT (... inner query with values from _subs.result_vars ...) ) AS "_fld_resp" ON ('true')内层查询通过 JSON path 运算符读取每个订阅者的会话变量与游标状态:
("_subs"."result_vars" #>> '{session,x-hasura-user-id}')::uuid关键设计点(源码注释见 multiplexed.go):
- GraphQL
$variable占位符保持编号参数($3、$4……),因为它们在 cohort 内完全相同;只有会话变量与游标状态按订阅者变化; - 会话变量识别纯按标记类型(
core.SessionVarValue、core.CursorValue、core.FunctionSessionArgument),绝不嗅探参数字符串值——用户提供的恰好以x-hasura-开头的字面量仍是静态参数,与 Hasura 的结构化追踪语义一致; buildResultVarsJSON为每个订阅者打包{"session": {...}, "cursor": {...}}JSON 对象,PrepareParams产出[subIDs, resultVarsJSON]两个固定参数,静态 GraphQL 参数由调用方追加在其后;- 边界情况显式报错而非静默出错:会话变量标记若被困在多元素数组参数内(如
_in: ["x-hasura-user-id", "<literal>"]),无法改写为单个result_vars查找,直接拒绝该订阅(ErrSessionVarInMultiElementArray),避免把字面量名当会话值求值导致错误/空结果集。
多路复用目前要求UNNEST+ LATERAL +::type[],全部是 Postgres 特性。Postgres 侧的执行入口是 connector/sql/postgres/postgres.go 的ExecuteMultiplexedOperation(pgx 查询并逐行扫描成{SubscriptionID, Data}对);SQLite handler 采用不同的多路复用方案,stream 订阅在 SQLite 上也可用但 cohort SQL 形态不同。
8. 拆除(Tear-down):四种触发路径
| 触发 | 路径 |
|---|---|
客户端发送complete | webSocketHandler.OnComplete→subscriptionState.handler.Stop(ctx, id) |
| 客户端关闭连接 | OnClose→ 遍历全部 subs,逐个调用Stop |
| 元数据重载 | controllerState.shutdown关闭state.done,每个 per-state handler 以 30s 预算执行Shutdown(ctx) |
| 服务关闭 | 与重载相同,使用进程退出上下文 |
当一个 cohort 的最后一名订阅者离开后,下一轮 poll-tick 在 manager 锁内观察到c.isEmpty()(cohort_manager.go),删除 cohort 条目并关闭其 stop 通道。这个 TOCTOU 模式是刻意的:在 manager 锁内检查空性,防止与addSubscription产生「先判空后新增」的竞态。
forwardUpdates在任一条件满足时干净退出:其订阅的stopCh关闭、连接请求上下文取消、或 handler 关闭更新通道。
9. 关键文件速查表
| 文件 | 用途 |
|---|---|
| services/constellation/controller/handlers.go | HandlerWebsocket入口(检测Upgrade: websocket头并创建连接) |
| services/constellation/controller/websocket/ | 纯协议层(读/写泵、帧封装),架构见 doc.go |
| services/constellation/controller/websocket.go | 每连接桥接、webSocketHandler、会话提取、订阅注册表 |
| services/constellation/subscription/types.go | Handler接口、Request、Update及错误哨兵 |
| services/constellation/connector/sql/subscription/handler.go | SQL connector 的Handler、stream/实时路由、QueryExecutor/QueryBuilder接口 |
| services/constellation/connector/sql/subscription/cohort_manager.go | 实时查询 cohort 生命周期、轮询循环、结果分发、SQL 缓存 |
| services/constellation/connector/sql/subscription/cohort.go | cohortKey、cohort、cohortSubscription、背压与变更检测 |
| services/constellation/connector/sql/subscription/stream_cohort_manager.go | stream cohort、每轮重建、游标提取与合并 |
| services/constellation/connector/sql/subscription/stream_cohort.go | streamCohortKey、游标推进时的 cohort 合并 |
| services/constellation/connector/sql/subscription/doc.go | 包级架构图与多路复用 SQL 深挖 |
| services/constellation/connector/sql/graphql/queries/multiplexed/multiplexed.go | SQL 改写:UNNEST+ JSON path 运算符 |
| services/constellation/connector/sql/postgres/postgres.go | pgx 的ExecuteMultiplexedOperation |
| internal/lib/syncmap/syncmap.go | 类型化并发 map,webSocketHandler.subs所用(位于仓库根,不在 constellation 的internal/下) |
延伸阅读
- connector/sql/subscription/doc.go——包级架构图与多路复用 SQL 的进一步说明;
- controller/websocket/doc.go——纯协议包与其 goroutine 模型;
- architecture.md——每连接状态快照如何与元数据重载联动;
- 想了解 HTTP 查询(
Resolve)路径与订阅路径的异同,可对照 query-execution.md 阅读。
小结:Constellation 的订阅设计把「共享」做到极致——cohort 键把可批量的订阅聚成一类,多路复用 SQL 把一类订阅压成一次查询,轮询循环与任何单一订阅者解耦,latest-wins 背压保证慢消费者不错过最新状态,游标哈希则让 stream 订阅在数据推进中自然合流。这套分层(纯协议层 → 连接桥接 → Handler 接缝 → 双 manager → 驱动多路复用)既保证了协议与实现的解耦,也为未来接入 CDC、消息总线等推送策略预留了干净的扩展点。
【免费下载链接】nhostThe Open Source Firebase Alternative with GraphQL.项目地址: https://gitcode.com/GitHub_Trending/nh/nhost
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考