
Opik Backend 指标插桩规范用 OpenTelemetry 为 LLM 工作流构建按阶段、按工作区的可观测性【免费下载链接】comet-llmDebug, evaluate, and monitor your LLM applications, RAG systems, and agentic workflows with comprehensive tracing, automated evaluations, and production-ready dashboards.项目地址: https://gitcode.com/GitHub_Trending/co/comet-llm导读本文讲解 comet-llm 仓库Opik后端apps/opik-backend中操作指标operational metrics插桩的规范性实践如何用 OpenTelemetry 将一条 LLM 工作流以在线评分/在线评测 online scoring 为典型示例分解为有序阶段为每个阶段埋设吞吐、延迟、错误三类指标并以workspace客户/工作区作为一等公民维度。读完本文你将掌握 counter / native histogram / gauge / upDownCounter 的选型规则、workspace_idworkspace_name双标签的解析与回退策略、响应式Reactive上下文中惰性Mono的插桩纪律以及如何在仓库源码中找到这些约定的真实落地实现。本文是.agents/skills/metrics-instrumentation/SKILL.md的完整展开所有代码与路径均来自当前仓库可与 opik-backend 技能通用约定与日志规则配套使用。1. 背景指标插桩与其它可观测性的边界在 Opik 后端可观测性分为两半操作指标本技能主题面向 SRE 与后端工程师回答管线现在快不快、稳不稳、堵在哪使用 OpenTelemetry 指标metrics最终落到按流程顺序编排的 Grafana 看板。分析指标analytics-instrumentation面向产品分析使用 PostHog 产品事件回答用户在干什么。参见 analytics-instrumentation 技能。两者互不替代。此外构建 Grafana 看板布局、查询契约、按客户名称解析、依赖面板、校验是另一个独立技能comet 监控工具链中的 dashboard-authoring 技能负责的本技能只规定要发射哪些指标不规定看板怎么画。这种分工保证了埋点契约与展示消费解耦——看板面板在新指标的后端 PR 部署前保持为空两个 PR 可以独立、以任意顺序合入见 §5.3。指标名是稳定契约类的位置会移动。因此在仓库中检索约定时请按**指标族metric families**而非具体类名搜索。2. 设计模型把工作流拆成阶段2.1 按阶段分解工作流任何需要可观测性的工作流评分 scoring、摄入 ingestion、实验 experiments、后台任务 jobs第一步都是把它分解为有序阶段ordered stages每个阶段明确四件事要素说明输入input本阶段消费什么消息、事件、请求输出output本阶段产出什么决策、入队、处理结果失败模式failure mode会以何种方式失败解码失败、Redis 不可用、重试耗尽既有插桩existing instrumentation该阶段是否已有指标/日志避免重复埋点分解完成后每一阶段都用RED覆盖流经它的工作Rate 速率 / Errors 错误 / Duration 耗时用USE覆盖它运行所在的资源Utilization 利用率 / Saturation 饱和度 / Errors 错误——两者互补RED 讲工作、USE 讲资源。2.2 按要回答的问题选仪表类型绝大多数阶段只需要counter发生了没有多久一次和histogram花了多久gauge只用于无法从前面两者推导出来的水平量。仪表类型语义何时使用何时禁用Counter单调递增的总量按 rate 读取事件吞吐、决策、结果、错误、累计体量消息数/字节数每秒多少启动以来多少个绝不可用于会下降的值Native histogram延迟/大小的分布用histogram_quantile读取处理时间、队列延迟、单次依赖操作耗时、端到端延迟、载荷大小字节/字符——凡 p95/p99 比均值更有意义之处计数就能说清的地方如纯成功/失败计数不要用直方图Gauge瞬时水平原样读取会升会降且无法由 counter 重建的量当前队列深度/积压、在飞in-flight数量饱和度、批/读/认领claim大小、锁等待者、堆用量想要趋势就用 counter 数事件而不是用 gauge 采样水平UpDownCounterOTel有符号 1/−1 增量维护的水平天然以开始 1、结束 −1维护的量in-flight 计数比每次观测读一次 size 更廉价、更少竞态—经验法则发生了多少次→ counter多久 / 多大→ histogram此刻有多少→ gauge或 UpDownCounter。关于 histogram 的选型还有一条强制性偏好优先使用 native histogram无显式 bucket§2.1只有当下游消费者例如既有 exporter强制要求时才退回经典lebucket 直方图。2.3 定义身份维度基数必须受限维度label是让指标可聚合、可下钻的关键但基数必须保持有界workspace_id与workspace_name—— 客户下钻customer drill维度必须成对出现name 缺失时回退为 id见 §3.3。阶段/类型标签—— 这是哪种工作evaluator_type、decision、Redis 操作、内容/MIME 类型。它让看板可以按阶段sum by(...)。result∈ {success,error} —— 挂在结果计数器上。error_type—— 挂在每一个错误计数器上取异常类名/失败类别当一个计数器服务于多个调用点listener/subscriber/endpoint时再补一个 component 标签使错误面板能同时按原因和来源下钻。基数红线基数必须被#workspaces × #types × #error_types约束。trace id、用户输入、原始消息、完整 URL 这类无界值严禁上标签。3. 后端指标OpenTelemetry 落地约定3.1 通过 OTel API 创建 Meter一个工作流一个命名空间所有仪表必须通过 OTel API 的GlobalOpenTelemetry.getMeter(namespace)创建每个工作流使用独立的命名空间。规范给出的模板private static final String METRIC_NAMESPACE workflow; var meter GlobalOpenTelemetry.getMeter(METRIC_NAMESPACE); meter.counterBuilder(%s_stage_total.formatted(METRIC_NAMESPACE)).setDescription(…).build(); meter.histogramBuilder(%s_stage_time.formatted(METRIC_NAMESPACE)).setUnit(ms).ofLongs().build(); // native histogram meter.gaugeBuilder(%s_stage_size.formatted(METRIC_NAMESPACE)).build(); meter.upDownCounterBuilder(%s_stage_in_flight.formatted(METRIC_NAMESPACE)).build(); // signed level (inc/dec)要点Counter 在 Prometheus 侧以_total后缀呈现。Native histogram 严禁定义显式 bucket——不得出现_bucket/_sum/_count/le相关显式构造。仓库中的真实落地示例——在线评分发布器OnlineScorePublisher.javaprivate static final String METRIC_NAMESPACE online_scoring; private static final AttributeKeyString EVALUATOR_TYPE_KEY AttributeKey.stringKey(evaluator_type); private static final AttributeKeyString RESULT_KEY AttributeKey.stringKey(result); ... this.enqueueCounter GlobalOpenTelemetry.getMeter(METRIC_NAMESPACE) .counterBuilder(%s_enqueue_total.formatted(METRIC_NAMESPACE)) .setDescription(Messages pushed to the online-scoring Redis stream, by evaluator type, workspace and result (success|error). resulterror counts publish failures that were previously only logged.) .build();同样地在线评分采样器OnlineScoringSampler.java注册了online_scoring_sampler_decisions_totalMeter meter GlobalOpenTelemetry.getMeter(ONLINE_SCORING_NAMESPACE); this.samplingDecisions meter.counterBuilder(online_scoring_sampler_decisions_total) .setDescription(Online-scoring sampling decisions, by workspace, evaluator type and outcome ...) .build();3.2 阶段/类型维度必须用 label不要编码进指标名online_scoring_enqueue_total{evaluator_type, result}这样的形态是规范要求stage/type 维度是 label而不是指标名的一部分这样看板才能写sum by(evaluator_type)(rate(...))。仅在扩展现有指标族这一种情况下允许把维度编码进名字但规范明确警示名字编码的维度无法在单个 selector 里跨名字做rate()会让看板聚合复杂化因此默认总是优先 label。3.3 workspace 双标签从响应式上下文读取name 回退到 idworkspace_id与workspace_name必须从响应式请求上下文读取workspace-id / workspace-name 上下文键并复用共享的 workspace 属性键常量禁止在每个调用点重新声明。这些常量集中在 ErrorMetricsResolver.javapublic static final AttributeKeyString ERROR_TYPE_KEY AttributeKey.stringKey(error_type); public static final AttributeKeyString WORKSPACE_ID_KEY AttributeKey.stringKey(workspace_id); public static final AttributeKeyString WORKSPACE_NAME_KEY AttributeKey.stringKey(workspace_name); public static final AttributeKeyString USER_NAME_KEY AttributeKey.stringKey(user_name); public static final AttributeKeyString STREAM_KEY AttributeKey.stringKey(stream); public static final String UNKNOWN unknown;规范的读取与回退逻辑在线评分发布器 OnlineScorePublisher.java 中的真实实现return Flux.deferContextual(ctx - { var workspaceId StringUtils.defaultIfBlank(ctx.getOrDefault(RequestContext.WORKSPACE_ID, UNKNOWN), UNKNOWN); String ctxWorkspaceName ctx.getOrDefault(RequestContext.WORKSPACE_NAME, null); var workspaceName StringUtils.defaultIfBlank(ctxWorkspaceName, workspaceId); var successAttrs Attributes.of(EVALUATOR_TYPE_KEY, type.getType(), WORKSPACE_ID_KEY, workspaceId, WORKSPACE_NAME_KEY, workspaceName, RESULT_KEY, success); var errorAttrs Attributes.of(EVALUATOR_TYPE_KEY, type.getType(), WORKSPACE_ID_KEY, workspaceId, WORKSPACE_NAME_KEY, workspaceName, RESULT_KEY, error); ... });对应的三条规范性要求workspace_name缺失时必须回退为workspace_id——注意实现中回退到 id 而不是回退到字面量 unknown保证 name 永远不会变成无意义的unknown。在响应式上下文已经携带 name 的路径上禁止再做 name-service 查找避免无谓的 RPC/数据库调用。如果某个事件尚未携带 name为它增加可空的workspaceName字段并在发布点从 workspace-name 上下文键填充它——正如本工作流所消费的 entity-created 事件已经做的那样。在线评分发布器还演示了从上下文把 name 戳进消息的做法对实现了RedisSubscriberMessage接口、但workspaceName为空的消息用上下文中的 name 填充使异步消费者及其按工作区指标能拿到真实名称OnlineScorePublisher.java。从源码看BaseRedisSubscriber通过读取该接口的 workspace 信息将*_processing_errors_total与*_processing_time直方图归因到来源工作区/用户见 RedisSubscriberMessage.java 的注释。3.4 被插桩的操作必须响应式用 deferContextual 读上下文规范给出的方法签名与骨架public MonoVoid enqueue(List? messages, Type type) { return Flux.deferContextual(ctx - { /* resolve workspace per 2.3 */ return Flux.fromIterable(messages).flatMap(m - redisAdd(m) .doOnNext(id - counter.add(1, successAttrs)) .doOnError(e - { counter.add(1, errorAttrs); log.error(Error … id{}, id, e); })); }).then().subscribeOn(Schedulers.boundedElastic()); }配套的硬性纪律方法必须返回MonoVoid且不得自行订阅调用方必须组合或订阅它否则就是 §4.1 的静默 no-op。请求作用域的响应式调用方必须用.then(...)/flatMap组合它使其继承 workspace 上下文。事件驱动 / 同步调用方必须通过单一共享 helper以 fire-and-forget 方式订阅helper 内部显式.contextWrite(ctx - ctx.put(WORKSPACE_ID, id).put(WORKSPACE_NAME, defaultIfBlank(name, id)))并带一个记录错误的 consumer——一个 helper不要在每处调用点复制。阻塞查询JDBC /findById必须通过Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic())执行enqueue/IO 工作必须在有界调度器上运行不能占用调用方线程例如 EventBus 线程。在线评分发布器正是如此enqueueMessage返回MonoVoidRedis stream 写入通过subscribeOn(Schedulers.boundedElastic())挪出调用方线程OnlineScorePublisher.javaenqueueThreadMessage里对规则的findById阻塞查找也用Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic())包裹随后flatMap委托给持有已解析规则的另一个重载OnlineScorePublisher.java。后一个重载还体现了调用方已持有解析结果时提供跳过查找的重载这一约定OnlineScorePublisher.java。3.5 日志格式与构建门禁日志占位符必须单引号包裹如evaluator{} workspaceId{}遵循 opik-backend 技能 的日志规则enqueue 日志还应当包含批量大小。交付前mvn -o compile必须通过spotless clean。4. 指标全景在线评分工作流的完整指标族规范指出要观察这些约定的实际形态直接 grepapps/opik-backend中的既有指标族。以在线评分为例可以从源码中梳理出生产者侧、消费者侧、基础设施依赖三块完整族谱。4.1 生产者侧Producer metrics在源头计数在产出端计数的原则——每个产出工作的阶段都在源头计数使生产者什么都没产出导致的下游饥饿变得可解释而非神秘采样决策online_scoring_sampler_decisions_total{decision}OnlineScoringSampler.java按 workspace 与结果下钻。入队到 Redisonline_scoring_enqueue_total{evaluator_type, workspace_id, workspace_name, result}其中resulterror计数的是一次真实的发布失败push 失败 真实丢失——源码注释明确指出这补上了以前只打日志不计数的盲区OnlineScorePublisher.java。4.2 消费者侧Consumer metrics吞吐与分阶段计时BaseRedisSubscriberBaseRedisSubscriber.java是消费者侧指标族的集中实现。它通过构造函数接收metricNamespace与metricsBaseName从而以${namespace}_${baseName}_...模式为每个具体订阅者subscriber生成一套指标。完整族谱指标类型说明${ns}_${base}_processing_timenative histogramms处理一条消息的耗时scorer/LLM 工作按 workspace 归因${ns}_${base}_queue_delaynative histogramms消息入队到处理结束之间的延迟从消息 ID 中的时间戳推算${ns}_${base}_processing_errorscounter处理消息出错次数按error_type workspace user 下钻${ns}_${base}_undecodable_messages_totalcounter无法解码的流条目可重试不视为 drop${ns}_${base}_backpressure_drops_totalcounter轮询 tick 因消费者忙而被丢弃的次数良性非丢工作${ns}_${base}_claim_errors/${ns}_${base}_claim_time/${ns}_${base}_claim_sizecounter / histogram / gauge每次XAUTOCLAIM认领操作的错误、耗时、认领条数${ns}_${base}_read_errors/${ns}_${base}_read_time/${ns}_${base}_read_sizecounter / histogram / gauge每次readGroup读取操作的错误、耗时、返回条数${ns}_${base}_ack_and_remove_errors/${ns}_${base}_ack_and_remove_timecounter / histogramack 并移除操作的错误与耗时${ns}_${base}_list_pending_errors/${ns}_${base}_list_pending_timecounter / histogram列出 pending 消息操作的错误与耗时${ns}_${base}_unexpected_errorscounter主循环捕获的意外异常按error_type下钻关键设计点均可在源码中印证queue_delay 与 processing_time 分离积压backlog与scorer 慢由此可区分端到端延迟 queue_delay processing_time。recordQueueDelay从 Redis 消息 ID 中提取插入时间戳来推算延迟并且特意让提前返回的路径无法解码、无 payload 字段也记录 queue delay——因为可重试的坏条目会循环投递其不断增长的年龄正是 PEL 循环的早期信号BaseRedisSubscriber.java。claim_size/read_size使用 gauge记录每次调用的批大小水平量空结果时置 0Objects.requireNonNullElse(messages, Map.of())避免空Mono的doOnSuccess空指针BaseRedisSubscriber.java。时间度量用doFinally收尾claimTime/readTime/messageProcessingTime都记录System.currentTimeMillis() - startMillisdoFinally保证无论成功失败都记账。backpressure 计数Flux.interval上的onBackpressureDrop回调累加backpressure_drops_totalBaseRedisSubscriber.java——这是 §4.4 约束的落地。processing_time直方图按 workspace 归因messageContext(message)一次解析 workspace/user成功与失败路径复用同一组属性BaseRedisSubscriber.java。processing_errors_total的完整维度error_type异常类/类别workspace_idworkspace_nameuser_name缺省时一律回退到unknownBaseRedisSubscriber.java。4.3 入口 RED 与饱和度入口 RED工作流的前门HTTP 摄入路由从http_server_request_duration_seconds读取 Rate/Errors/Duration5xx 按 endpoint ×error_type× workspace 下钻从而摄入问题永远不会被误判为评分问题。饱和度与资源水位USE 方法用 gauge 暴露在飞工作量每 pod 上限与 JVM 堆 used-vs-limit每 pod使管线在开始失败之前就暴露逼近天花板的迹象——这是对 RED 的补充。体量与载荷大小字节/字符计数器带宽、总字节数与载荷大小分布按内容类型与 workspace 下钻用于成本与影响归因。基础设施依赖工作流依赖的数据存储Redis streams、ClickHouse、锁、MySQL通过其 exporter 与system.query_log呈现在看板上让管线慢能定位到具体拖后腿的依赖。4.4 成功是推导出来的成功永不重复计数成功 吞吐 − 错误以一个成功率头条磁贴呈现。这避免了成功计数 失败计数两套数字对不上的问题。5. 测试引入惰性 Mono 后的既有测试修复规范 §3.1 指出一个实际痛点把方法改造成返回MonoVoid惰性会破坏那些为了副作用而调用的既有测试。恢复绿灯的三板斧生产代码组合它的地方用宽松桩替换lenient().when(pub.enqueue(any(),any())).thenReturn(Mono.empty())单元测试中需要断言下游效果的对返回的Mono调用.block()原先用verifyNoInteractions(mock)的地方改为verify(mock, never()).method(...)——因为现在存在一个宽松桩verifyNoInteractions会误报。这三条处理方式相互配合宽松桩避免未使用桩异常.block()显式触发惰性链verify(never())保留未调用语义。6. 规范性约束必须遵守以下是本技能列出的硬约束任何指标插桩改动都必须满足4.1返回的Mono在订阅前是惰性的未订阅的 enqueue 是静默 no-op。每个调用点必须组合或订阅它。4.2上下文已携带 workspace 时必须从响应式上下文取不得走 name-service 查找§3.3。4.3IO/enqueue 工作必须在有界调度器上运行离开调用方线程§3.4。4.4背压 / 轮询 tick 计数器如backpressure_drops_total计的是消费者忙时被跳过的调度 tick它不是丢失的工作禁止单独据此告警。要发射它但要在文档中向看板构建者说明其良性语义。4.5工作必须基于origin/main从它创建 worktree不能用可能过期的本地检出。7. 交付流程一次指标插桩改动以两个 PR交付各在自己的分支user/TASK-name上指标 PR本仓库 opik必须不含任何客户、集群或基础设施标识符必须遵循.github/pull_request_template.md包括## Documentation章节PR linter 没有它会失败、## IssuesResolves OPIK-XXXX与## AI-WATERMARK: yes含 tools/model/scope/human-verification标题格式[OPIK-XXXX] [BE] …合入前按.agents/commands/comet/send-code-review-slack.md发送代码评审通知。配套看板 PRcomet monitoring按 dashboard-authoring 技能构建。独立性新指标的看板面板在该后端 PR 部署前保持为空所以两个 PR 相互独立可任意顺序合入。8. 在仓库中快速定位指标实现的清单按指标名稳定契约而非类名检索是定位实现的最可靠方式# 在 apps/opik-backend 中按指标族检索在线评分工作流 grep -rn sampler_decisions_total\|enqueue_total\|processing_time_milliseconds\|queue_delay_milliseconds\|processing_errors_total\|unexpected_errors_total apps/opik-backend/src本文涉及的源码锚点相对仓库根目录SKILL.md 原文 —— 规范正文OnlineScorePublisher.java —— 生产者online_scoring_enqueue_total、MonoVoidenqueue 契约、workspace 上下文解析与消息 name 戳印OnlineScoringSampler.java —— 采样决策计数、workspace 上下文携带BaseRedisSubscriber.java —— 消费者侧指标族processing_time / queue_delay / 各 Redis 操作时间 / 错误计数器 / 背压计数ErrorMetricsResolver.java —— 共享的workspace_id/workspace_name/error_type/user_name/stream属性键常量与error_type(Throwable)解析RedisSubscriberMessage.java —— 消费端指标按 workspace/user 归因所依赖的消息接口opik-backend 技能 —— 配套的通用约定与日志规则9. 小结Opik 后端的指标插桩约定可以浓缩为一句话把工作流拆成有序阶段每个阶段用 counter 数事件、用 native histogram 度量延迟、用 gauge 暴露水平把workspace_id/workspace_name作为一等标签从响应式上下文解析并把成功作为吞吐减错误推导出来。生产者侧计数让饥饿可解释消费者侧将 queue_delay 与 processing_time 分离让积压与慢 scorer 可区分USE 方法在失败之前就暴露饱和度。指标名是稳定契约、类的位置会移动——按指标族检索、按看板契约消费就能让按流程顺序从上读到下、数字一断就知道是哪一阶段挂了的目标在每个工作流上可复制。【免费下载链接】comet-llmDebug, evaluate, and monitor your LLM applications, RAG systems, and agentic workflows with comprehensive tracing, automated evaluations, and production-ready dashboards.项目地址: https://gitcode.com/GitHub_Trending/co/comet-llm创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考