Opik Backend 指标插桩规范:用 OpenTelemetry 为 LLM 工作流构建按阶段、按工作区的可观测性
2026/9/13 14:41:03 网站建设 项目流程

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_idworkspace_name—— 客户下钻(customer drill)维度,必须成对出现;name 缺失时回退为 id(见 §3.3)。
  • 阶段/类型标签—— 这是"哪种工作":evaluator_typedecision、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_idworkspace_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"); ... });

对应的三条规范性要求:

  1. workspace_name缺失时必须回退为workspace_id——注意实现中回退到 id 而不是回退到字面量 "unknown",保证 name 永远不会变成无意义的unknown
  2. 在响应式上下文已经携带 name 的路径上,禁止再做 name-service 查找(避免无谓的 RPC/数据库调用)。
  3. 如果某个事件尚未携带 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)是消费者侧指标族的集中实现。它通过构造函数接收metricNamespacemetricsBaseName,从而以${namespace}_${baseName}_...模式为每个具体订阅者(subscriber)生成一套指标。完整族谱:

指标类型说明
${ns}_${base}_processing_timenative histogram(ms)处理一条消息的耗时(scorer/LLM 工作),按 workspace 归因
${ns}_${base}_queue_delaynative histogram(ms)消息入队到处理结束之间的延迟(从消息 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记录每次调用的批大小(水平量),空结果时置 0(Objects.requireNonNullElse(messages, Map.of())),避免空MonodoOnSuccess空指针(BaseRedisSubscriber.java)。
  • 时间度量用doFinally收尾claimTime/readTime/messageProcessingTime都记录System.currentTimeMillis() - startMillisdoFinally保证无论成功失败都记账。
  • 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>(惰性)会破坏那些"为了副作用而调用"的既有测试。恢复绿灯的三板斧:

  1. 生产代码组合它的地方,用宽松桩替换:lenient().when(pub.enqueue(any(),any())).thenReturn(Mono.empty())
  2. 单元测试中需要断言下游效果的,对返回的Mono调用.block()
  3. 原先用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发送代码评审通知。
  • 配套看板 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_totalMono<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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询