我给 OpenTelemetry Collector 加了一层 WasmProcessor,把采样决策跑在沙箱里:30 秒接入自定义算子
我给 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 模块,模块返回 keep 或 drop。决策逻辑用你熟悉的语言写(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。如果你有类似的实践经验,欢迎评论区交流。
更多推荐



所有评论(0)