十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

PostHog SQLV2 节点结果交付机制解析:从短轮询到 pub/sub 推送与结果分页存储设计

PostHog SQLV2 节点结果交付机制解析:从短轮询到 pub/sub 推送与结果分页存储设计 PostHog SQLV2 节点结果交付机制解析从短轮询到 pub/sub 推送与结果分页存储设计【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthogSQLV2 是 PostHog 新版 Notebooknotebooks中 SQL 节点的执行引擎本文以其后端设计文档 sql_v2_result_delivery.md 为主体结合products/notebooks/backend/下的真实源码完整讲解一条 SQLV2 运行run从执行完成到结果回传前端的全过程当前已实现的前端短轮询 持久化行方案、中断interrupt如何保持单一事实来源、未来可选的 Redis pub/sub 推送设计以及结果存储与分页的三层架构。读完本文你将理解为何 SQLV2 选择轮询而非 SSE、直接通道direct lane与内核通道kernel lane的差异、result_id指针的妙用以及NotebookNodeRun一行对应一次执行的模型边界。核心问题一条 run 如何知道自己跑完了SQLV2 节点的执行结果被持久化写入NotebookNodeRunPostgres前端节点通过一个廉价读端点轮询直到 run 进入终态。文档开篇即点明结果是一个单一终值single terminal value而非实时事件流——这是整个交付机制设计取舍的总纲。从 models.py 可以看到NotebookNodeRun的状态机定义running、done、failed以及用户主动停止产生的interrupted与 FAILED 不同interrupted 的 envelope 中捕获的 stdout/stderr 仍然会呈现给 UI。结果如何写入NotebookNodeRun取决于 run 走哪条通道lane通道适用场景结果写入方式内核通道kernel lanepython / duckdb 节点sandbox 回调callback写入 envelope直接通道direct lanehogql 节点无 sandboxsql_v2_direct.py结果轮询自身推进状态行当前方案前端短轮询已实现后端一个索引查询返回状态后端暴露GET .../sql_v2/runs/run_id返回{status, result, error}直接通道还会附带临时的rows。实现上就是一次索引查询读取NotebookNodeRun行直接通道再多一次 Redis 状态读取不持有连接、不忙等。直接通道direct lane异步查询管理器 受保护的转移这是与旧方案差异最大的部分。纯 HogQL runnode_type hogql不需要内核其内联查询直接复用数据面data plane所用的异步查询管理器async query manager执行。但异步查询管理器没有完成回调因此由结果轮询本身推进状态行sync_direct_run轮询按query_id run_id读取查询状态并应用受保护的状态转移filter(statusRUNNING).update(...)。这在并发轮询下是幂等的且已完成的查询永远无法覆盖一次中断。查询 id 并非 run id 本身而是由 SECRET_KEY 派生的 HMACnotebook_direct_query_idrun id 是客户端可见的出现在 notebook 文档与查询日志中直接用作 query id 会把 run 的 rows 暴露进共享命名空间绕过 notebook 与按用户的仓库检查派生一个未发布的 id 让轮询无需存储任何额外状态即可重算。一个 RUNNING 的 run 若查询状态消失管理器状态 TTL 为 20 分钟会在宽限窗口后判为失败——这就是内核通道一直缺失的看门狗。对应常量DIRECT_RUN_RESULT_GRACE_SECONDS 600sql_v2_direct.py。管理器结果存活期间轮询响应还携带rows完整封顶集合≤300 行供客户端本地分页。这里也体现了查询的行数封顶策略RESULT_CACHE_ROWS 300DISPLAY_PAGE_LIMIT 50sql_v2.py。enqueue 时多取一行3001用于探测has_more镜像内核的封顶抓取enqueue_direct_runenvelope 只携带前 50 行作为预览。对 HogQL 查询还做 LIMIT pushdown 优化apply_page_bounds无自带 LIMIT 的查询直接在源文本上追加limit N让 ClickHouse 把限制下推到聚合视图如persons避免全表去重无法解析或带 OFFSET 的查询则回退到外层包装select * from (...) as posthog_notebook_page limit/offset。前端runId running 状态驱动的轮询定时器前端节点在持有runId且 run 处于running时以约 1s 间隔轮询该端点done/failed/interrupted时停止。进行中状态由 run 状态派生因此跨重挂载、刷新都能恢复。轮询定时器存放在cache.disposables自动清理 隐藏标签页时暂停。中断Journey 9两条通道、一个事实来源内核通道POST .../sql_v2/runs/run_id/interrupt将 run 级停止代理给 kernel-serverinterrupt_sql_v2_run终态interrupted仍经由 回调 → run 行 → 轮询 到达保持单一事实来源。被中断的 envelope 携带停止前捕获的 stdout/stderr结果端点像处理doneenvelope 一样呈现它。当内核不可达时端点自行将 run 标记为interrupted——用户总能从永远 RUNNING的行中脱身内核通道没有后端看门狗。直接通道hogql没有内核可发信号端点自行将行标记为 abandoned受保护的轮询转移无法覆盖中断然后在异步管理器侧停止查询cancel_query 若查询仍在排队则撤销 Celery 任务否则在 ClickHouse 上 KILL 该 run 派生 query id 对应的查询。该取消是尽力而为best effort因为行已是终态拒绝 KILL 的集群只会让查询运行到自身有界完成、结果被丢弃。两个通道都提供 Stop。选型的核心理由结果是单一终值轮询持久化行比 SSE 更简单且天然免疫连接断开、刷新、节点重挂载这正是此前踩过的失败模式代价是每次 run 只有寥寥几次 1-查询读。未来方案Redis pub/sub 推送降低延迟时再实现当轮询延迟约 1s变得明显或大量并发 run 使轮询变得浪费时可叠加 Redis pub/sub保持 Postgres 持久化写入作为事实来源新增 pub/sub 通道让结果在回调落地的瞬间被推送。设计要点回调处理器sql_v2_callback.py保留 DB 写入envelope →NotebookNodeRun为权威来源写入后尽力PUBLISH notebook:sql_v2:run:run_id一个小信号如{status: done}。payload 保持微小——SSE 处理器会从 Postgres 重读权威 envelope。发布失败绝不丢数据由连接时读取兜底。SSE 流处理器用先订阅再检查subscribe-then-check取代忙轮询先订阅notebook:sql_v2:run:run_id在 DB 读之前——这消除了回调恰好在读与订阅之间触发的竞态读取一次 run若已是done/failed则发出结果/错误、退订、返回覆盖在我连接前就已结束否则在有界超时 周期性 SSE 心跳下阻塞订阅收到消息 → 重读 run → 发出 → 返回超时 → 报错。无轮询循环、无time.sleep。Redis async复用现有异步 Redis 客户端处理器必须是async的使被持有的 SSE 连接是一个空闲协程而非工作线程。前端节点保留其持久化 run 恢复逻辑挂载时result缺失则按持久化runIdfetch/poll作为回退路径——SSE 推送是快路径连接时的 DB 读是安全网。边缘情况清单订阅先于读的顺序唯一要紧的竞态、Postgres 权威而 pub/sub 只是唤醒漏消息回退到连接时读绝无数据丢失、超时 心跳连接不悬挂、代理不杀空闲流、所有退出路径结果/错误/超时/断连都要退订清理。为何不用 Redis Stream不像 PostHog DesktopPostHog Desktop 流式传输大量增量 agent 事件重连时需要重放基于裁剪过的流上的Last-Event-ID游标。SQLV2 只交付一个终态结果pub/sub 持久化行已足够且简单得多。仅当 SQLV2 开始流式传输增量输出stdout、逐行到达的 rows、进度时才需要重新考虑 Stream 方案。结果存储与分页设计尚未实现三层架构——切勿混为一谈层存放位置大小NotebookNodeRun.envelopePostgres JSONFieldrun 行上小——元数据 first_page预览完整结果所有页沙箱在运行时填充的result store可能很大NotebookNodeRun.result_idPostgresUUID 列取自 envelope极小——指向 store 的句柄envelope 携带{columns, row_count, first_page, result_id}——有界的预览加一个指针永不持有超过第一页的内容。完整数据集存于 result store由result_id引用。这就是大数据按引用传递而非按值传递规则Temporal payload 封顶同理——一个多 MB 的 dataframe 绝不能落进 JSONField 或回调体。一次 run 一行而不是每页一行NotebookNodeRun是每次执行每次 Run 点击一行产生一个result_id与一个存储的结果集。取页不是 run——它是对已物化结果的读因此不创建新的NotebookNodeRun、不启动 Temporal workflow、不产生 DB 写入。run 保持审计/历史记录的角色页读取保持无状态1 次 Run 点击 → 1 个 NotebookNodeRun → 1 个 result_id → 1 个存储的结果集 ↑ page 1, 2, 3, download all ──只读────┘ (0 个新 run)完整结果存放在哪里现在hogql run 客户端侧。直接 run 在一次 ClickHouse 查询中抓取最多RESULT_CACHE_ROWS300行完整封顶集合随结果轮询rows在管理器 Redis 结果存活期间约 20 分钟传递浏览器本地分页——服务端 page 端点拒绝hogql run。过期或刷新后只剩 envelope 的first_page配一个重跑入口——与旧版 SQL 编辑器同模型一次钳制抓取、客户端分页。kernel-server 的按 run hogql 缓存与/page重查询仍存在但只服务于 direct lane 之前的旧 run二者都计划移除。内核 runpython/duckdb从沙箱内 result store 分页。内核将每个产出的 frame 写入/data/results/result_id.arrow携带result_id的/page请求在服务器进程内切片该文件result_store.py 使用 pyarrow mmappa.memory_mappa.ipc.open_file切片table.slice(offset, limit)。没有数据面回退——run 的代码不是 HogQL 查询——因此丢失 frame沙箱死亡意味着重跑即下述仅存活时可用的取舍。安全细节result_id先经uuid.UUID校验再拼接路径把路径 join 限定在 UUID 文件名内。未来持久化 store对象存储 / Parquet / 结果表。沙箱在运行时物化完整结果result_id作为存储键消除仅存活时可用的限制。当结果必须在内核拆除/刷新后存活时切换到该方案。内核驻留 store 的取页数据在沙箱内取页必须往返沙箱——但不能使用 run 的异步回调。回调存在是因为 run 的延迟无界从内存 frame 切片行则快速且有界所以分页是普通的同步请求/响应Run (无界): FE → POST /sql_v2/run → Temporal → kernel-server POST /run → 202 (稍后) → POST /callback → DB → FE 轮询 Page (有界): FE → GET /sql_v2/results/result_id?page2page_size100 → 后端找到运行中的内核 → HTTP POST kernel-server /page {result_id, page, page_size} → kernel-server 切片驻留 frame在 200 响应中返回 rows → 后端将 rows 返回 FE无需 docker 控制面。run 走 Temporal 部分是因为 ensure_sql_v2_server 通过write_file/executedocker socket引导服务器。到取页时 kernel-server 已启动分页只是对其的一次网络 HTTP 调用——即使在 dev 的 Seatbelt 沙箱下也允许只有 dockersocket被拒绝网络出口不受限。因此同步的 web→kernel-server 调用是合理且低延迟的。对应常量可见 sql_v2.py_PAGE_POST_TIMEOUT_SECONDS 60以及按用户的页取锁PAGE_LOCK_TTL_SECONDS 70一个用户同时只有一个进行中的取页。待构建在 kernel-server 加同步/page路由server.py 的do_POST已处理/page由 fetch_page 支撑 一个代理到它的薄后端读端点。服务端 fetch_sql_v2_page 的现有实现区分了两种取页路径hogql 走code重查询数据面内核 run 走result_id切片。内核驻留的取舍仅存活时可用result_id切片只在对应内核存活期间存在。若内核空闲超时、重启或 notebook 重载到新内核result_id过期 →/page返回 not-found → UI 必须重跑以获得新的result_id。持久化存储正是以后消除这一限制的手段。相关模型注释KernelRuntime.server_url/server_connect_token以明文TextField存储与 PostHog Desktop 将sandbox_url/sandbox_connect_token明文存于TaskRun.state普通 JSONField一致。连接令牌是临时的 Modal 隧道令牌而非持久账户机密账户机密如环境变量tasks 会通过EncryptedJSONStringField加密。若需要纵深防御加密字段是既有模式。NotebookNodeRun对KernelRuntime没有外键也不需要分发定位当前运行的内核结果按run_id返回后续 run 只是在当时运行的内核上创建新行。若要为调试提供 run→沙箱可追溯性优先在 run 上存储sandbox_id字符串而非硬 FK——KernelRuntime行是瞬态的starting/stopped/discarded/errorFK 会把 run 历史耦合到频繁变动的行。源码中kernel_runtime_id models.UUIDField(nullTrue, blankTrue)models.py正是这一决策的落地普通 id 而非 FK回调据此把 run 的 frame 快照归档到正确的KernelRuntime。结语一条路径、两类通道、三层存储回顾全文SQLV2 结果交付的设计哲学清晰可辨持久化行是唯一事实来源一切传输机制轮询、未来的 pub/sub都只是它的读取或唤醒方式。直接通道用轮询推进受保护状态转移弥补了异步查询管理器无回调的缺陷并顺带为内核通道补上了看门狗中断在两条通道上都以行先终态、查询尽力取消的方式收敛回同一事实来源。而结果存储的envelope 预览 result_id 句柄分层则为大结果集的物化与分页预留了从沙箱内存到对象存储的演进路径——在结果需要跨内核存活之前沙箱驻留的 Arrow 文件加同步/page切片已经是延迟与复杂度之间的最优解。【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表