学习目标
学完本章,你应该能够:
- 用自己的话讲清 OpenTelemetry 中 Trace / Metric / Log 三者各自回答什么问题、分别适合什么场景,面试时能讲成一张图。
- 在 Go 服务里完整接入 OTel:初始化
TracerProvider+MeterProvider、打 Span 记录业务、记四种 Metric,并正确处理优雅退出(Shutdown)。 - 读懂并写出生产级 Collector 配置:Receiver / Processor / Exporter 流水线,以及那条绝不能乱的 processor 顺序。
- 实现跨服务、跨协议的 trace 串联(HTTP / gRPC / Kafka),讲清 propagator 为什么是串联成败的关键。
- 基于采集到的可观测数据,写一段多维度互相印证的异常检测脚本,并说清楚为什么单一维度误报率高。
前置知识:
- 能把 Go 服务跑起来(
go run/go build),看得懂包、接口、struct。 - 懂
context.Context是怎么在多个函数之间传递的——本文大量代码都靠它串上下文。 - 知道 Docker Compose 的基本用法(起几个容器、映射端口)。
- 大致知道 Prometheus(指标存储)、Jaeger(链路存储)、Kafka(消息队列)是干嘛的,不需要很熟。
本章你会动手做的事:
- 用 docker-compose 一把拉起 Collector + Prometheus + Jaeger 全家桶。
- 把示例 Go 代码里的 Provider 初始化、Span、Metric 跑通,并在 Jaeger 看到一条真实 trace。
- 写一个 3-sigma 异常检测脚本,对 Prometheus 的查询数据跑一轮检测,观察告警阈值怎么影响误报。
先建立直觉:可观测性三件套像什么
类比:把你的分布式系统想成一家大医院的检验科。
Trace是"一份标本的流转单"——从抽血、化验到出报告,每一站谁经手、花了多久,全程可追溯;Metric是"墙上那些实时滚动的仪表"——每分钟来了多少标本、平均耗时多少、仪器故障率,一眼看全局趋势;Log是"每道工序手写的备注"——某次异常的具体细节和上下文。OpenTelemetry 就是帮你在每家分店统一发"流转单 + 仪表 + 备注"的标准,最后汇总到总院(后端存储)统一分析。
本文要打通的,正是下面这条从应用到后端的完整链路。这张图建议先记住整体骨架,后面每一节都是在填充其中的某一块:
flowchart LR
App[Go 应用
埋点产生 Trace / Metric / Log] -->|OTLP 协议| Collector[OpenTelemetry Collector
接收 处理 转发]
Collector -->|Remote Write| Prom[(Prometheus
存储 Metrics)]
Collector -->|OTLP| Jaeger[(Jaeger
存储 Traces)]
Collector -->|HTTP Push| Loki[(Loki
存储 Logs)]实战目标
在前一篇笔记中,我们了解了 OpenTelemetry 的架构与核心概念。本篇将完成一个完整的开发实战:从 Golang 应用埋点,到 Collector 收集转发,再到 Prometheus 和 Jaeger 后端存储,最后基于采集到的数据进行简单的 AIOps 异常检测。
我写这篇笔记的时候假设你已经能把 Go 服务跑起来、能看懂 context.Context 在函数间怎么传。如果这俩还卡壳,先回去补一下基础再来,否则代码里的 ctx 传开会让你懵。
环境准备
需要部署以下组件:
- OpenTelemetry Collector:负责接收和转发数据。
- Prometheus:存储 Metrics。
- Jaeger:存储 Traces。
- 示例 Golang 应用:产生 Metrics 和 Traces。
可以使用 Docker Compose 或 Helm 在本地快速搭建:
# docker-compose 简化示例
services:
otel-collector:
image: otel/opentelemetry-collector-contrib:latest
volumes:
- ./otel-collector-config.yaml:/etc/otelcol-contrib/config.yaml
ports:
- "4317:4317"
- "4318:4318"
- "8889:8889"
prometheus:
image: prom/prometheus:latest
volumes:
- ./prometheus.yml:/etc/prometheus/prometheus.yml
ports:
- "9090:9090"
jaeger:
image: jaegertracing/all-in-one:latest
ports:
- "16686:16686"
- "14250:14250"
这套 docker-compose 我本地跑过 N 次,几个常见坑先说一下:
- Collector 镜像一定要用
contrib版本,基础版没带 prometheusremotewrite exporter,会报"unknown exporter"。 - Jaeger all-in-one 只适合开发,生产别用,重启数据全没。
- 4317 是 gRPC、4318 是 HTTP,应用和 Collector 在同一台机器用哪个都行,跨网络建议 gRPC(长连接省握手)。
- 如果 Collector 启动失败,九成是配置文件 YAML 缩进错了,看日志先看
service.pipelines那段。
Golang 应用接入 OpenTelemetry
建立直觉:Provider 是"总部",Span 是"工单"
类比:接入 OTel 之前,先分清两个角色。
TracerProvider和MeterProvider像是你们公司的"总部"——它管着跟 Collector 的连接、采样策略、资源属性,全局只有一个;而每次业务调用里你tracer.Start(...)出来的Span,像是总部派发给一线的一张"工单",工单上写着这次操作用了多久、带了哪些业务标签、有没有出错。一线干完活把工单交回总部,总部攒一批批量寄给 Collector。理解了这个"总部—工单"关系,下面这段代码就是在搭总部、发工单、最后关门(Shutdown)别把没寄出的工单弄丢。
初始化这两个 Provider 时最容易乱的就是"连接复用"和"优雅退出",下面用一张流程图把顺序理清楚:
flowchart TD
Start[Setup 被调用] --> Res[创建 Resource
service.name / version / env]
Res --> Dial[拨一次 gRPC 连接
conn 复用给 trace 和 metric]
Dial --> TP[创建 TraceExporter
+ TracerProvider
BatchSpanProcessor]
Dial --> MP[创建 MetricExporter
+ MeterProvider
PeriodicReader]
TP --> Prop[设置全局 TextMapPropagator]
MP --> Prop
Prop --> Ret[返回 Providers
调用方负责 Shutdown]初始化 TracerProvider 与 MeterProvider
原文那段代码是个骨架,能跑但缺一堆东西:导出失败怎么办、资源属性怎么带、应用退出时怎么 flush。我把完整版贴出来,这段代码我基本是直接复制到每个新项目里改改就用:
// pkg/otel/setup.go
package otel
import (
"context"
"errors"
"fmt"
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc"
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
"go.opentelemetry.io/otel/propagation"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/resource"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
semconv "go.opentelemetry.io/otel/semconv/v1.21.0"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
// Config OTel 初始化配置
type Config struct {
ServiceName string // 必填,所有数据都会带这个标签
ServiceVersion string // 建议填,版本信息方便排查
Environment string // production/staging/dev
OTLPEndpoint string // Collector 地址,比如 localhost:4317
SampleRatio float64 // 采样率 0~1,生产建议 0.1~0.3
}
// Providers 把初始化结果打包返回,方便调用方在退出时 shutdown
type Providers struct {
TracerProvider *sdktrace.TracerProvider
MeterProvider *sdkmetric.MeterProvider
}
// Setup 一把初始化 TracerProvider + MeterProvider
// 调用方必须在退出时调用 Shutdown,否则 buffer 里的数据会丢
func Setup(ctx context.Context, cfg Config) (*Providers, error) {
if cfg.ServiceName == "" {
return nil, errors.New("service name is required")
}
if cfg.OTLPEndpoint == "" {
cfg.OTLPEndpoint = "localhost:4317"
}
if cfg.SampleRatio <= 0 {
cfg.SampleRatio = 1.0 // 默认全采样,开发环境用
}
// 资源属性:每个 span/metric/log 都会带上这些
// service.name 必填,其他都是强烈建议
res, err := resource.New(ctx,
resource.WithAttributes(
semconv.ServiceName(cfg.ServiceName),
semconv.ServiceVersion(cfg.ServiceVersion),
semconv.DeploymentEnvironment(cfg.Environment),
),
)
if err != nil {
return nil, fmt.Errorf("create resource failed: %w", err)
}
// ---- Trace 初始化 ----
// 关键:dial 一次 gRPC 连接复用给 trace 和 metric,别 new 两次
// 每次都 dial 会在高 QPS 下耗光连接
conn, err := grpc.DialContext(ctx, cfg.OTLPEndpoint,
grpc.WithTransportCredentials(insecure.NewCredentials()),
// 关键:block=true 等连接建立,否则 exporter 第一次发送会失败
grpc.WithBlock(),
)
if err != nil {
return nil, fmt.Errorf("dial collector failed: %w", err)
}
traceExp, err := otlptracegrpc.New(ctx,
otlptracegrpc.WithGRPCConn(conn),
)
if err != nil {
return nil, fmt.Errorf("create trace exporter failed: %w", err)
}
// 采样器:ParentBased 让子 span 跟随父 span 决策
// TraceIDRatioBased 按比例采样,简单实用
sampler := sdktrace.ParentBased(
sdktrace.TraceIDRatioBased(cfg.SampleRatio),
)
tp := sdktrace.NewTracerProvider(
sdktrace.WithResource(res),
sdktrace.WithSampler(sampler),
// BatchSpanProcessor:批量发送,生产必用
// 关键参数:BatchTimeout 越小延迟越低但吞吐越差
sdktrace.WithBatcher(traceExp,
sdktrace.WithBatchTimeout(2*time.Second),
sdktrace.WithMaxExportBatchSize(512),
),
// 关键:OnExport 失败时记日志,否则 exporter 静默丢数据你都不知道
// sdktrace.WithRawExporter 之类的高级用法这里不展开
)
otel.SetTracerProvider(tp)
// ---- Metric 初始化 ----
metricExp, err := otlpmetricgrpc.New(ctx,
otlpmetricgrpc.WithGRPCConn(conn),
)
if err != nil {
return nil, fmt.Errorf("create metric exporter failed: %w", err)
}
// PeriodicReader:每隔一段时间把累积的 metric 推出去
// 关键:间隔太短会增加 Collector 压力,太长则告警延迟
reader := sdkmetric.NewPeriodicReader(metricExp,
sdkmetric.WithInterval(15*time.Second),
sdkmetric.WithTimeout(5*time.Second),
)
mp := sdkmetric.NewMeterProvider(
sdkmetric.WithResource(res),
sdkmetric.WithReader(reader),
)
otel.SetMeterProvider(mp)
// 全局 propagator:让 trace context 能跨进程传递
// baggage 是用来传业务自定义键值对的,跟 trace 上下文一起走
otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator(
propagation.TraceContext{}, // W3C traceparent
propagation.Baggage{}, // W3C baggage
))
return &Providers{TracerProvider: tp, MeterProvider: mp}, nil
}
// Shutdown 优雅退出,必须调用
// 不调的话 BatchSpanProcessor 里没发出去的 span 会丢
// 我踩过坑:以为进程退出 OS 会自动 flush,实际上 OTel 没这逻辑
func (p *Providers) Shutdown(ctx context.Context) error {
var errs []error
if p.TracerProvider != nil {
if err := p.TracerProvider.Shutdown(ctx); err != nil {
errs = append(errs, fmt.Errorf("tracer shutdown: %w", err))
}
}
if p.MeterProvider != nil {
if err := p.MeterProvider.Shutdown(ctx); err != nil {
errs = append(errs, fmt.Errorf("meter shutdown: %w", err))
}
}
return errors.Join(errs...)
}
应用 main 里这么用:
// cmd/api/main.go
package main
import (
"context"
"log"
"os"
"os/signal"
"syscall"
"time"
"yourapp/pkg/otel"
)
func main() {
ctx, cancel := signal.NotifyContext(context.Background(),
syscall.SIGINT, syscall.SIGTERM)
defer cancel()
providers, err := otel.Setup(ctx, otel.Config{
ServiceName: "order-service",
ServiceVersion: "v1.2.3",
Environment: "production",
OTLPEndpoint: "otel-collector:4317",
SampleRatio: 0.2,
})
if err != nil {
log.Fatalf("otel setup failed: %v", err)
}
// 关键:defer shutdown,但要用独立的 context,因为主 ctx 可能已被 cancel
defer func() {
shutdownCtx, c := context.WithTimeout(context.Background(), 5*time.Second)
defer c()
if err := providers.Shutdown(shutdownCtx); err != nil {
log.Printf("otel shutdown failed: %v", err)
}
}()
// 业务代码...
<-ctx.Done()
log.Println("service shutting down")
}
⚠️ 新手必踩的坑:Shutdown 用了主 ctx。主 ctx 已经被信号 cancel 了,如果直接拿它去调
Shutdown,会立刻返回一个错误,buffer 里没发出去的数据就这么丢了。必须用上面代码里那种"独立的、带 timeout 的 context"。另外defer是 LIFO(后进先出)执行,注意别让别的 defer 把 shutdown 挤到后面、或抢先把连接关了。
踩坑提示:shutdown 用主 ctx 会导致一个微妙问题——主 ctx 已经被信号 cancel 了,shutdown 立刻返回错误,buffer 里数据没发出去。必须用一个独立的、带 timeout 的 context。还有 defer 的执行顺序是 LIFO,注意别让别的 defer 把 shutdown 抢在前面。
创建 Span 并记录业务操作
光初始化不算完,你得知道怎么在实际业务里写 span。下面这段代码覆盖了所有常用 API:baggage、attribute、event、错误记录、status。
类比:一个
Span就是一次"操作档案"。attribute是写死在档案上的固定标签(比如 user.id、订单数),随档案一起走;event是档案里按时间记录的"过程中发生了啥"(比如"缓存查询失败"),带时间戳;baggage比较特殊——它不是写在这张档案上的,而是跟着档案一起被"随件寄往下游",下游开新档案时还能看到。理解这三者的区别,下面代码就不绕了。
一条 trace 里多个 span 是怎么通过 parent_span_id 串成树的,先看这张图(这就是后面 cache.get、db.create 自动挂到 order.create 下的原因):
flowchart TD
Root["order.create (Internal)"] --> C1["cache.get (Internal)"]
Root --> C2["db.create (Internal)"]
C1 -. parent_span_id .-> Root
C2 -. parent_span_id .-> Root// internal/order/service.go
package order
import (
"context"
"errors"
"fmt"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/baggage"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
)
type Service struct {
tracer trace.Tracer
repo Repo
cache Cache
}
func NewService(repo Repo, cache Cache) *Service {
// tracer 名字建议用包名或者服务名,方便筛选
return &Service{
tracer: otel.Tracer("order-service"),
repo: repo,
cache: cache,
}
}
// CreateOrder 创建订单,演示一个完整的 span 用法
func (s *Service) CreateOrder(ctx context.Context, userID string, items []string) (*Order, error) {
// 1. 创建 span。ctx 一定要传进来,否则跟父 span 断了
// span 名建议用动词+名词,比如 "order.create",别用 "CreateOrder" 函数名
ctx, span := s.tracer.Start(ctx, "order.create",
// 关键:业务属性在这里塞,排查问题就靠这些
trace.WithAttributes(
attribute.String("user.id", userID),
attribute.Int("items.count", len(items)),
),
// span kind:server/internal/client/producer/consumer
// 内部业务逻辑用 Internal,接收外部请求用 Server
trace.WithSpanKind(trace.SpanKindInternal),
)
defer span.End()
// 2. baggage:跨 span 跨进程传业务键值对
// 跟 attribute 的区别:baggage 会跟着 trace context 传到下游服务
// 比如把 user.id 一直传到 payment-service
if m, err := baggage.NewMember("user.id", userID); err == nil {
if b, err := baggage.New(m); err == nil {
ctx = baggage.ContextWithBaggage(ctx, b)
}
}
// 3. 缓存查询子 span
cacheHit, err := s.checkCache(ctx, userID)
if err != nil {
// 关键:错误不要直接 span.SetStatus,先记 event 再标记
// event 比 status 信息更丰富,带时间戳和属性
span.AddEvent("cache.check.failed", trace.WithAttributes(
attribute.String("error.type", "redis_error"),
attribute.String("error.message", err.Error()),
))
// 缓存失败不阻断主流程,继续走 DB
}
if cacheHit != nil {
span.SetAttributes(attribute.Bool("cache.hit", true))
return cacheHit, nil
}
span.SetAttributes(attribute.Bool("cache.hit", false))
// 4. 写 DB
order, err := s.repo.Create(ctx, userID, items)
if err != nil {
// 错误处理标准三件套:record error + set status + return
return s.recordError(ctx, span, "db.create.failed", err)
}
// 5. 发事件,记录关键业务节点
span.AddEvent("order.created", trace.WithAttributes(
attribute.String("order.id", order.ID),
attribute.Float64("order.amount", order.Amount),
))
return order, nil
}
// checkCache 演示一个子 span
func (s *Service) checkCache(ctx context.Context, userID string) (*Order, error) {
// 关键:子 span 用 ctx 启动,自动建立 parent 关系
ctx, span := s.tracer.Start(ctx, "cache.get")
defer span.End()
span.SetAttributes(attribute.String("cache.key", "user:"+userID))
order, err := s.cache.Get(ctx, "user:"+userID)
if err != nil {
// 这里别 span.SetStatus(Error),缓存失败不是业务错误
// 加个 attribute 标记一下就行,主 span 自己决定要不要降级
span.SetAttributes(attribute.String("cache.error", err.Error()))
return nil, err
}
return order, nil
}
// recordError 统一错误处理,避免每个地方都写三行
func (s *Service) recordError(ctx context.Context, span trace.Span, event string, err error) (*Order, error) {
// 1. RecordError:标准 OTel API,会自动加 exception.type/message/stacktrace 属性
span.RecordError(err)
// 2. SetStatus:把 span 标记为 ERROR,Jaeger 里会显示红色
span.SetStatus(codes.Error, err.Error())
// 3. AddEvent:再加个自定义事件,方便按事件名筛选
span.AddEvent(event, trace.WithAttributes(
attribute.String("error.class", fmt.Sprintf("%T", err)),
))
// 4. 返回包装后的错误,业务层照常处理
return nil, fmt.Errorf("create order failed: %w", err)
}
// Order 订单结构
type Order struct {
ID string `json:"id"`
UserID string `json:"user_id"`
Amount float64 `json:"amount"`
Items []string `json:"items"`
}
// Repo 接口抽象
type Repo interface {
Create(ctx context.Context, userID string, items []string) (*Order, error)
}
// Cache 接口抽象
type Cache interface {
Get(ctx context.Context, key string) (*Order, error)
}
var _ = errors.New // 防 import 报错
几个我必须强调的点:
- span 名不要用函数名。
CreateOrder在 Jaeger 里看到的就是一坨大写驼峰,根本不知道是干啥的。用order.create这种"领域.动作"格式,一眼就懂。 - 错误处理三件套:
RecordError+SetStatus(Error)+ 业务层返回带 wrap 的 error。少一个都不行,少了 RecordError 就没堆栈,少了 SetStatus trace 不变红,少了 wrap 业务层就丢了上下文。 - baggage 不要塞大字段。baggage 是跟着每个请求头传的,塞个 1KB 的 JSON 进去,每个下游调用都多 1KB 流量。一般只放 user.id、tenant.id 这种短字段。
- attribute 数量控制。一个 span 别加 50 个 attribute, indexing 会拖慢后端。重要的几个就够了,详细信息写 event 或者 log。
⚠️ 新手必踩的坑:把缓存失败当成业务错误标红。
checkCache返回 err 时,千万不要顺手span.SetStatus(codes.Error, ...)——缓存挂了不代表"创建订单"这个业务失败,主流程还会继续走 DB。正确做法是只挂一个attribute标记,让上层order.create自己决定怎么降级。一旦错标成 Error,Jaeger 里这支 trace 整条变红,你排查时会被误导。
踩坑提示:span.RecordError(err) 默认会捕获 stacktrace,但只在 panic 或者显式 RecordError 时才会触发。runtime error 的 stacktrace 可能不全,因为 Go 的栈在 error 创建时已经回溯过了。如果你一定要完整 stack,可以用 pkg/errors 或者 fmt.Errorf 配合 %+v。
记录 Metrics
metric 三种类型 counter / histogram / observable gauge 我都贴一下,每种用法不一样,混着用会出问题。
类比:Metric 类型选错,就像用"温度计"去数"今天来了多少客人"。
Counter是只能往上翻的计数器(客人总数);Histogram是把一堆数值丢进桶里看分布(客人年龄分布、请求耗时分位);UpDownCounter是能加能减的当前值(店里现在还有多少客人);ObservableGauge是去问一个"外部数据源"拿到的当前快照(鱼缸里现在还剩多少水)。选之前先问自己:我要的是"累计" “分布” 还是"当前值"?下面这张决策图帮你选:
flowchart TD
Q[我要统计什么] --> A{是累计次数吗}
A -->|是| C[Counter
请求数 错误数 订单数]
A -->|否| B{是当前值吗}
B -->|可增可减| U[UpDownCounter
队列长度 活跃连接]
B -->|数据源在外部| G[ObservableGauge
连接池大小 goroutine 数]
B -->|是分布| H[Histogram
延迟分位 响应大小]// internal/order/metrics.go
package order
import (
"context"
"time"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/metric"
)
// Metrics 把所有指标集中起来,方便管理
type Metrics struct {
ordersCreated metric.Int64Counter
orderDuration metric.Float64Histogram
ordersInFlight metric.Int64UpDownCounter
cacheSize metric.Int64ObservableGauge
cache Cache
}
func NewMetrics(cache Cache) (*Metrics, error) {
meter := otel.Meter("order-service")
m := &Metrics{cache: cache}
var err error
// 1. Counter:只能加不能减,用于累计计数
// 关键:description 和 unit 别省,Grafana 配置时要用
m.ordersCreated, err = meter.Int64Counter(
"orders_created_total",
metric.WithDescription("Total number of orders created"),
metric.WithUnit("{order}"),
)
if err != nil {
return nil, err
}
// 2. Histogram:分桶统计,用于延迟分布
// 关键:boundaries 要根据业务设,默认的边界偏 Web 通用
// 订单创建场景 5ms~10s 比较合理,别用默认的 0~1000s
m.orderDuration, err = meter.Float64Histogram(
"order_duration_seconds",
metric.WithDescription("Order creation duration"),
metric.WithUnit("s"),
metric.WithExplicitBucketBoundaries(0.005, 0.01, 0.05, 0.1, 0.3, 0.5, 1, 3, 5, 10),
)
if err != nil {
return nil, err
}
// 3. UpDownCounter:可加可减,用于当前值(队列长度、活跃请求数)
m.ordersInFlight, err = meter.Int64UpDownCounter(
"orders_in_flight",
metric.WithDescription("Number of orders being processed"),
metric.WithUnit("{order}"),
)
if err != nil {
return nil, err
}
// 4. ObservableGauge:异步指标,SDK 采集时回调你提供的函数拿值
// 适合那些"已经有现成数据源"的场景,比如 Redis 连接池当前大小
// 关键:回调函数会被 SDK 周期性调用,别在里面做重活
m.cacheSize, err = meter.Int64ObservableGauge(
"cache_size_bytes",
metric.WithDescription("Current cache size in bytes"),
metric.WithUnit("By"),
)
if err != nil {
return nil, err
}
// RegisterCallback 注册采集回调
// 关键:回调返回的 Observation 必须用 observable 个体的 ObserveInt64 方法
_, err = meter.RegisterCallback(func(ctx context.Context, observer metric.Observer) error {
size := m.cache.Size(ctx) // 业务侧自己实现
observer.ObserveInt64(m.cacheSize, size)
return nil
}, m.cacheSize)
if err != nil {
return nil, err
}
return m, nil
}
// RecordOrder 在创建订单的业务函数里调用
func (m *Metrics) RecordOrder(ctx context.Context, duration time.Duration, success bool) {
// Counter:按结果加标签,方便分别看成功率
attrs := metric.WithAttributes()
if !success {
// 关键:标签值不要用高基数字段,比如 order_id
// 高基数会把 Prometheus 内存撑爆
attrs = metric.WithAttributes() // 实际用 attribute.String("result", "error")
}
m.ordersCreated.Add(ctx, 1, attrs)
m.orderDuration.Record(ctx, duration.Seconds())
}
// StartInFlight 标记订单开始处理
func (m *Metrics) StartInFlight(ctx context.Context) {
m.ordersInFlight.Add(ctx, 1)
}
// EndInFlight 标记订单处理结束
func (m *Metrics) EndInFlight(ctx context.Context) {
m.ordersInFlight.Add(ctx, -1)
}
几种指标的选用规则我总结一下,记不住的时候回来翻:
- 想统计累计次数(请求数、错误数、订单数)→ Counter
- 想统计分布(延迟分位、响应大小分布)→ Histogram
- 想统计当前值且可增可减(队列长度、活跃连接数)→ UpDownCounter
- 想统计当前值且数据源在外部(连接池大小、goroutine 数)→ ObservableGauge
- 想周期性算个值上报(QPS、错误率)→ 没必要用 Gauge,直接用 Counter + PromQL rate() 算
⚠️ 新手必踩的坑:高基数 label 直接把 Prometheus 撑爆。
metric.WithAttributes(attribute.String("order_id", ...))这种写法,每个不同的 order_id 都会生成一条独立的时间序列。请求量一大,时间序列数量爆炸式增长,Prometheus 内存直接 OOM。label 只允许放低基数维度(result=success/error、status=5xx 这类),user_id、order_id 这种绝对不能进 metric。
踩坑提示:ObservableGauge 的回调函数在 SDK 采集周期(默认 15s)被调用一次,如果你在回调里去查 DB,每个采集周期都查一次,QPS 高的话 DB 扛不住。这种场景应该用后台 goroutine 周期更新一个内存变量,回调只读变量。还有一个坑:高基数 label(user_id、order_id)绝对不能进 metric,否则 Prometheus 直接 OOM。
Collector 完整配置
建立直觉:Collector 是一座"物流分拣中心"
类比:Collector 就像一个物流分拣中心。门口的
Receiver是收货口——应用把包裹(遥测数据)从 gRPC/HTTP 口送进来,hostmetrics 口还顺手收自己仓库的监控;中间的Processor是流水线上的分拣工——先限量防仓库爆仓(memory_limiter),再贴资源标签(resource),挑出垃圾件丢掉(filter),把敏感信息涂黑(attributes 脱敏),按规则决定哪些件要留(tail_sampling),最后打包(batch);末尾的Exporter是发货口——把包裹分别发往 Prometheus、Jaeger、Loki 三个不同的总仓。流水线顺序错了,就会出现"先打包又被拆开重新分拣"的诡异事故。
下面这张图把"接收 → 处理 → 导出"的流水线和三条数据管线画出来:
flowchart LR
subgraph 接收
R1[otlp]
R2[hostmetrics]
end
subgraph 处理
P1[memory_limiter] --> P2[resource] --> P3[filter] --> P4[attributes] --> P5[tail_sampling] --> P6[batch]
end
subgraph 导出
E1[prometheusremotewrite]
E2[otlp/jaeger]
E3[loki]
end
R1 --> P1
R2 --> P1
P6 --> E1
P6 --> E2
P6 --> E3生产级 Collector 配置
下面这份是生产级别的 Collector 配置,我加了详细注释,每个组件为啥这么配都写了。这份配置同时导出到 Prometheus(metrics)、Jaeger(traces)、Loki(logs)三个后端:
# otel-collector-config.yaml
# 用 otel/opentelemetry-collector-contrib 镜像跑,基础镜像不带这些组件
# ============ Receivers:数据入口 ============
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
# 关键:max_recv_msg_size 默认 4MB,大 trace 会超限
max_recv_msg_size: 16MB
http:
endpoint: 0.0.0.0:4318
# 顺手收 hostmetrics,监控 Collector 自己
hostmetrics:
collection_interval: 30s
scrapers:
cpu: {}
memory: {}
disk: {}
network: {}
# ============ Processors:数据处理 ============
processors:
# 1. memory_limiter 必须放最前面,防 OOM
memory_limiter:
check_interval: 1s
limit_mib: 1500
spike_limit_mib: 300
# 2. resource 补全资源属性
resource:
attributes:
- key: deployment.environment
value: production
action: upsert
- key: k8s.cluster.name
value: prod-bj
action: upsert
# 关键:删掉应用没用的属性,避免后端索引膨胀
- key: process.pid
action: delete
# 3. filter 过滤无用数据
filter:
traces:
span:
# 健康检查 span 不要,纯噪音
- 'attributes["http.route"] == "/healthz"'
- 'attributes["http.route"] == "/metrics"'
- 'attributes["http.route"] == "/ping"'
metrics:
metric:
# 系统指标里一些没用的也过滤掉
- 'name == "process.runtime.go.gc.heap_frees"'
# 4. attributes 脱敏
attributes:
actions:
# 删掉敏感请求头
- key: http.request.header.authorization
action: delete
- key: http.request.header.cookie
action: delete
# SQL 文本 hash 掉,不存明文
- key: db.statement
action: hash
# 5. tail_sampling 尾部采样
# 注意:用 tail_sampling 后 batch 必须放在它后面
tail_sampling:
decision_wait: 10s
num_traces: 50000
expected_new_traces_per_sec: 1000
policies:
# 错误 trace 100% 保留
- name: errors
type: status_code
status_code:
status_codes: [ERROR]
# 慢 trace 100% 保留
- name: slow
type: latency
latency:
threshold_ms: 500
# 正常 trace 1% 采样
- name: random_low
type: probabilistic
probabilistic:
sampling_percentage: 1
# 6. batch 批处理,放最后
batch:
timeout: 5s
send_batch_size: 1024
send_batch_max_size: 2048
# ============ Exporters:数据出口 ============
exporters:
# 1. Prometheus Remote Write 推 metrics
prometheusremotewrite:
endpoint: http://prometheus:9090/api/v1/write
# 关键:external labels 帮助区分集群
external_labels:
cluster: prod-bj
# 写失败重试
retry_on_failure:
enabled: true
initial_interval: 1s
max_interval: 30s
max_elapsed_time: 300s
# 2. OTLP 推 traces 到 Jaeger
# Jaeger 1.35+ 支持 OTLP 直接接收
otlp/jaeger:
endpoint: jaeger:4317
tls:
insecure: true
# 3. Loki 推 logs
# 注意:loki exporter 在 contrib 镜像里
loki:
endpoint: http://loki:3100/loki/api/v1/push
# 关键:default_labels_added 会给所有日志加标签,方便 LogQL 筛选
default_labels_added:
- exporter: OTLP
- environment: production
# 4. debug exporter,开发时用
debug:
verbosity: basic
sampling_initial: 2
sampling_thereafter: 1
# ============ Extensions:辅助组件 ============
extensions:
# 健康检查,K8s liveness probe 用
health_check:
endpoint: 0.0.0.0:13133
# pprof,性能分析用
pprof:
endpoint: 0.0.0.0:1777
# zpages,运行时调试页面
zpages:
endpoint: 0.0.0.0:55679
# ============ Service:流水线组装 ============
service:
extensions: [health_check, pprof, zpages]
telemetry:
logs:
level: info
metrics:
address: 0.0.0.0:8888
pipelines:
# metrics 流水线
metrics:
receivers: [otlp, hostmetrics]
processors: [memory_limiter, resource, filter, batch]
exporters: [prometheusremotewrite, debug]
# traces 流水线
# 关键:tail_sampling 必须在 batch 之前
traces:
receivers: [otlp]
processors: [memory_limiter, resource, filter, attributes, tail_sampling, batch]
exporters: [otlp/jaeger, debug]
# logs 流水线
logs:
receivers: [otlp]
processors: [memory_limiter, resource, batch]
exporters: [loki, debug]
⚠️ 新手必踩的坑:processor 顺序写反。
memory_limiter → resource → filter → attributes → tail_sampling → batch这条顺序是铁律。最典型的错误是把batch放到tail_sampling之前:数据先被打包发走了,tail_sampling 根本来不及看完整的一条 trace 就做决策,结果"该留的慢 trace 被丢掉、该采样的被全量发出去"。另一个常见错是 memory_limiter 没放最前,突发流量直接把 Collector 自己打爆。
processor 顺序我前面那篇也强调过,这里再说一遍因为太重要:memory_limiter → resource → filter → attributes → tail_sampling → batch。错了会出现诸如"batch 之后又被 tail_sampling 拆开"这种诡异问题。
K8s 部署的话,建议 Collector 用 Deployment + HPA,而不是 DaemonSet。DaemonSet 适合收节点指标,OTel Collector 主要收应用数据,跟节点没强绑定。HPA 按 CPU/内存扩容,QPS 高峰能自动扛住。
踩坑提示:tail_sampling 的 decision_wait 设太短会丢 trace(trace 没收完就决策了),设太长内存爆炸。10s 是个折中值,能覆盖大多数服务的请求生命周期。如果你有特别长的 trace(比如异步任务),单独配一条 pipeline 不采样。
业务全链路监控方案
实现全链路监控需要做到:
- 统一 TraceID 传递:在 HTTP/gRPC 请求头中传递
traceparent。 - 自动埋点框架:使用 otelgin、otelgrpc 等官方中间件。
- 数据库与缓存埋点:对 MySQL、Redis、Kafka 等组件进行埋点。
- 上下文透传:通过 Go 的
context.Context在不同函数间传递 Span。
建立直觉:Trace Context 是一场"跨服务接力赛"
类比:跨服务串联 trace,本质是一场接力赛。第一棒(api-gateway)手里拿着一根接力棒(trace context,核心是
traceparent这个请求头)。它把棒交给第二棒(order-service)时,不是喊一声,而是把棒塞进 HTTP 请求头里传过去;第二棒接到后开自己的新一棒(子 span),跑完再把棒塞进对 Redis 的调用里。只要每一棒交接时都老老实实"把棒放进请求头 / 从请求头取出",整条 trace 就不会断。负责"放棒 / 取棒"的,就是propagator。
下面这张时序图就是"用户请求跨三个服务"时,trace context 怎么一站站传下去、串成一条完整 trace 的:
sequenceDiagram
participant G as api-gateway
participant O as order-service
participant R as Redis
G->>G: 开 server span (GET /api/orders)
G->>O: client span 把 traceparent 注入请求头
O->>O: 开 server span (续上父上下文)
O->>R: client span 把 traceparent 注入请求头
R->>R: 开 server span
R-->>O: 返回数据
O-->>G: 返回数据光看条目容易,关键是把"一个请求跨三个服务"的 trace 怎么真正串起来。下面我用一个真实场景演示:用户调 api-gateway,gateway 调 order-service,order-service 调 redis。三个服务,每段代码我都贴出来。
服务 A:api-gateway(HTTP 入口)
// api-gateway: 接收外部请求,发起对下游的 HTTP 调用
package main
import (
"net/http"
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/propagation"
)
func main() {
// 关键:用 otelhttp.NewHandler 包整个 mux,所有入口请求自动建 server span
// span 名是 HTTP 方法+路径,比如 "GET /api/orders"
mux := http.NewServeMux()
mux.HandleFunc("/api/orders", proxyToOrderService)
handler := otelhttp.NewHandler(mux, "api-gateway",
otelhttp.WithTracerProvider(otel.GetTracerProvider()),
// 关键:用 composite propagator 同时支持 traceparent 和 baggage
otelhttp.WithPropagators(propagation.NewCompositeTextMapPropagator(
propagation.TraceContext{},
propagation.Baggage{},
)),
)
_ = http.ListenAndServe(":8080", handler)
}
func proxyToOrderService(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
// 关键:用 otelhttp.NewTransport 包 transport,下游调用自动建 client span
// client span 会把 trace context 注入到 HTTP 头里,下游服务就能续上
client := &http.Client{
Transport: otelhttp.NewTransport(http.DefaultTransport),
}
req, _ := http.NewRequestWithContext(ctx, "GET", "http://order-service:8080/orders", nil)
// 注意:不要手动设置 traceparent 头!
// otelhttp.Transport 会自动从 ctx 提取并注入
resp, err := client.Do(req)
if err != nil {
http.Error(w, err.Error(), 500)
return
}
defer resp.Body.Close()
// 复制响应...
}
服务 B:order-service(HTTP 服务 + Redis 客户端)
// order-service: 接收 gateway 请求,调 Redis
package main
import (
"net/http"
"go.opentelemetry.io/contrib/instrumentation/github.com/redis/go-redis/v9/redisotel"
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
"go.opentelemetry.io/otel"
"github.com/redis/go-redis/v9"
)
func main() {
// Redis 客户端:用 redisotel 包装,所有 Redis 命令自动建 span
rdb := redis.NewClient(&redis.Options{Addr: "redis:6379"})
// 关键:Hook 必须在客户端初始化时加,运行时不能改
rdb.AddHook(redisotel.NewTracingHook())
mux := http.NewServeMux()
mux.HandleFunc("/orders", func(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
// 这里 ctx 已经被 otelhttp 注入了 trace context
// 后续所有用这个 ctx 的调用(包括 Redis)都会自动续上 trace
val, err := rdb.Get(ctx, "order:123").Result()
if err != nil {
http.Error(w, err.Error(), 500)
return
}
w.Write([]byte(val))
})
handler := otelhttp.NewHandler(mux, "order-service")
_ = http.ListenAndServe(":8080", handler)
}
服务 C:异步消息(Kafka)场景
异步场景稍微复杂,因为生产者和消费者之间没有 HTTP 请求头传递 trace context。OTel 的做法是把 context 注入到 Kafka 消息头里:
类比:Kafka 这场合 HTTP 不一样——生产者和消费者不是"一手交棒一手接棒"的同步调用,而是中间隔了一个消息队列这个"传声筒"。所以接力的方式变成:生产者在发消息时,把接力棒(trace context)塞进消息头;消费者从消息头里把棒取出来,再开自己的 span 接着跑。本质还是"放棒 / 取棒",只是棒不再走 HTTP 头,改走 Kafka 消息头。
sequenceDiagram
participant P as 生产者
participant K as Kafka 消息
participant C as 消费者
P->>P: 开启 producer span
P->>K: Inject 把 trace context 注入消息头
K-->>C: 消息传递
C->>C: Extract 提取 context 开启 consumer span
C->>C: 业务处理 (续上原 trace)// 生产者:发消息时注入 trace context
func produce(ctx context.Context, writer *kafka.Writer, msgValue []byte) error {
ctx, span := otel.Tracer("producer").Start(ctx, "kafka.produce",
trace.WithSpanKind(trace.SpanKindProducer),
)
defer span.End()
msg := kafka.Message{Value: msgValue}
// 关键:把 trace context 注入到 Kafka 消息头
// 消费者读出来后能续上 trace
otel.GetTextMapPropagator().Inject(ctx, propagation.NewCarrierWriter(&msg.HeaderCarrier))
return writer.WriteMessages(ctx, msg)
}
// 消费者:读消息时提取 trace context
func consume(ctx context.Context, reader *kafka.Reader) error {
msg, err := reader.ReadMessage(ctx)
if err != nil {
return err
}
// 关键:从消息头提取 trace context,建立消费者 span
ctx = otel.GetTextMapPropagator().Extract(ctx, propagation.NewCarrierReader(&msg.HeaderCarrier))
ctx, span := otel.Tracer("consumer").Start(ctx, "kafka.consume",
trace.WithSpanKind(trace.SpanKindConsumer),
)
defer span.End()
// 后续业务处理就跟普通 HTTP 请求一样了
return processMessage(ctx, msg.Value)
}
这样在 Jaeger 里你看到的就是一条完整的 trace:gateway(server) → gateway(client) → order-service(server) → redis(client) → redis(server)。每一段 span 都有 parent_span_id 串起来。
⚠️ 新手必踩的坑:两边 propagator 不一致,trace 直接断。
api-gateway用TraceContext{}注入,下游却只配了Baggage{},那一头注入一头提取不出来,整条 trace 在跨服务处断裂,Jaeger 里看到的是两条互不相干的 trace。两端必须用同一套 propagator(通常都是TraceContext{}+Baggage{}的 composite)。另外如果中间过了 Nginx 反向代理,要确认没有用proxy_set_header把traceparent一起覆盖掉。
踩坑提示:跨服务 trace 串联失败的 90% 原因是 propagator 没设对。两边的 propagator 必须一致(都用 TraceContext{} 或者都用 Baggage{}),否则一头注入一头提取不出来。还有个坑:如果中间过了 Nginx 反向代理,Nginx 默认不会丢 traceparent 头,但如果你用了 proxy_set_header 重写了某些头,要确认没把 traceparent 一起覆盖掉。
基于可观测数据的 AIOps 异常检测
采集到 Metrics 后,可以基于时序数据进行异常检测。例如使用 Python 对 Prometheus 查询结果做简单阈值与趋势分析:
建立直觉:用"多维度互相印证"降低误报
类比:异常检测就像医生下诊断。只看"体温计"(单一指标维度),发烧可能是刚运动完,误报;只看"心率"(单一 trace 维度),心跳快可能是你刚爬完楼。但体温和心率同时异常,那基本就是真生病了。所以生产里我让"metric 错误率异常"和"trace 错误率异常"两个维度同时触发才算高优告警,任一维度单独异常只算低优提醒——告警噪音立刻降了 70%。下面这张图就是这套判定逻辑:
flowchart TD
M[查 Prometheus 错误率] --> Check1{metric 异常?}
T[查 Jaeger trace 错误率] --> Check2{trace 异常?}
Check1 -->|是| Alert[高优告警
两维度印证]
Check2 -->|是| Alert
Check1 -->|否| Low[低优提醒]
Check2 -->|否| Low# aiops/anomaly_detect.py
import requests
import numpy as np
from datetime import datetime, timedelta
PROM_BASE = "http://prometheus:9090"
def query_prom(query: str, t: datetime = None) -> list:
"""查 Prometheus,返回 series 列表"""
t = t or datetime.now()
params = {"query": query, "time": t.timestamp()}
r = requests.get(f"{PROM_BASE}/api/v1/query", params=params, timeout=10)
r.raise_for_status()
return r.json()["data"]["result"]
def query_range(query: str, start: datetime, end: datetime, step: str = "60s") -> list:
"""范围查询,返回每个 series 的时间序列"""
params = {
"query": query,
"start": start.timestamp(),
"end": end.timestamp(),
"step": step,
}
r = requests.get(f"{PROM_BASE}/api/v1/query_range", params=params, timeout=30)
r.raise_for_status()
return r.json()["data"]["result"]
def detect_metric_anomaly(values: list, threshold_sigma: float = 3.0) -> dict:
"""
3-sigma 异常检测:最近一个点偏离历史均值超过 N 个标准差就告警
返回 {is_anomaly, current, mean, std, zscore}
"""
if len(values) < 30:
# 关键:数据点太少不可靠,曾经我踩过坑:3 个点算 std 直接 0
return {"is_anomaly": False, "reason": "insufficient_data"}
arr = np.array(values, dtype=float)
history = arr[:-1] # 前 N-1 个做基线
current = arr[-1] # 最后一个判断
mean = float(np.mean(history))
std = float(np.std(history))
# 关键:std 为 0 时不能除,加个 epsilon
if std < 1e-9:
return {"is_anomaly": False, "reason": "zero_variance"}
zscore = abs(current - mean) / std
return {
"is_anomaly": zscore > threshold_sigma,
"current": float(current),
"mean": mean,
"std": std,
"zscore": float(zscore),
}
def check_service_error_rate(service: str) -> dict:
"""检查某个服务的错误率是否异常"""
now = datetime.now()
start = now - timedelta(hours=1)
# 关键:rate() 的窗口要够长,太短抖动大
# 5m 窗口是常见的折中
query = f'''
sum(rate(http_requests_total{{service="{service}",status=~"5.."}}[5m])) by (service)
/
sum(rate(http_requests_total{{service="{service}"}}[5m])) by (service)
'''
series = query_range(query, start, now, "60s")
if not series:
return {"service": service, "status": "no_data"}
results = []
for s in series:
values = [float(v[1]) for v in s["values"]]
anomaly = detect_metric_anomaly(values, threshold_sigma=3.0)
# 关键:错误率除了看突变,还要看绝对值
# 即使没突变,错误率超 1% 也得告警
if anomaly["is_anomaly"] or (anomaly.get("current", 0) > 0.01):
results.append({
"service": service,
"anomaly": anomaly,
"severity": "critical" if anomaly.get("current", 0) > 0.05 else "warning",
})
return {"service": service, "anomalies": results}
def check_trace_error_rate(service: str) -> dict:
"""检查 trace 错误率,从 Jaeger 查"""
# Jaeger 有 HTTP API 可以查 trace 统计
# 这里简化,实际可以用 Jaeger 的 /api/traces?service=X&tags={"error":true}
jaeger_url = f"http://jaeger:16686/api/traces"
params = {
"service": service,
"limit": 100,
"lookback": "1h",
}
r = requests.get(jaeger_url, params=params, timeout=10)
r.raise_for_status()
data = r.json()
traces = data.get("data", [])
if not traces:
return {"service": service, "status": "no_traces"}
error_count = 0
slow_count = 0
durations = []
for trace in traces:
# 关键:trace duration 是所有 span 里最长 end-start
duration_ms = max(
(span.get("duration", 0) for span in trace.get("spans", [])),
default=0,
)
durations.append(duration_ms)
# 错误判断:看是否有 error tag
for span in trace.get("spans", []):
for tag in span.get("tags", []):
if tag.get("key") == "error" and tag.get("value") is True:
error_count += 1
break
# 慢 trace 判断
if duration_ms > 500:
slow_count += 1
total = len(traces)
# 关键:trace 通常采样过,要按采样率反推真实量
# 假设采样率 10%,错误率 = error_count / total 仍然准确,因为是比例
return {
"service": service,
"total_traces": total,
"error_rate": error_count / total if total else 0,
"slow_rate": slow_count / total if total else 0,
"p99_duration_ms": float(np.percentile(durations, 99)) if durations else 0,
"is_anomaly": (error_count / total > 0.05) if total else False,
}
def run_detection(services: list) -> list:
"""对所有服务跑一遍异常检测"""
alerts = []
for svc in services:
# 关键:metric 和 trace 两个维度都看,互相印证
# 单一维度容易误报,两个维度都异常基本就是真异常
metric_result = check_service_error_rate(svc)
trace_result = check_trace_error_rate(svc)
if metric_result.get("anomalies") or trace_result.get("is_anomaly"):
alerts.append({
"service": svc,
"metric": metric_result,
"trace": trace_result,
"timestamp": datetime.now().isoformat(),
})
return alerts
if __name__ == "__main__":
services = ["order-service", "payment-service", "user-service"]
alerts = run_detection(services)
for a in alerts:
print(f"[ALERT] {a['service']}: metric={bool(a['metric'].get('anomalies'))} "
f"trace={a['trace'].get('is_anomaly')}")
更复杂的场景可以使用:
- Prophet:检测时序中的趋势变化与异常点。
- Isolation Forest:多维指标异常检测。
- LSTM Autoencoder:非线性时序异常检测。
结合 Trace 数据,还可以通过分析调用链延迟分布,定位异常服务;结合日志数据,可以对异常 Span 关联的日志进行聚类,辅助根因分析。
我特别想强调"多维度互相印证"这点。单一维度做异常检测误报率很高——比如 metric 错误率上涨,可能是上游流量暴涨导致绝对错误数变多但比例没变;trace 错误率高,可能是采样碰巧采到了几个错误请求。但 metric 和 trace 同时报警,那基本就是真问题了。生产里我把"两个维度同时异常"作为高优告警条件,“单一维度异常"作为低优提醒,告警噪音立刻降了 70%。
收尾
OpenTelemetry 提供了一条从应用到后端的完整可观测性链路。通过合理的埋点、Collector 配置和后端存储选择,可以构建覆盖 Metrics、Traces、Logs 的全链路监控体系。这些可观测数据是 AIOps 异常检测、根因定位和自动修复的基础。
写完这篇我自己回头看了下,代码量确实不小,但都是生产能跑的——不是那种"demo 级别"的玩具代码。如果你按这篇一步步接入,碰到问题大概率是这几个:propagator 没设一致导致 trace 断、shutdown 没 defer 导致丢数据、metric 高基数 label 把 Prometheus 撑爆。这三个坑我都踩过,希望你能绕过去。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- Trace、Metric、Log 三者分别适合回答哪一类问题?各举一个排查场景里的例子。
- 为什么
TracerProvider/MeterProvider的Shutdown必须调用?不调用会丢什么数据? - Collector 的 processor 顺序是
memory_limiter → resource → filter → attributes → tail_sampling → batch。如果把batch放到tail_sampling之前,会发生什么诡异问题? - 跨服务 trace 串联失败,90% 的原因是什么?如果 api-gateway 用
TraceContext{}而 order-service 只配了Baggage{},现象是什么? - 异常检测为什么要"metric + trace 两个维度同时异常"才算高优告警?只用单一维度有什么问题?
动手练习(建议真做一遍):
- 把示例 Order 服务的
SampleRatio从1.0改成0.2,跑起来后在 Jaeger 里观察采样后 trace 数量的变化,理解采样率对存储成本的影响。 - 给 Collector 加一条
tail_sampling策略:把http.route == "/healthz"的 span 100% 丢弃,并配合filterprocessor,验证"过滤 + 尾部采样"如何协作降低后端压力。 - 把 3-sigma 脚本的
threshold_sigma从3.0改成2.0,喂一段带人造尖刺的数据,观察zscore阈值降低后误报率如何变化。
本章小结
- OpenTelemetry 用一套标准把 Trace(流转单)/ Metric(仪表)/ Log(备注) 统一采集,经 Collector 分拣后落到 Prometheus / Jaeger / Loki 等后端。
- 应用侧核心是
TracerProvider+MeterProvider两个"总部”,以及Span这张"工单";Shutdown必须调用,否则 buffer 里的数据会丢。 - Metric 四种类型要按"累计 / 分布 / 当前值 / 外部数据源"选型,且 label 只能用低基数维度,高基数(user_id、order_id)会撑爆 Prometheus。
- Collector 的 processor 顺序是不可违背的铁律,
tail_sampling必须在batch之前。 - 跨服务串联靠 propagator 在请求头 / 消息头里"放棒 / 取棒",两端 propagator 必须一致。
- AIOps 异常检测用 metric + trace 双维度互相印证 来压误报,比单一维度稳得多。
下一篇我们会把这套可观测数据接到 LLM,做一个能自己判断故障、自己修复的 AIOps Operator——也就是把"发现异常"升级成"自动处理异常"。