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_id+workspace_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 采样水平 |
| UpDownCounter(OTel) | 有符号 +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.java):
private static final String METRIC_NAMESPACE = "online_scoring"; private static final AttributeKey<String> EVALUATOR_TYPE_KEY = AttributeKey.stringKey("evaluator_type"); private static final AttributeKey<String> 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). result=error counts publish failures that were previously only logged.") .build();同样地,在线评分采样器(OnlineScoringSampler.java)注册了online_scoring_sampler_decisions_total:
Meter 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 回退到 id
workspace_id与workspace_name必须从响应式请求上下文读取(workspace-id / workspace-name 上下文键),并复用共享的 workspace 属性键常量,禁止在每个调用点重新声明。这些常量集中在 ErrorMetricsResolver.java:
public static final AttributeKey<String> ERROR_TYPE_KEY = AttributeKey.stringKey("error_type"); public static final AttributeKey<String> WORKSPACE_ID_KEY = AttributeKey.stringKey("workspace_id"); public static final AttributeKey<String> WORKSPACE_NAME_KEY = AttributeKey.stringKey("workspace_name"); public static final AttributeKey<String> USER_NAME_KEY = AttributeKey.stringKey("user_name"); public static final AttributeKey<String> 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 Mono<Void> 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()); }配套的硬性纪律:
- 方法必须返回
Mono<Void>,且不得自行订阅;调用方必须组合或订阅它(否则就是 §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返回Mono<Void>,Redis stream 写入通过subscribeOn(Schedulers.boundedElastic())挪出调用方线程(OnlineScorePublisher.java);enqueueThreadMessage里对规则的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 与结果下钻)。 - 入队到 Redis:
online_scoring_enqueue_total{evaluator_type, workspace_id, workspace_name, result},其中result="error"计数的是一次真实的发布失败(push 失败 = 真实丢失)——源码注释明确指出这补上了"以前只打日志不计数"的盲区(OnlineScorePublisher.java)。
4.2 消费者侧(Consumer metrics):吞吐与分阶段计时
BaseRedisSubscriber(BaseRedisSubscriber.java)是消费者侧指标族的集中实现。它通过构造函数接收metricNamespace与metricsBaseName,从而以${namespace}_${baseName}_...模式为每个具体订阅者(subscriber)生成一套指标。完整族谱:
| 指标 | 类型 | 说明 |
|---|---|---|
${ns}_${base}_processing_time | native histogram(ms) | 处理一条消息的耗时(scorer/LLM 工作),按 workspace 归因 |
${ns}_${base}_queue_delay | native histogram(ms) | 消息入队到处理结束之间的延迟(从消息 ID 中的时间戳推算) |
${ns}_${base}_processing_errors | counter | 处理消息出错次数,按error_type+ workspace + user 下钻 |
${ns}_${base}_undecodable_messages_total | counter | 无法解码的流条目(可重试,不视为 drop) |
${ns}_${base}_backpressure_drops_total | counter | 轮询 tick 因消费者忙而被丢弃的次数(良性,非丢工作) |
${ns}_${base}_claim_errors/${ns}_${base}_claim_time/${ns}_${base}_claim_size | counter / histogram / gauge | 每次XAUTOCLAIM认领操作的错误、耗时、认领条数 |
${ns}_${base}_read_errors/${ns}_${base}_read_time/${ns}_${base}_read_size | counter / histogram / gauge | 每次readGroup读取操作的错误、耗时、返回条数 |
${ns}_${base}_ack_and_remove_errors/${ns}_${base}_ack_and_remove_time | counter / histogram | ack 并移除操作的错误与耗时 |
${ns}_${base}_list_pending_errors/${ns}_${base}_list_pending_time | counter / histogram | 列出 pending 消息操作的错误与耗时 |
${ns}_${base}_unexpected_errors | counter | 主循环捕获的意外异常,按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记录每次调用的批大小(水平量),空结果时置 0(Objects.requireNonNullElse(messages, Map.of())),避免空Mono的doOnSuccess空指针(BaseRedisSubscriber.java)。- 时间度量用
doFinally收尾:claimTime/readTime/messageProcessingTime都记录System.currentTimeMillis() - startMillis,doFinally保证无论成功失败都记账。 - backpressure 计数:
Flux.interval上的onBackpressureDrop回调累加backpressure_drops_total(BaseRedisSubscriber.java)——这是 §4.4 约束的落地。 processing_time直方图按 workspace 归因:messageContext(message)一次解析 workspace/user,成功与失败路径复用同一组属性(BaseRedisSubscriber.java)。processing_errors_total的完整维度:error_type(异常类/类别)+workspace_id+workspace_name+user_name,缺省时一律回退到unknown(BaseRedisSubscriber.java)。
4.3 入口 RED 与饱和度
- 入口 RED:工作流的前门(HTTP 摄入路由)从
http_server_request_duration_seconds读取 Rate/Errors/Duration;5xx 按 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 指出一个实际痛点:把方法改造成返回Mono<Void>(惰性)会破坏那些"为了副作用而调用"的既有测试。恢复绿灯的三板斧:
- 生产代码组合它的地方,用宽松桩替换:
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 没有它会失败)、## Issues(Resolves OPIK-XXXX)与## AI-WATERMARK: yes(含 tools/model/scope/human-verification); - 标题格式
[OPIK-XXXX] [BE] …; - 合入前按
.agents/commands/comet/send-code-review-slack.md发送代码评审通知。
- 配套看板 PR(comet 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、Mono<Void>enqueue 契约、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),仅供参考