)
OmO 原生遥测本地技能查询库扩展实战基于 HogQL 构建工具并行度汇总查询parallelism / parallelism_daily【免费下载链接】oh-my-openagentOmO: Just type mass ulw keyword with your prompt. Now you are the master of graph engineering.项目地址: https://gitcode.com/gh_mirrors/oh/oh-my-openagent本文以 OmOoh-my-openagent原生遥测技能开发中的 Task 7 证据文档为主体完整讲解如何在本地技能脚本fetch_data.py的QUERIES查询库中新增parallelism与parallelism_daily两条 HogQL 聚合查询用于按会话汇总工具调用并行化native tool-call parallelism带来的往返次数节省与墙钟延迟节省。读完本文你将掌握如何以纯增量 diff 扩展 PostHog/HogQL 查询库、如何用--only定向运行与 JSON payload 校验验证空结果语义、如何用 relaxed-filter 双层证明排除SQL 静默无效以及如何守住聚合语义合规eval 桶隔离、比率必须下游推导、上界指标降级为参考列等硬性约定。1. 背景parallelism_summary事件与并行度遥测管线在深入查询库扩展之前需要先理解被查询数据的来源。OmO 的原生遥测组件在每个会话关闭session_shutdown时恰好发出一次parallelism_summary事件其发送点在 omo-native-parallel-summary.tspi.on(session_shutdown, (_payload: unknown, eventContext: unknown): void { const sessionId extractSessionId(eventContext) if (sessionId undefined) return const snapshot registry.snapshot(sessionId) registry.clear(sessionId) if (snapshot undefined) return const properties buildParallelismSummary(snapshot, hashSessionId(sessionId)) if (properties undefined) return captureEvent(parallelism_summary, properties) })该文件头部注释明确说明了事件注册顺序是承重load-bearing的必须在会话组件之前注册且只在session_shutdown时发出以保证每会话一次的体量上界——在turn_end上发送会导致每轮都发体量爆炸只在首个 turn 发送则会静默丢弃后续轮次的 wave。事件消费一次后立即清空 registry因此同一会话的第二次 shutdown 不会重复发送两个会话也永远不会混在一起。事件的属性 schema 定义在 parallelism-schema.ts其中与本文两条查询直接相关的 13 个聚合列属性名必须与查询中properties.X字节级一致查询列属性名语义saved_round_tripsnon_eval_saved_round_trips非 eval 桶节省的往返次数modeled_saved_msmodeled_wallclock_saved_ms头条节省指标各调用时长之和减去 wave 墙钟跨度waves_totalnon_eval_waves_total非 eval 桶 wave 总数waves_multinon_eval_waves_multi并发 1 的非 eval wave 数joined_callsnon_eval_joined_calls合并进 wave 的非 eval 调用数eval_only_waveseval_only_waves纯 eval 桶 wave 数独立列mixed_wavesmixed_waves混合桶 wave 数独立列incomplete_callsincomplete_calls不完整调用数质量计数器clock_anomaliesclock_anomalies时钟异常数质量计数器dropped_callsdropped_calls被丢弃调用数质量计数器measured_turn_msmeasured_turn_duration_ms_total实测 turn 总时长upper_bound_saved_ms_refupper_bound_saved_ms上界参考值非估算仅作参考列1.1 三个节省指标的计算语义为什么modeled_wallclock_saved_ms是头条、upper_bound_saved_ms只能做参考底层算法在 savings-math.tsexport function modeledWallClockSavedMs(wave: MeasurableWave): ModeledSavedMs { const durations usableDurations(wave.calls) if (durations.length 1 || !Number.isFinite(wave.spanMs)) { return { label: modeled, valueMs: 0 } } return { label: modeled, valueMs: sum(durations) - wave.spanMs } } export function upperBoundSavedMs(wave: MeasurableWave): UpperBoundSavedMs { const durations usableDurations(wave.calls) if (durations.length 1) return { label: upper_bound, valueMs: 0 } const mean sum(durations) / durations.length return { label: upper_bound, valueMs: (durations.length - 1) * mean } }modeled模型化估算按 wave 跨度建模sum(durations) - spanMs。注释里给了反例——链式 waveA 0-5、B 4-9、C 8-12实际耗时 12msspan 公式只报 2ms 节省而max(duration)变体会报 9ms4.5 倍夸大真正同时发起的批次spanMs max(duration)数值不变。因此诚实批次不会被惩罚夸大路径被结构性排除。upper_bound上界(N - 1) * mean(duration)只是上界而非测量两个结果带有不同的字面量标签modeled/upper_bound调用方不经过显式类型转换无法混用。往返次数按maxConcurrency而非 wave 大小计算——链式 wave 有 3 个调用但最多 2 个并发只算节省 1 次往返。负值不钳制span 宽于时长总和说明观测互相矛盾隐藏它等于隐藏异常。1.2 事件覆盖与属性白名单的 QA 锚点parallelism_summary已进入原生事件覆盖集合与属性白名单仓库中的 omo-native-telemetry-assertions.mjs 将其与daily_active、turn_completed、delegation_started等并列为expectedNativeEvents任何拼写错误或删除都会导致 QA 失败同时OMO_NATIVE_PROPERTY_ALLOWLISTS白名单保证每个属性都是文档化或 SDK 新增的 key。这意味着本文两条查询里出现的properties.X名称是经过 schema 与白名单双重约束的。2. 查询库扩展纯增量 diff 添加两条查询Task 7 的目标文件位于仓库之外、不受版本控制~/.agents/skills/omo-native-telemetry/scripts/fetch_data.py。由于技能目录不是 git 仓库diff 按计划记录在本证据文件中对照/tmp/fetch_data.py.bak生成。仓库工作树内唯一的改动就是这份证据文件本身其余路径全部未动。2.1 新增的查询定义diff 的核心是在QUERIES字典末尾追加两个键 -58,6 58,16 # Marginal cost of one extra round trip: the median session-median turn among sessions that # actually delegate. Uses observed cost_usd, which already prices cache reads correctly. parallel_savings_basis: select quantile(0.5)(p50_prefix) med_prefix, quantile(0.5)(p50_cost) med_cost, count() sessions from (select properties.$session_id sid, quantile(0.5)(toFloat(properties.input_tokens)toFloat(properties.cache_read_tokens)toFloat(properties.cache_write_tokens)) p50_prefix, quantile(0.5)(toFloat(properties.cost_usd)) p50_cost from events where eventturn_completed and properties.$session_id in (select distinct properties.$session_id from events where eventdelegation_started) group by sid), # Native tool-call parallelism, one parallelism_summary per session. The non_eval_* properties # ALREADY exclude eval-tool waves, so eval_only_waves/mixed_waves are reported as SEPARATE # buckets and are never folded into the non_eval totals - one eval cell running many internal # operations would otherwise masquerade as fleet-wide parallelism. modeled_wallclock_saved_ms # (sum of call durations minus the waves wall-clock span) is the headline savings figure; # upper_bound_saved_ms is a BOUND, not an estimate, so it is carried only as a secondary # reference column. Fleet ratios must be derived downstream as sum(numerator)/sum(denominator) # - these queries deliberately emit raw sums only, never per-session averages or medians. parallelism: select count() sessions, sum(toFloat(properties.non_eval_saved_round_trips)) saved_round_trips, sum(toFloat(properties.modeled_wallclock_saved_ms)) modeled_saved_ms, sum(toFloat(properties.non_eval_waves_total)) waves_total, sum(toFloat(properties.non_eval_waves_multi)) waves_multi, sum(toFloat(properties.non_eval_joined_calls)) joined_calls, sum(toFloat(properties.eval_only_waves)) eval_only_waves, sum(toFloat(properties.mixed_waves)) mixed_waves, sum(toFloat(properties.incomplete_calls)) incomplete_calls, sum(toFloat(properties.clock_anomalies)) clock_anomalies, sum(toFloat(properties.dropped_calls)) dropped_calls, sum(toFloat(properties.measured_turn_duration_ms_total)) measured_turn_ms, sum(toFloat(properties.upper_bound_saved_ms)) upper_bound_saved_ms_ref from events where eventparallelism_summary, parallelism_daily: fselect toDate({KST}) d, count() sessions, sum(toFloat(properties.modeled_wallclock_saved_ms)) modeled_saved_ms from events where eventparallelism_summary group by d order by d, }parallelism返回 13 列聚合parallelism_daily则按 KST 自然日分组输出每日趋势且只携带头条节省指标modeled_wallclock_saved_ms。2.2 约定合规性Convention Compliance纯标准库diff 零新增 import。单行字符串两条查询都是QUERIES字典中的单行字符串与其他条目格式一致。固定查询通道通过既有的hogql()/fetch_one()路径查询 PostHog 项目552066。数值属性统一toFloat每个数值属性都套用toFloat(properties.X)。复用模块级KST常量parallelism_daily用 f-string 插值模块级KST常量toTimeZone(timestamp,Asia/Seoul)与dau/turns_daily/prompts_daily的写法完全一致而不是重新声明时区。注释风格对齐新增条目上方的大段注释解释了 eval 桶分离与upper_bound_saved_ms降级理由与既有parallel_savings的注释风格一致。2.3 语义合规性Semantic Complianceeval 桶隔离eval_only_waves与mixed_waves各自独立成列任何位置都没有 eval 列与non_eval_*列相加的运算符。这与发送端 buildParallelismSummary 的实现一致——节省只对non_eval桶求和mixedwave 只有自己的计数器因为把一个 eval 调用的跨度计入节省会让单一 eval 单元的内部多操作伪装成全舰队并行。头条指标固定modeled_saved_ms是parallelism的第二列也是parallelism_daily唯一的节省指标upper_bound_saved_ms仅出现一次、位于最后一列、显式命名为upper_bound_saved_ms_ref。SQL 中零比率只使用count()与sum(toFloat(...))没有任何avg、quantile、median或除法运算符。全舰队比率必须在下游推导为sum(分子)/sum(分母)杜绝avg(a/b)与按会话取中位数。属性名字节级一致全部 13 个聚合列与 parallelism schema见上表的属性名逐字节一致。3. 验证一语法检查与定向运行验收标准计划的验收标准是一条组合命令先做 AST 语法检查再用--only只跑两条新查询$ cd /Users/yeongyu/.agents/skills/omo-native-telemetry/scripts python3 -c import ast,sys; ast.parse(open(fetch_data.py).read()); print(syntax ok) python3 fetch_data.py /tmp/t7-qa --only parallelism,parallelism_daily; echo EXIT$? syntax ok parallelism: 1 rows parallelism_daily: 0 rows EXIT0退出码为 0两个 JSON 文件都已写出。--only这种定向开关让开发者可以在不动全量查询库的情况下单测新增条目是查询库迭代的关键工作流。3.1 JSON payload 的正确空答案语义$ cd /tmp/t7-qa for f in parallelism.json parallelism_daily.json; do echo --- $f; cat $f; echo; python3 -c import json,sys; json.load(open($f)); print(valid json); done --- parallelism.json [[0, null, null, null, null, null, null, null, null, null, null, null, null]] valid json --- parallelism_daily.json [] valid json今天是正确结果——parallelism_summary尚未部署因此查询天然为空parallelism.json非分组聚合永远恰好返回一行。sessions 0每个sum()都是nullClickHouse/HogQL 对空集合求和返回 null而不是伪造的0。这是诚实的空答案头部sessions0让下游卡片一目了然不可能被误读为真实会话节省了 0ms。parallelism_daily.json[]是真正的空序列因为GROUP BY d没有产生任何分组。两者都没有崩溃、没有伪造聚合、都是合法 JSON。4. 验证二全量查询库回归sibling-regression只验证新查询还不够必须确认没有破坏既有条目$ cd /Users/yeongyu/.agents/skills/omo-native-telemetry/scripts python3 fetch_data.py /tmp/t7-full /tmp/t7-full.log 21; echo EXIT$?; cat /tmp/t7-full.log; echo --- ERROR lines:; grep -c ERROR /tmp/t7-full.log EXIT0 cohort_0812: 360 rows user_features: 682 rows parallelism: 1 rows funnel_stages: 1 rows parallel_savings_basis: 1 rows parallel_cache: 558 rows depth_cache: 5 rows ulw_funnel: 1 rows skill_pairs_onboarding: 12 rows ulw: 6 rows models: 25 rows token_split: 1 rows parallel_savings: 30 rows deleg_funnel: 1 rows d1_retention_0812: 1 rows skill_diversity: 1 rows depth_ladder: 1 rows parallelism_daily: 0 rows hourly: 100 rows token_pcts: 1 rows turn_depth: 1 rows ordinal: 5 rows prompt_len: 4 rows prompts_daily: 6 rows per_user: 653 rows queue_mode: 4 rows delegation: 15 rows sessions_with_delegation: 1 rows version_daily: 19 rows headline: 1 rows new_users_daily: 6 rows skills: 22 rows prompts_total: 1 rows features: 3 rows platform: 4 rows batch_size: 3 rows hour_profile: 24 rows turns_daily: 6 rows versions: 5 rows dau: 6 rows delegation_bg: 2 rows --- ERROR lines: 0fetch_data.py退出码为 041 条查询的日志中 ERROR 行为 0grep -c输出 0 时自身退出码为 1这是期望行为脚本自身的EXIT0已在其上方打印。尤其重要的是parallel_savings30 行与parallel_savings_basis1 行仍返回既有形态没有兄弟查询被破坏——这正是 §2.1 中 diff 纯增量的结果QUERIES字典中没有任何条目被删除、重排或修改。5. 验证三relaxed-filter 双层证明——SQL 有效而非静默为空空结果集本身无法区分两种情形查询有效但无数据与查询因错误而匹配不到任何东西。Task 7 用两级递进证明排除了后者。5.1 结构证明只把事件过滤器换成真实存在的事件$ cd /Users/yeongyu/.agents/skills/omo-native-telemetry/scripts python3 - PY import fetch_data as fd tok fd.read_token() for name in (parallelism, parallelism_daily): q fd.QUERIES[name].replace(eventparallelism_summary, eventturn_completed) print(f--- RELAXED {name}\nQUERY: {q}) rows fd.hogql(tok, q) print(fROWS: {len(rows)}) print(rows[:8]) PY --- RELAXED parallelism QUERY: select count() sessions, sum(toFloat(properties.non_eval_saved_round_trips)) saved_round_trips, sum(toFloat(properties.modeled_wallclock_saved_ms)) modeled_saved_ms, sum(toFloat(properties.non_eval_waves_total)) waves_total, sum(toFloat(properties.non_eval_waves_multi)) waves_multi, sum(toFloat(properties.non_eval_joined_calls)) joined_calls, sum(toFloat(properties.eval_only_waves)) eval_only_waves, sum(toFloat(properties.mixed_waves)) mixed_waves, sum(toFloat(properties.incomplete_calls)) incomplete_calls, sum(toFloat(properties.clock_anomalies)) clock_anomalies, sum(toFloat(properties.dropped_calls)) dropped_calls, sum(toFloat(properties.measured_turn_duration_ms_total)) measured_turn_ms, sum(toFloat(properties.upper_bound_saved_ms)) upper_bound_saved_ms_ref from events where eventturn_completed ROWS: 1 [[1633540, None, None, None, None, None, None, None, None, None, None, None, None]] --- RELAXED parallelism_daily QUERY: select toDate(toTimeZone(timestamp,Asia/Seoul)) d, count() sessions, sum(toFloat(properties.modeled_wallclock_saved_ms)) modeled_saved_ms from events where eventturn_completed group by d order by d ROWS: 6 [[2026-08-11, 1666, None], [2026-08-12, 262489, None], [2026-08-13, 465343, None], [2026-08-14, 415600, None], [2026-08-15, 274439, None], [2026-08-16, 214010, None]]两条查询都在服务端解析并执行分别返回 1 行和 6 行。parallelism_daily中的 KST 日期分桶演示了真实的 6 个 KST 日桶。这里 sum 仍为 null 只是因为turn_completed不携带parallelism_*属性。5.2 聚合证明结构不变属性映射到turn_completed真实存在的属性$ cd /Users/yeongyu/.agents/skills/omo-native-telemetry/scripts python3 - PY import fetch_data as fd tok fd.read_token() q (fd.QUERIES[parallelism] .replace(eventparallelism_summary, eventturn_completed) .replace(properties.non_eval_saved_round_trips, properties.turn_index) .replace(properties.modeled_wallclock_saved_ms, properties.total_tokens) .replace(properties.non_eval_waves_total, properties.input_tokens) .replace(properties.upper_bound_saved_ms, properties.output_tokens)) print(QUERY:, q) print(ROWS:, fd.hogql(tok, q)) q2 (fd.QUERIES[parallelism_daily] .replace(eventparallelism_summary, eventturn_completed) .replace(properties.modeled_wallclock_saved_ms, properties.total_tokens)) print(QUERY:, q2) for r in fd.hogql(tok, q2): print(r) PY QUERY: select count() sessions, sum(toFloat(properties.turn_index)) saved_round_trips, sum(toFloat(properties.total_tokens)) modeled_saved_ms, sum(toFloat(properties.input_tokens)) waves_total, sum(toFloat(properties.non_eval_waves_multi)) waves_multi, sum(toFloat(properties.non_eval_joined_calls)) joined_calls, sum(toFloat(properties.eval_only_waves)) eval_only_waves, sum(toFloat(properties.mixed_waves)) mixed_waves, sum(toFloat(properties.incomplete_calls)) incomplete_calls, sum(toFloat(properties.clock_anomalies)) clock_anomalies, sum(toFloat(properties.dropped_calls)) dropped_calls, sum(toFloat(properties.measured_turn_duration_ms_total)) measured_turn_ms, sum(toFloat(properties.output_tokens)) upper_bound_saved_ms_ref from events where eventturn_completed ROWS: [[1633594, 80327925.0, 253928243169.0, 32854976150.0, None, None, None, None, None, None, None, None, 747578824.0]] QUERY: select toDate(toTimeZone(timestamp,Asia/Seoul)) d, count() sessions, sum(toFloat(properties.total_tokens)) modeled_saved_ms from events where eventturn_completed group by d order by d [2026-08-11, 1666, 202514905.0] [2026-08-12, 262489, 38811264641.0] [2026-08-13, 465343, 70675922083.0] [2026-08-14, 415600, 67649333555.0] [2026-08-15, 274439, 42810892867.0] [2026-08-16, 214062, 33779245301.0]这是决定性证明查询结构保持不变、只把属性名映射到turn_completed真实携带的字段后每个sum(toFloat(properties.X))槽位都返回实数每日序列也返回了按 KST 日填充的趋势。因此toFloat(...) → sum(...) → 列别名链路与 KST 分组都在正常工作生产运行中的 null/空完全由事件尚未存在导致。这些 relaxed 查询只在子进程中 ad hoc 运行磁盘上的QUERIES字典仍过滤eventparallelism_summary见 §2.1 的 diff没有任何松弛化内容被持久化。6. 对抗性类目清单Adversarial Classes证据文档用一张表格系统性地排查了十类对抗场景这里完整收录类目结论证据畸形/缺失数据零匹配事件已处理§3.1——parallelism.json是sessions0 null 求和的一行parallelism_daily.json是[]。无崩溃、无异常、退出码 0。关键点空聚合不会伪造一个可归因于真实会话的0节省数字——开头的sessions0让下游卡片对空性一目了然误导性成功输出空结果被误认为成功直接证明排除§5.1/§5.2——相同查询结构对turn_completed返回 1 行与 6 行且 §5.2 证明 sum 列在属性存在时产出实数。语法或别名错误会触发 HogQL 错误并被fetch_one以 ERROR 行暴露不稳定行为 / PostHog 限流未绕过fetch_one、hogql、3 次重试、4 * (attempt 1)退避、max_workers4池全部未改动。新查询与所有兄弟查询一样走fetch_one41 条查询一次跑完且 0 ERROR 行新增负载未触发限流兄弟查询回归重排/删除/修改既有条目排除§1 diff 在字典内纯增量§2.3 显示全部 39 条既有查询仍返回行包括parallel_savings30与parallel_savings_basis1eval 桶污染构造上排除eval_only_waves/mixed_waves独占独立列两条查询字符串中任何位置都不存在将 eval 列与non_eval_*列结合的算术运算符平均值之比 / 按会话中位数泄漏构造上排除两条新查询都不含avg、quantile、median或除法运算符只有count()与sum(toFloat(...))因此所有舰队比率必然在下游以 sum/sum 计算upper_bound_saved_ms被提升为头条排除它只出现一次、作为最后一列、别名为upper_bound_saved_ms_ref且完全不出现在parallelism_daily趋势序列只携带modeled_wallclock_saved_ms错误的时区分桶排除parallelism_daily插值模块级KST常量toTimeZone(timestamp,Asia/Seoul)而非重新声明§5.1/§5.2 显示toDate(toTimeZone(timestamp,Asia/Seoul))正确产出 KST 日桶非标准库依赖蔓延排除import 行零改动diff 只触及QUERIES字典字面量并发 worker 工作树污染仓库排除工作树下唯一写入的路径就是这份证据文件技能编辑位于~/.agents/skills/...不在任何 git 树内7. 清理凭据与最终结论$ rm -rf /tmp/t7-qa /tmp/t7-full /tmp/t7-full.log /tmp/t7.diff /tmp/fetch_data.py.bak $ ls -d /tmp/t7-qa /tmp/t7-full 21 ls: /tmp/t7-full: No such file or directory ls: /tmp/t7-qa: No such file or directory两个临时输出目录、全量库日志、diff 临时文件与用于生成 diff 的编辑前备份副本全部删除无任何临时产物残留。最终结果PASS。python3 fetch_data.py /tmp/t7-qa --only parallelism,parallelism_daily→ 退出码 0/tmp/t7-qa/parallelism.json与/tmp/t7-qa/parallelism_daily.json均存在且可被解析为合法 JSON空/零会话结果符合当前预期python3 fetch_data.py /tmp/t7-full→ 退出码 041 条查询中0条 ERRORrelaxed-filter 运行证明 SQL 真实有效、聚合链路真实求和。8. 可复用的方法论要点从 Task 7 中可以提炼出一套可迁移到其他查询库扩展的工程纪律纯增量 diff新条目追加在QUERIES字典末尾先做 AST 语法检查再--only定向运行最后全量回归确认 0 ERROR 与兄弟查询形态不变。空结果即诚实结果让count()与 sum 列同时出现在 payload 中空聚合的sessions0 null 求和天然自证其空避免伪造0数值误导下游可视化。relaxed-filter 双证明先只换事件过滤器证明结构可解析、分组正确再把属性映射到真实存在字段证明聚合链路真实求和两步即可把无数据与查询错误彻底区分。比率语义防线前置在 SQL 中只发原始求和把 sum/sum 比率推导留给下游并让 eval 桶独立成列、上界指标显式标注_ref后缀从查询文本层面杜绝指标误用。对照发送端 schema 逐字节核对属性名查询中的properties.X必须与事件 schema如 parallelism-schema.ts和发送端实现如 omo-native-parallel-summary.ts保持一致属性白名单与事件覆盖测试omo-native-telemetry-assertions.mjs会在部署侧兜底。这套流程的边界条件同样值得注意目标文件位于仓库之外的技能目录因此 diff 与证据记录在.omo/evidence/telemetry-parallel-latency-v2/task-7.md中一旦parallelism_summary事件随parallelism_v2schema见 parallelism-schema.ts 的schema_kind枚举部署上线这两条查询将直接从空结果切换为真实的会话级与按日并行度节省数据无需再改查询文本。【免费下载链接】oh-my-openagentOmO: Just type mass ulw keyword with your prompt. Now you are the master of graph engineering.项目地址: https://gitcode.com/gh_mirrors/oh/oh-my-openagent创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考