
Hasura live queries 架构解析如何将 GraphQL 订阅扩展到百万级活跃连接【免费下载链接】graphql-engineBlazing fast, instant realtime GraphQL APIs on all your data with fine grained access control, also trigger webhooks on database events.项目地址: https://gitcode.com/gh_mirrors/gr/graphql-engine本篇文章基于 architecture/live-queries.md深入解析 Hasura GraphQL Engine 中 live queries实时查询的架构设计从每个客户端订阅一份查询结果的朴素思路出发逐步拆解 Hasura 如何通过GraphQL 编译为单条 SQL、声明式授权、多查询多路复用三大核心思路配合间隔轮询的刷新策略在单个 Postgres 之上承载百万级活跃 GraphQL 订阅。读完本文你将掌握 live queries 的完整实现链路、关键配置参数batch_size、refetch interval以及对应的服务端源码位置可以直接在 server/src-lib/Hasura/GraphQL/Execute/Subscription 目录下对照阅读。TL;DR一次百万级实时订阅的压测测试环境设定每个客户端Web/移动应用携带一个 auth token 订阅数据实时结果数据存放在 Postgres 中测试过程中每秒在 Postgres 中更新 1 百万行数据确保每个客户端都能收到一条新结果推送。Hasura 作为带授权能力的 GraphQL API 提供方。测试目标Hasura 能同时承载多少个并发实时订阅客户端Hasura 能否纵向垂直与横向水平扩展文档记录的单实例纵向扩展结果如下单实例配置活跃 live query 数CPU 平均负载1xCPU, 2GB RAM500060%2xCPU, 4GB RAM1000073%4xCPU, 8GB RAM2000090%横向扩展当扩展至 1 百万个活跃 live queries 时Postgres 负载约为 28%峰值连接数约 850 个。测试配置说明均为默认配置、未做任何调优AWS RDS Postgres16xCPU、64GB RAMPostgres 11Hasura 运行于 AWS Fargate每实例 4xCPU、8GB RAM使用默认配置负载均衡使用 AWS ELB 默认配置。这套数字说明live queries 的瓶颈不在连接数而在于如何把为每个客户端独立刷新的高昂代价降下来——这正是下文架构要解决的问题。GraphQL 与 subscriptionsGraphQL 让应用开发者可以精确地向 API 查询自己需要的数据。以文档中的外卖配送应用为例Postgres 中的模式大致如下orders订单、order_status订单状态、delivery_agent配送员等表应用的一个订单状态界面需要为当前用户拉取最新订单状态与配送员位置。对应的 GraphQL 查询会返回订单状态和配送员位置数据。在底层这个查询以字符串形式发送到 Web 服务器服务器解析它、应用授权规则并调用数据库等后端获取数据最后以请求时指定的形状JSON把数据返回给应用。Live queries 的定义订阅某条特定查询的最新结果。当底层数据发生变化时服务器应把最新结果主动推送给客户端。这与 GraphQL 天然契合GraphQL 客户端原生支持 subscriptions能够处理繁琐的 websocket 连接。把一条查询变成实时查询在客户端侧看起来就像把关键字query换成subscription那么简单——前提是 GraphQL 服务器能实现它。实现 GraphQL live-queries为什么难实现 live queries 是痛苦的。难点主要有二增量计算的困难一旦拿到包含全部授权规则的数据库查询理论上可以随着事件发生增量地计算结果但在 Web 服务层做这件事在实践上极具挑战。对于 Postgres 这类数据库这等价于随着底层表变化持续维护一个物化视图的难题。可扩展的 websocket 服务器构建一个可水平扩展的 websocket 服务器本身也不轻松不过某些框架和语言能让所需的并发编程更容易一些。Hasura 当前采取的替代方案是针对特定查询携带特定客户端的授权规则重新抓取refetch全部数据。本节的剩余部分解释为什么重新抓取也不简单以及 Hasura 如何把它做到极致。重新抓取一条 GraphQL 查询为什么困难看一条 GraphQL 查询典型的处理过程授权 数据抓取逻辑必须对查询中的**每一个节点**执行。这意味着即使一条稍大的查询比如抓取一个集合也能轻易拖垮数据库N1 查询问题与实现糟糕的 ORM 一样N1 问题对数据库不友好且难以对 Postgres 做最优查询。Data loader 之类的模式能缓解问题但仍会多次查询底层 Postgres把响应中的条目数退化成了查询中的节点数级别的多次查询。live queries 让问题更糟每个客户端的查询都会变成一次独立的重新抓取。即使查询相同由于授权规则产生不同的 session 变量也必须为每个客户端独立抓取。Hasura 的三大思路能否做得更好答案是利用数据模型到 GraphQL schema 的声明式映射为数据库生成一条SQL 查询。这样无论响应中的条目数多少、GraphQL 查询节点数多少都能避免对数据库的多次访问。思路一把一条 GraphQL 查询编译成一条 SQLHasura 内部有一个transpiler转译器它借助数据模型 → GraphQL schema的映射元数据把 GraphQL 查询编译成从数据库取数的 SQL 查询GraphQL query → GraphQL AST → SQL AST → SQL这消灭了 N1 查询问题并让数据库能够看到完整查询从而优化数据抓取。但这还不够——resolver 还要通过只抓取被允许的数据来执行授权规则因此需要把授权规则嵌入到生成的 SQL 中。思路二让授权声明式化数据访问层面的授权本质上是一个约束它同时依赖被抓取数据行的值以及应用用户特有的、动态提供的session 变量。例如最朴素的情况一行里含有user_id表示数据归属或者用户可见的文档记录在关联表document_viewers中再或者session 变量本身携带数据归属信息例如一个账户管理员可以访问账户 [1,2,3…]该信息不在当前数据库中而在 session 变量里——通常由其他数据系统提供。为此Hasura 在 API 层实现了一个类似 Postgres RLS 的授权层提供声明式的访问控制框架。类比一下SQL 查询中的 current session 变量变成了来自 cookie、JWT 或 HTTP header 的 HTTP session-variables。文档中提到一个趣事由于 Hasura 的工程起步很早团队在 Postgres 正式支持 RLS 之前就在应用层实现了 Postgres RLS 的特性——甚至踩过与 Postgres RLS 修复的INSERT ... RETURNING子句相同的 bug。把授权实现在 Hasura 应用层而不是使用 RLS 并通过 Postgres 连接的 current session settings 传递用户信息带来了显著收益稍后就会看到。总结授权现在是声明式的并且可以施加在表、视图、甚至函数若函数返回SETOF级别因此可以构造一条内嵌授权规则的单一 SQL 查询GraphQL query → GraphQL AST → Internal AST with authorization rules → SQL AST → SQL思路三把多条 live queries 批量合进一条 SQL如果只实现思路一和思路二那么 10 万个已连接的客户端仍可能导致约 10 万条 Postgres 查询的负载假设发生 10 万次更新、每次更新与一个客户端相关。但既然 API 层持有所有应用用户级别的 session 变量就可以用一条 SQL 查询同时为一批客户端重新抓取数据例如客户端们订阅最新订单状态 配送员位置。可以在查询内部创建一个关系relation把各客户端的查询变量如 order IDs和 session 变量如 user IDs作为不同的行放进这个关系再把这个关系与取数查询join从而在一次响应中拿到多个客户端的最新数据。响应中的每一行就是每个用户的最终结果。这样即使参数和 session 变量完全动态、只在查询时才知道也能同时为多个用户抓取最新结果。什么时候重新抓取Hasura 团队实验过多种从底层 Postgres 捕获事件、以决定何时重新抓取查询的方法Listen/Notify需要给所有表都装上触发器消费者Web 服务器在重启或网络中断时可能丢失事件。WAL可靠的流但 LR slots 代价高昂使水平扩展变难且托管数据库厂商往往不提供高写入负载会污染 WAL需要在应用层做节流。经过这些实验后当前方案回退为基于时间间隔的轮询interval based polling不是有事件就重抓而是按时间间隔重抓查询。两个主要原因把数据库事件映射到特定客户端的 live query在声明式权限和 live query 条件平凡时如order_id 1 且user_id cookie.session_id勉强可行一旦复杂化例如查询使用status ILIKE failed_%就变得不可处理声明式权限有时还横跨多张表。团队在该方向配合基础增量更新投入了大量调研并有几个生产小项目采用了类似 Postgres 增量物化视图维护的思路。对任何应用除非写入吞吐量极小否则最终都会在某个时间间隔上做节流/防抖debounce事件。该方案的代价是写入负载小时存在延迟——本可以立即重抓却要等 X 毫秒。这个代价可以通过合理调节重抓间隔与批大小轻松缓解。目前团队优先移除了最昂贵的瓶颈——查询重抓未来会继续改进特别是在适用场景下用事件依赖来减少每个时间间隔内被重抓的 live query 数量。文档还提到 Hasura 内部有基于事件方法的其他驱动若现有方案不满足特定场景可以联系团队跑基准测试。源码中的实现证据poller、cohort 与响应哈希在 server/src-lib/Hasura/GraphQL/Execute/Subscription 目录下可以找到与上述架构一一对应的实现。轮询间隔与批大小默认值与配置项Options.hs 中定义了SubscriptionsOptionslive queries 与 streaming queries 复用同一类型batch_size默认100refetch_delay轮询间隔默认1秒。对应的 CLI 参数与环境变量定义在 Serve.hsCLI 参数环境变量默认值说明--live-queries-multiplexed-refetch-intervalHASURA_GRAPHQL_LIVE_QUERIES_MULTIPLEXED_REFETCH_INTERVAL1000ms可多路复用的 live queries 在此间隔内至多推送一次结果--live-queries-multiplexed-batch-sizeHASURA_GRAPHQL_LIVE_QUERIES_MULTIPLEXED_BATCH_SIZE100多路复用的 live queries 按此大小拆分批次--streaming-queries-multiplexed-refetch-intervalHASURA_GRAPHQL_STREAMING_QUERIES_MULTIPLEXED_REFETCH_INTERVAL1000msstreaming queries 的多路复用推送间隔--streaming-queries-multiplexed-batch-sizeHASURA_GRAPHQL_STREAMING_QUERIES_MULTIPLEXED_BATCH_SIZE100streaming queries 的多路复用批大小这就是文档中通过调节重抓间隔和批大小来权衡延迟的具体旋钮。轮询线程pollLiveQueryPoll/LiveQuery.hs 中的pollLiveQuery是每个 poller 周期性执行的核心动作与文档中的三思路对应快照 cohort 并按 batch size 切分从 cohort map 取出所有 cohort拥有相同查询、相同 session/query 变量的订阅者分组用chunksOf按 batch size 切成批次BatchId标识并发执行批查询每个批次调用runDBSubscription一次执行数据库查询——这正是把多条 live queries 合进一条 SQL的实现按响应哈希决定是否推送每个 cohort 都保存上一次结果的哈希respRef只有执行出错或本次结果哈希与上次不同时才把新结果推给订阅者否则跳过——避免把未变化的数据反复推给客户端显著降低 WebSocket 推送开销。Cohort多路复用的基本单元Poll/Common.hs 对Cohort的注释直接印证了文档的描述A batched group of Subscribers who are not only listening to the same query but also have identical session and query variables. … In SQL, each Cohort corresponds to a single row in the laterally-joined _subs table (and therefore a single row in the query result).即每个 cohort 对应生成 SQL 中横向连接lateral join的_subs表里的一行——这就是把各客户端的查询变量与 session 变量作为不同行放进一个关系再 join的落地。同时Cohort区分_cExistingSubscribers已推送过、仅在结果变化时推送与_cNewSubscribers新订阅者无论结果是否变化都立即推送配合ResponseHashBlake2b_256 结果哈希用于判断结果是否变化构成了推送去重的完整机制。Poller 生命周期State.hs 中每个 poller 线程以forkImmortal启动循环执行pollLiveQuery b ...之后sleep $ unrefine $ unRefetchInterval refetchInterval——即按配置的 refetch interval 休眠后进入下一轮轮询这正是文档所述基于时间间隔的轮询的直接实现。此外 Prometheus.hs 中暴露了submActiveLiveQueryPollersInError等指标pollLiveQuery中会在批次查询出错/恢复时切换 poller 的PollerResponseState并记录对应监控指标方便运维观测轮询健康状态。测试如何验证百万级实时订阅的可靠性与可扩展性文档指出用 websocket 做 live queries 的可扩展性与可靠性测试本身就是一个挑战测试套件与基础设施自动化工具花了几周时间搭建。测试架构如下客户端模拟一个 Node.js 脚本运行大量 GraphQL live-query 客户端把事件日志记录在内存中之后摄入数据库对应公开的subscription-benchmark项目写入负载一个脚本对数据库制造写负载使变化覆盖所有运行 live query 的客户端每秒更新 1 百万行结果校验测试套件运行完毕后在校验脚本对摄入日志/事件的数据库执行检查确认无错误且所有事件都被收到测试有效性标准收到的 payload 中0 错误从事件产生到客户端收到结果的平均延迟小于 1000ms。正是这套严格标准支撑了开头 TL;DR 中 1 百万活跃 live queries 的基准数字。这套方案的收益Hasura 让 live queries 变得简单易得查询的概念可以零额外成本地扩展为实时查询开发者无需为此学习新的东西——这是团队最看重的一点。具体收益包括表达力强、功能完整的 live queries完整支持 Postgres 的运算符、聚合、视图、函数等可预测的性能查询被编译为单条 SQL 并多路复用负载可预估纵向与横向扩展单实例可从 50001xCPU扩展到 200004xCPU活跃订阅多实例可线性扩展到百万级适用于所有云/数据库厂商不依赖 WAL、Listen/Notify 等厂商差异大的能力只依赖标准 Postgres。未来工作文档列出的未来优化方向核心是进一步降低 Postgres 负载把事件映射到活跃 live queries在适用的场景下用事件依赖减少每个轮询间隔内被重抓的查询数量结果集的增量计算向随事件增量维护结果的方向演进减少全量重抓。这两点与 architecture/streaming-subscriptions.md 描述的 streaming subscriptions流式订阅方向互补live queries 适合不需要持续数据/事件流的场景如当前在线用户集合、某个聚合的最新值而 streaming subscriptions 则是面向持续消费大结果集或事件流的另一种订阅字段。两者共享同一套多路复用、声明式授权与轮询基础设施。延伸阅读architecture/streaming-subscriptions.md流式订阅的架构、批处理多路复用与性能基准architecture/sql-server.mdSQL Server 后端的接入架构server/src-lib/Hasura/GraphQL/Execute/Subscription/Options.hs订阅相关配置的默认值与 JSON 序列化server/src-lib/Hasura/GraphQL/Execute/Subscription/Poll/LiveQuery.hs多路复用轮询线程的核心实现server/src-lib/Hasura/GraphQL/Execute/Subscription/Poll/Common.hsCohort 与响应哈希的定义translations/live-queries.chinese.md本文档的中文翻译版本。【免费下载链接】graphql-engineBlazing fast, instant realtime GraphQL APIs on all your data with fine grained access control, also trigger webhooks on database events.项目地址: https://gitcode.com/gh_mirrors/gr/graphql-engine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考