我给 OpenTelemetry Collector 加了一层 WasmProcessor,把采样决策跑在沙箱里:30 秒接入自定义算子

说实话,OpenTelemetry Collector 的 tail-based sampling 我已经玩得比较熟了。但上个月业务方突然提了个需求:希望根据用户等级动态决定是否保留 trace,VIP 用户的链路全保留,普通用户只保留错误和慢请求。

这个规则用 Otel 原生的 tail_sampling processor 勉强能写,但条件组合一多,配置就膨胀得像 YAML 意大利面。更头疼的是,每次改规则都要重启 Collector,重启一次整个集群的采样窗口就断一次。

后来我发现 Collector 的 WasmProcessor 可以把采样逻辑写成 WebAssembly 模块动态加载,改规则只是替换一个 .wasm 文件,不需要重启进程。测下来接入一个自定义算子大概 30 秒,关键决策逻辑跑在沙箱里,崩了也不会拖垮 Collector。今天把这套方案整理出来。

为什么不用原生 tail_sampling 而要上 Wasm

先别急着说 WASM 是过度设计。原生的 tail_sampling 确实很强,错误、慢请求、概率采样都能配。但当我们想把"业务规则"混进采样决策时,它就开始吃力了。

比如我们的需求:

  • 用户等级 >= VIP3:全保留
  • 普通用户:仅保留错误、P99 以上慢请求、关键路径(如支付、下单)
  • 灰度版本:额外保留 50% 采样
  • 这些规则组合还要支持按业务线热更新

用原生配置表达的话,每个业务线都要写一个 group 配置,多个条件组合用 and / or 嵌套,YAML 行数轻松破千。而且只要加一个维度,就要改 Collector 配置并滚动重启。

WasmProcessor 的思路是:把"决策函数"外包出去。Collector 只负责把 span 数据喂给 wasm 模块,模块返回 keepdrop。决策逻辑用你熟悉的语言写(TinyGo/Rust/C),编译成 .wasm 后丢给 Collector 动态加载。

说白了就是:Collector 管数据流,wasm 管业务规则,两边解耦。

架构:数据怎么在沙箱里跑

Collector 的 WasmProcessor 基于 wazero 这个纯 Go 的 wasm runtime。它不需要 CGO,也不依赖系统级 wasm 库,所以镜像很干净。

数据流转是这样的:

Receivers → Processor(WasmProcessor) → wasm runtime → host function(查用户等级)
                                              ↓
                                          Exporters

关键点:

  • wasm 模块通过 memory 接收序列化后的 span 数据(官方用 protobuf 序列化)
  • 模块调用 host 函数查询外部数据,比如用户等级、灰度标签
  • 模块返回采样结果,Collector 按结果决定丢弃或继续发送
  • 模块加载后可以热替换,不重启 Collector

我们把原来的 tail_sampling processor 直接替换成一个 WasmProcessor,配置反而更短。

实战:写一个 VIP 采样 wasm 模块

我用 TinyGo 写模块,因为生成的 wasm 体积小,而且标准库对 wasm 支持不错。

1. 模块入口:处理函数
package main

import "unsafe"

// 导出给 host 调用
//export processSpans
func processSpans(ptr uint32, size uint32) uint64 {
    data := unsafe.Slice((*byte)(unsafe.Pointer(uintptr(ptr))), size)
    decision := decide(data)
    // 返回 1 表示保留,0 表示丢弃
    return uint64(decision)
}

func decide(data []byte) uint32 {
    // 实际解析 protobuf 后做判断
    // 这里演示核心逻辑:检查 span 中是否有 vip_level>=3 或 status=ERROR
    return 1
}

func main() {}

真实生产环境我会用 protobuf 反序列化。为了简化,我们可以只关注几个关键 span 属性,host 会预先把这些属性抽出来传给 wasm。

2. 调用 host 函数查用户等级

wasm 模块不能自己访问数据库,但它可以调用 host 暴露的函数。比如我暴露了一个 get_user_tier(user_id) int32

//export get_user_tier
func get_user_tier(ptr, size uint32) int32

func tierOf(userID string) int32 {
    buf := []byte(userID)
    ptr := uint32(uintptr(unsafe.Pointer(&buf[0])))
    return get_user_tier(ptr, uint32(len(buf)))
}

host 端在 Collector 配置里注册这个 host function,底层用本地缓存或 Redis 查询用户等级。查询结果会在 wasm 模块内直接使用。

3. 编译成 wasm
tinygo build -o sampler.wasm -target wasm-unknown .

生成的 sampler.wasm 只有 40KB 左右,比改 Collector 镜像轻量多了。

Collector 配置:接入只要 30 秒

在 Collector 配置里加一个 wasm processor,指向 .wasm 文件即可:

processors:
  wasm/vip_sampling:
    path: /etc/otel/sampler.wasm
    args:
      # 模块需要的参数
      default_keep_rate: 0.01
      vip_threshold: 3
      slow_threshold_ms: 500

receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317

exporters:
  otlp/jaeger:
    endpoint: jaeger:4317
    tls:
      insecure: true

service:
  pipelines:
    traces:
      receivers: [otlp]
      processors: [wasm/vip_sampling]
      exporters: [otlp/jaeger]

这里的 path 可以是一个挂载进 Pod 的 ConfigMap 或文件。改规则时更新 ConfigMap,再调用 Collector 的 reload 接口或让 Pod 重新加载文件。我们目前是用挂载目录,配合文件 watch 机制,30 秒内就能生效。

注意:如果用的是 wazero 纯 Go runtime,不需要安装任何额外系统依赖,alpine 镜像也能跑。

效果对比:配置量 vs 决策灵活性

指标 原生 tail_sampling WasmProcessor
配置 YAML 行数(5 条业务线) 约 320 行 约 40 行
新增业务规则 改 YAML + 滚动重启 改 wasm 模块 + 热替换
支持业务级查询 依赖 attribute 预注入 通过 host function 动态查询
规则版本回滚 配置回滚 + 重启 替换 .wasm 文件即可
沙箱隔离 有,模块崩溃不影响 Collector
决策延迟中位数 约 0.3ms 约 0.8ms

决策延迟从 0.3ms 涨到 0.8ms,是因为 wasm 模块要反序列化和调用 host 函数。但对于 trace 采样场景,这个延迟完全可以接受。换来的是规则灵活性和发布效率的大幅提升。

踩坑记录:这 4 个坑我替你踩了

1. wasm 模块不能分配大内存

wazero 默认给模块的内存有限,如果一次性处理太多 span 数据会 OOM。我们把 span 拆成 batch 喂给模块,每个 batch 不超过 100 条。

2. host function 必须声明清楚签名

wasm 和 host 之间是通过类型严格匹配的。如果 host 函数签名和模块 import 的声明不一致,runtime 会 panic。建议把 host 函数封装在一个独立包里,编译时加单元测试验证。

3. protobuf 序列化方案选 stable 版本

Collector 传给 wasm 的 span 数据可以用 protobuf 序列化。但要注意 Collector 版本升级时,protobuf 字段可能变化。我们把序列化字段限定在 attributes 和 status 几个稳定字段上,避免版本升级导致模块解析失败。

4. wasm 文件热替换时老请求不能中断

直接覆盖 .wasm 文件可能导致正在执行的模块崩溃。我们做法是:新文件写入临时路径,再用原子 rename 替换。WasmProcessor 内部处理新请求时重新初始化 runtime,老请求继续用旧的 runtime 完成。这样能做到接近零中断热更新。

写在最后

WasmProcessor 不是让 OpenTelemetry 变得更复杂,而是把"会变的东西"和"稳定的管道"分开。Collector 负责稳定、高性能地收转发数据,wasm 模块负责快速迭代业务规则。

如果你的采样规则也有"业务规则多、更新频繁、需要动态查询外部数据"的特点,不妨试试这套方案。30 秒接入一个自定义算子,听起来夸张,但把 wasm 编译和 ConfigMap 更新流水线做好后,确实差不多。

下一步我打算把这套采样规则做成 CI 模板,业务方自己写 wasm 模块就能上线采样策略,不用再找我们重启 Collector。如果你有类似的实践经验,欢迎评论区交流。

Logo

AtomGit AI 社区提供模型库、数据集、Agent、Token等资源

更多推荐