OpenTelemetry 开发实战

2025-11-12T14:11:02+08:00 | 27分钟阅读 | 更新于 2025-11-12T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 用自己的话讲清 OpenTelemetry 中 Trace / Metric / Log 三者各自回答什么问题、分别适合什么场景,面试时能讲成一张图。
  2. 在 Go 服务里完整接入 OTel:初始化 TracerProvider + MeterProvider、打 Span 记录业务、记四种 Metric,并正确处理优雅退出(Shutdown)。
  3. 读懂并写出生产级 Collector 配置:Receiver / Processor / Exporter 流水线,以及那条绝不能乱的 processor 顺序。
  4. 实现跨服务、跨协议的 trace 串联(HTTP / gRPC / Kafka),讲清 propagator 为什么是串联成败的关键。
  5. 基于采集到的可观测数据,写一段多维度互相印证的异常检测脚本,并说清楚为什么单一维度误报率高。

前置知识

  • 能把 Go 服务跑起来(go run / go build),看得懂包、接口、struct。
  • context.Context 是怎么在多个函数之间传递的——本文大量代码都靠它串上下文。
  • 知道 Docker Compose 的基本用法(起几个容器、映射端口)。
  • 大致知道 Prometheus(指标存储)、Jaeger(链路存储)、Kafka(消息队列)是干嘛的,不需要很熟。

本章你会动手做的事

  1. 用 docker-compose 一把拉起 Collector + Prometheus + Jaeger 全家桶。
  2. 把示例 Go 代码里的 Provider 初始化、Span、Metric 跑通,并在 Jaeger 看到一条真实 trace。
  3. 写一个 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 之前,先分清两个角色。TracerProviderMeterProvider 像是你们公司的"总部"——它管着跟 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.getdb.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_samplingdecision_wait 设太短会丢 trace(trace 没收完就决策了),设太长内存爆炸。10s 是个折中值,能覆盖大多数服务的请求生命周期。如果你有特别长的 trace(比如异步任务),单独配一条 pipeline 不采样。

业务全链路监控方案

实现全链路监控需要做到:

  1. 统一 TraceID 传递:在 HTTP/gRPC 请求头中传递 traceparent
  2. 自动埋点框架:使用 otelgin、otelgrpc 等官方中间件。
  3. 数据库与缓存埋点:对 MySQL、Redis、Kafka 等组件进行埋点。
  4. 上下文透传:通过 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-gatewayTraceContext{} 注入,下游却只配了 Baggage{},那一头注入一头提取不出来,整条 trace 在跨服务处断裂,Jaeger 里看到的是两条互不相干的 trace。两端必须用同一套 propagator(通常都是 TraceContext{} + Baggage{} 的 composite)。另外如果中间过了 Nginx 反向代理,要确认没有用 proxy_set_headertraceparent 一起覆盖掉。

踩坑提示:跨服务 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 撑爆。这三个坑我都踩过,希望你能绕过去。

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. Trace、Metric、Log 三者分别适合回答哪一类问题?各举一个排查场景里的例子。
  2. 为什么 TracerProvider / MeterProviderShutdown 必须调用?不调用会丢什么数据?
  3. Collector 的 processor 顺序是 memory_limiter → resource → filter → attributes → tail_sampling → batch。如果把 batch 放到 tail_sampling 之前,会发生什么诡异问题?
  4. 跨服务 trace 串联失败,90% 的原因是什么?如果 api-gateway 用 TraceContext{} 而 order-service 只配了 Baggage{},现象是什么?
  5. 异常检测为什么要"metric + trace 两个维度同时异常"才算高优告警?只用单一维度有什么问题?

动手练习(建议真做一遍)

  1. 把示例 Order 服务的 SampleRatio1.0 改成 0.2,跑起来后在 Jaeger 里观察采样后 trace 数量的变化,理解采样率对存储成本的影响。
  2. 给 Collector 加一条 tail_sampling 策略:把 http.route == "/healthz" 的 span 100% 丢弃,并配合 filter processor,验证"过滤 + 尾部采样"如何协作降低后端压力。
  3. 把 3-sigma 脚本的 threshold_sigma3.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——也就是把"发现异常"升级成"自动处理异常"。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!