引言
在微服务架构中,一个用户请求可能跨越数十个服务。当出现性能问题或错误时,如何快速定位根因?分布式追踪提供了答案——它记录请求在系统中的完整路径,让我们能够可视化服务间的调用关系、识别性能瓶颈、追踪错误传播。
分布式追踪核心概念
追踪模型
Trace(追踪): 一次完整的请求链路
└─ Span(跨度): 单个操作或工作单元
├─ Trace ID: 全局唯一标识(如 4bf92f3577b34da6)
├─ Span ID: 当前Span的ID(如 00f067aa0ba902b7)
├─ Parent Span ID: 父Span的ID(如 463ac35c9f6413ad)
├─ Operation Name: 操作名称(如 "GET /api/users")
├─ Start Time: 开始时间
├─ Duration: 持续时间
├─ Tags: 键值对元数据
└─ Logs: 时间戳事件
追踪示例
用户请求: GET /api/orders/123
Trace ID: abc123
├─ Span 1: API Gateway (50ms)
│ ├─ Operation: "HTTP GET"
│ ├─ Tags: {http.method: "GET", http.url: "/api/orders/123"}
│ │
│ ├─ Span 2: Auth Service (10ms)
│ │ ├─ Operation: "validate_token"
│ │ └─ Tags: {user.id: "456"}
│ │
│ ├─ Span 3: Order Service (35ms)
│ │ ├─ Operation: "get_order"
│ │ │
│ │ ├─ Span 4: Database Query (20ms)
│ │ │ ├─ Operation: "SELECT"
│ │ │ └─ Tags: {db.system: "postgresql", db.statement: "SELECT * FROM orders"}
│ │ │
│ │ └─ Span 5: Cache Lookup (5ms)
│ │ ├─ Operation: "redis.get"
│ │ └─ Tags: {db.system: "redis"}
│ │
│ └─ Span 6: Response Serialization (5ms)
│ └─ Operation: "json.marshal"
OpenTelemetry集成
Go服务集成
package main
import (
"context"
"log"
"net/http"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp"
"go.opentelemetry.io/otel/propagation"
"go.opentelemetry.io/otel/sdk/resource"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
semconv "go.opentelemetry.io/otel/semconv/v1.21.0"
"go.opentelemetry.io/otel/trace"
)
func initTracer() (*sdktrace.TracerProvider, error) {
// 创建OTLP导出器
exporter, err := otlptracehttp.New(context.Background(),
otlptracehttp.WithEndpoint("otel-collector:4318"),
otlptracehttp.WithInsecure(),
)
if err != nil {
return nil, err
}
// 创建资源
res, err := resource.Merge(
resource.Default(),
resource.NewWithAttributes(
semconv.SchemaURL,
semconv.ServiceName("order-service"),
semconv.ServiceVersion("1.0.0"),
),
)
if err != nil {
return nil, err
}
// 创建追踪提供者
tp := sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exporter),
sdktrace.WithResource(res),
sdktrace.WithSampler(sdktrace.ParentBased(
sdktrace.TraceIDRatioBased(0.1), // 10%采样率
)),
)
// 设置全局追踪提供者
otel.SetTracerProvider(tp)
// 设置全局传播器(W3C Trace Context)
otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator(
propagation.TraceContext{},
propagation.Baggage{},
))
return tp, nil
}
func main() {
tp, err := initTracer()
if err != nil {
log.Fatal(err)
}
defer tp.Shutdown(context.Background())
http.HandleFunc("/api/orders", handleOrders)
http.ListenAndServe(":8080", nil)
}
HTTP中间件集成
package middleware
import (
"net/http"
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
)
// TracingMiddleware 自动追踪HTTP请求
func TracingMiddleware(serviceName string) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return otelhttp.NewHandler(next, serviceName,
otelhttp.WithSpanNameFormatter(func(operation string, r *http.Request) string {
return r.Method + " " + r.URL.Path
}),
otelhttp.WithPropagators(otel.GetTextMapPropagator()),
)
}
}
// CustomTracingMiddleware 自定义追踪逻辑
func CustomTracingMiddleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
tracer := otel.Tracer("order-service")
// 从请求上下文提取或创建新的Span
ctx, span := tracer.Start(r.Context(), "HTTP "+r.Method,
trace.WithAttributes(
attribute.String("http.method", r.Method),
attribute.String("http.url", r.URL.String()),
attribute.String("http.user_agent", r.UserAgent()),
),
trace.WithSpanKind(trace.SpanKindServer),
)
defer span.End()
// 添加用户信息(如果已认证)
userID := getUserIDFromContext(ctx)
if userID != "" {
span.SetAttributes(attribute.String("user.id", userID))
}
// 包装ResponseWriter以捕获状态码
wrapped := &responseWriter{ResponseWriter: w, statusCode: http.StatusOK}
// 调用下一个处理器
next.ServeHTTP(wrapped, r.WithContext(ctx))
// 记录响应状态码
span.SetAttributes(attribute.Int("http.status_code", wrapped.statusCode))
if wrapped.statusCode >= 500 {
span.SetStatus(trace.StatusCodeError, "Server error")
}
})
}
type responseWriter struct {
http.ResponseWriter
statusCode int
}
func (rw *responseWriter) WriteHeader(code int) {
rw.statusCode = code
rw.ResponseWriter.WriteHeader(code)
}
数据库追踪
package database
import (
"context"
"database/sql"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
)
type TracedDB struct {
*sql.DB
tracer trace.Tracer
}
func NewTracedDB(db *sql.DB) *TracedDB {
return &TracedDB{
DB: db,
tracer: otel.Tracer("database"),
}
}
func (db *TracedDB) QueryContext(ctx context.Context, query string, args ...interface{}) (*sql.Rows, error) {
ctx, span := db.tracer.Start(ctx, "SQL Query",
trace.WithAttributes(
attribute.String("db.system", "postgresql"),
attribute.String("db.statement", query),
attribute.Int("db.args_count", len(args)),
),
trace.WithSpanKind(trace.SpanKindClient),
)
defer span.End()
start := time.Now()
rows, err := db.DB.QueryContext(ctx, query, args...)
duration := time.Since(start)
span.SetAttributes(attribute.Int64("db.duration_ms", duration.Milliseconds()))
if err != nil {
span.RecordError(err)
span.SetStatus(trace.StatusCodeError, err.Error())
}
return rows, err
}
func (db *TracedDB) ExecContext(ctx context.Context, query string, args ...interface{}) (sql.Result, error) {
ctx, span := db.tracer.Start(ctx, "SQL Exec",
trace.WithAttributes(
attribute.String("db.system", "postgresql"),
attribute.String("db.statement", query),
),
)
defer span.End()
result, err := db.DB.ExecContext(ctx, query, args...)
if err != nil {
span.RecordError(err)
span.SetStatus(trace.StatusCodeError, err.Error())
} else if rowsAffected, _ := result.RowsAffected(); rowsAffected > 0 {
span.SetAttributes(attribute.Int64("db.rows_affected", rowsAffected))
}
return result, err
}
gRPC追踪
package grpc
import (
"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
"google.golang.org/grpc"
)
// 客户端拦截器
func NewClientConn(address string) (*grpc.ClientConn, error) {
return grpc.Dial(address,
grpc.WithUnaryInterceptor(otelgrpc.UnaryClientInterceptor()),
grpc.WithStreamInterceptor(otelgrpc.StreamClientInterceptor()),
)
}
// 服务器拦截器
func NewServer() *grpc.Server {
return grpc.NewServer(
grpc.UnaryInterceptor(otelgrpc.UnaryServerInterceptor()),
grpc.StreamInterceptor(otelgrpc.StreamServerInterceptor()),
)
}
追踪上下文传播
HTTP Header传播
package propagation
import (
"context"
"net/http"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/propagation"
)
// 在服务间传递追踪上下文
func CallDownstreamService(ctx context.Context, url string) error {
req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
if err != nil {
return err
}
// 注入追踪上下文到HTTP Header
otel.GetTextMapPropagator().Inject(ctx, propagation.HeaderCarrier(req.Header))
// Header中会包含:
// traceparent: 00-4bf92f3577b34da6a3b78e5c-00f067aa0ba902b7-01
// tracestate: congo=t61rcWkgMzE
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
return nil
}
// 从HTTP请求提取追踪上下文
func ExtractTraceContext(r *http.Request) context.Context {
ctx := r.Context()
return otel.GetTextMapPropagator().Extract(ctx, propagation.HeaderCarrier(r.Header))
}
消息队列传播
package messaging
import (
"context"
"github.com/segmentio/kafka-go"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/propagation"
)
// Kafka生产者:注入追踪上下文
func ProduceMessage(ctx context.Context, writer *kafka.Writer, key, value []byte) error {
msg := kafka.Message{
Key: key,
Value: value,
}
// 将追踪上下文注入到Kafka Header
carrier := propagation.MapCarrier{}
otel.GetTextMapPropagator().Inject(ctx, carrier)
for k, v := range carrier {
msg.Headers = append(msg.Headers, kafka.Header{
Key: k,
Value: []byte(v),
})
}
return writer.WriteMessages(ctx, msg)
}
// Kafka消费者:提取追踪上下文
func ConsumeMessage(ctx context.Context, msg kafka.Message) (context.Context, error) {
// 从Kafka Header提取追踪上下文
carrier := propagation.MapCarrier{}
for _, header := range msg.Headers {
carrier.Set(header.Key, string(header.Value))
}
ctx = otel.GetTextMapPropagator().Extract(ctx, carrier)
return ctx, nil
}
OpenTelemetry Collector部署
Collector配置
# otel-collector-config.yaml
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
http:
endpoint: 0.0.0.0:4318
jaeger:
protocols:
grpc:
endpoint: 0.0.0.0:14250
thrift_http:
endpoint: 0.0.0.0:14268
processors:
batch:
timeout: 1s
send_batch_size: 1024
memory_limiter:
check_interval: 1s
limit_mib: 4000
spike_limit_mib: 800
attributes:
actions:
- key: environment
value: production
action: upsert
exporters:
otlp:
endpoint: tempo:4317
tls:
insecure: true
jaeger:
endpoint: jaeger:14250
tls:
insecure: true
logging:
loglevel: debug
service:
pipelines:
traces:
receivers: [otlp, jaeger]
processors: [memory_limiter, batch, attributes]
exporters: [otlp, jaeger, logging]
Kubernetes部署
# otel-collector-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: otel-collector
spec:
replicas: 2
selector:
matchLabels:
app: otel-collector
template:
metadata:
labels:
app: otel-collector
spec:
containers:
- name: otel-collector
image: otel/opentelemetry-collector:0.88.0
args:
- --config=/etc/otel-collector-config.yaml
ports:
- containerPort: 4317 # OTLP gRPC
- containerPort: 4318 # OTLP HTTP
resources:
requests:
memory: "512Mi"
cpu: "500m"
limits:
memory: "2Gi"
cpu: "2000m"
volumeMounts:
- name: config
mountPath: /etc/otel-collector-config.yaml
subPath: otel-collector-config.yaml
volumes:
- name: config
configMap:
name: otel-collector-config
---
apiVersion: v1
kind: Service
metadata:
name: otel-collector
spec:
selector:
app: otel-collector
ports:
- name: otlp-grpc
port: 4317
targetPort: 4317
- name: otlp-http
port: 4318
targetPort: 4318
追踪后端部署
Jaeger部署
# jaeger-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: jaeger
spec:
replicas: 1
selector:
matchLabels:
app: jaeger
template:
metadata:
labels:
app: jaeger
spec:
containers:
- name: jaeger
image: jaegertracing/all-in-one:1.50
env:
- name: COLLECTOR_OTLP_ENABLED
value: "true"
- name: SPAN_STORAGE_TYPE
value: "elasticsearch"
- name: ES_SERVER_URLS
value: "http://elasticsearch:9200"
ports:
- containerPort: 16686 # UI
- containerPort: 14250 # gRPC
- containerPort: 14268 # HTTP
---
apiVersion: v1
kind: Service
metadata:
name: jaeger
spec:
selector:
app: jaeger
ports:
- name: ui
port: 16686
targetPort: 16686
- name: grpc
port: 14250
targetPort: 14250
Grafana Tempo部署
# tempo-config.yaml
server:
http_listen_port: 3200
distributor:
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
http:
endpoint: 0.0.0.0:4318
ingester:
trace_idle_period: 10s
max_block_bytes: 1_000_000
max_block_duration: 5m
compactor:
compaction:
compaction_window: 1h
max_block_bytes: 100_000_000
block_retention: 48h
storage:
trace:
backend: s3
s3:
bucket: tempo-traces
endpoint: minio:9000
access_key: minioadmin
secret_key: minioadmin
insecure: true
采样策略
采样配置
package tracing
import (
sdktrace "go.opentelemetry.io/otel/sdk/trace"
)
// 1. 基于比率的采样(简单场景)
func NewRatioSampler(rate float64) sdktrace.Sampler {
return sdktrace.TraceIDRatioBased(rate)
}
// 2. 基于父级的采样(推荐)
func NewParentBasedSampler(rootSampler sdktrace.Sampler) sdktrace.Sampler {
return sdktrace.ParentBased(rootSampler)
}
// 3. 自定义采样器(高级场景)
type CustomSampler struct {
// 错误请求100%采样
// 慢请求100%采样
// 普通请求1%采样
}
func (s *CustomSampler) ShouldSample(p sdktrace.SamplingParameters) sdktrace.SamplingResult {
// 检查是否包含错误属性
for _, attr := range p.Attributes {
if attr.Key == "error" && attr.Value.AsBool() {
return sdktrace.SamplingResult{
Decision: sdktrace.RecordAndSample,
}
}
}
// 检查是否慢请求
for _, attr := range p.Attributes {
if attr.Key == "http.duration_ms" && attr.Value.AsInt64() > 1000 {
return sdktrace.SamplingResult{
Decision: sdktrace.RecordAndSample,
}
}
}
// 普通请求1%采样
if rand.Float64() < 0.01 {
return sdktrace.SamplingResult{
Decision: sdktrace.RecordAndSample,
}
}
return sdktrace.SamplingResult{
Decision: sdktrace.Drop,
}
}
func (s *CustomSampler) Description() string {
return "CustomSampler"
}
动态采样(Collector配置)
# 在otel-collector-config.yaml中添加
processors:
tail_sampling:
policies:
# 错误请求100%保留
- name: errors-policy
type: status_code
status_code: {status_codes: [ERROR]}
# 慢请求100%保留
- name: latency-policy
type: latency
latency: {threshold_ms: 1000}
# 特定端点100%保留
- name: critical-endpoints
type: string_attribute
string_attribute:
key: http.url
values: [/api/checkout, /api/payment]
# 其他请求1%采样
- name: default
type: probabilistic
probabilistic: {sampling_percentage: 1}
Span 事件与结构化日志
Span 不仅记录起止时间,还可以通过 Events 记录内部关键节点:
func ProcessPayment(ctx context.Context, orderID string) error {
ctx, span := tracer.Start(ctx, "payment.process")
defer span.End()
// 记录关键事件(时间戳自动附加)
span.AddEvent("validation.started",
trace.WithAttributes(attribute.String("order.id", orderID)))
if err := validateOrder(orderID); err != nil {
span.RecordError(err)
span.SetStatus(trace.StatusCodeError, err.Error())
return err
}
span.AddEvent("validation.completed",
trace.WithAttributes(attribute.Int64("validation.ms", 15)))
span.AddEvent("payment.gateway.called",
trace.WithAttributes(attribute.String("gateway", "stripe")))
resp, err := callStripeAPI(orderID)
if err != nil {
span.RecordError(err)
return err
}
span.AddEvent("payment.gateway.completed",
trace.WithAttributes(
attribute.String("stripe.charge_id", resp.ChargeID),
attribute.Int64("gateway.latency_ms", resp.LatencyMS),
))
return nil
}
Span Events vs Logs:
- Span Events:与特定 Span 绑定,带时间戳,适合标记 Span 内部里程碑
- 结构化 Logs:独立输出,通过 Trace ID 关联,适合记录业务事件
理想方案:两者共存,OpenTelemetry Collector 统一收集后关联展示。
Baggage:跨服务上下文传递
Baggage 是 OpenTelemetry 提供的键值对传播机制,伴随请求在整个链路中传递,类似 HTTP Header 但由 SDK 自动处理注入和提取。
import "go.opentelemetry.io/otel/baggage"
// 服务 A:写入 Baggage
func HandleRequest(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
// 创建 Baggage 成员(自动注入到后续所有出站请求)
member, _ := baggage.NewMember("tenant.id", "acme-corp")
member2, _ := baggage.NewMember("user.tier", "enterprise")
b, _ := baggage.New(member, member2)
ctx = baggage.ContextWithBaggage(ctx, b)
// 后续所有 HTTP/gRPC 调用自动携带 tenant.id 和 user.tier
callServiceB(ctx)
}
// 服务 B:读取 Baggage
func serviceBHandler(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
b := baggage.FromContext(ctx)
tenantID := b.Member("tenant.id").Value() // "acme-corp"
userTier := b.Member("user.tier").Value() // "enterprise"
// 根据租户隔离数据或路由到不同数据库分片
db := getDBForTenant(tenantID)
// ...
}
适用场景:租户标识、用户等级、请求来源渠道、A/B 测试分组等需要在全链路透传但不在业务参数中显式传递的元信息。
安全注意:Baggage 通过 Header 传播,所有中间环节都可见,禁止传递敏感信息(如用户密码、Token)。
链路拓扑与依赖分析
分布式追踪的另一大价值是揭示服务间的调用拓扑和健康状态:
依赖图谱关键指标
| 指标 | 含义 | 健康阈值 |
|---|---|---|
| 调用量 (Call/S) | 每分钟服务间调用次数 | 无固定阈值,关注趋势 |
| 错误率 | 失败请求占比 | < 0.1% |
| P99 延迟 | 99% 请求的最大延迟 | 因服务而异 |
| 依赖深度 | 单次请求跨服务层数 | < 5 层为佳 |
Jaeger 依赖图
# Jaeger 自动生成服务依赖图(基于 span 的 service.name 和 parent/child 关系)
# 访问 http://jaeger:16686/dependencies 查看拓扑
# API 获取依赖数据
curl "http://jaeger:16686/api/dependencies?endTs=$(date +%s)000&lookback=86400000"
三大支柱统一:Trace × Metrics × Logs
可观测性三大支柱各自独立有局限,OpenTelemetry 统一语义后可以实现联动分析:
关联模型
Trace Context(Trace ID / Span ID)
│
├──▶ Traces(请求链路)— Jaeger/Tempo/Zipkin
│
├──▶ Metrics(指标聚合)— Prometheus + Grafana
│ Exemplar: Trace ID 附着到 histogram 桶
│
└──▶ Logs(结构化日志)— Loki/ELK
每条日志携带 trace_id / span_id 字段
Exemplar:从指标到追踪的跳转
// Prometheus Exemplar:在直方图桶上附加 Trace ID
import "github.com/prometheus/client_golang/prometheus"
histogram := prometheus.NewHistogramVec(prometheus.HistogramOpts{
Name: "http_request_duration_seconds",
Buckets: prometheus.DefBuckets,
}, []string{"method", "endpoint"})
func handler(w http.ResponseWriter, r *http.Request) {
start := time.Now()
ctx, span := tracer.Start(r.Context(), "http.handler")
defer span.End()
// 执行业务逻辑...
duration := time.Since(start).Seconds()
// 记录 Exemplar:将 Trace ID 附加到当前 bucket
traceID := span.SpanContext().TraceID().String()
histogram.WithLabelValues(r.Method, r.URL.Path).(prometheus.ExemplarAdder).AddWithExemplar(
duration,
prometheus.Labels{"trace_id": traceID},
)
}
在 Grafana 中点击 Prometheus 直方图的某个异常桶,可直接跳转到对应 Trace 详情。
统一日志输出
// Zap/Logrus 集成 OTel Trace ID,实现日志与追踪联动
import "go.uber.org/zap"
func WithTraceLogging(ctx context.Context, log *zap.Logger) *zap.Logger {
span := trace.SpanFromContext(ctx)
if span.SpanContext().IsValid() {
return log.With(
zap.String("trace_id", span.SpanContext().TraceID().String()),
zap.String("span_id", span.SpanContext().SpanID().String()),
zap.Bool("trace_sampled", span.SpanContext().IsSampled()),
)
}
return log
}
// 使用
logger.Info("payment processed",
zap.String("order_id", orderID),
zap.Float64("amount", amount),
)
// 输出:{"level":"info","ts":...,"msg":"payment processed","trace_id":"abc123","span_id":"def456","trace_sampled":true,"order_id":"ORD-789","amount":99.99}
性能优化
批处理与缓冲
// 配置批处理导出
tp := sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exporter,
sdktrace.WithMaxExportBatchSize(512),
sdktrace.WithBatchTimeout(5*time.Second),
sdktrace.WithMaxQueueSize(2048),
),
)
减少Span创建
// 避免在热路径创建过多Span
func ProcessItems(ctx context.Context, items []Item) error {
// ✗ 错误:每个item都创建Span
for _, item := range items {
ctx, span := tracer.Start(ctx, "process_item")
processItem(item)
span.End()
}
// ✓ 正确:为整个批次创建一个Span
ctx, span := tracer.Start(ctx, "process_items_batch",
trace.WithAttributes(
attribute.Int("items.count", len(items)),
),
)
defer span.End()
for _, item := range items {
processItem(item)
}
return nil
}
总结
分布式追踪的核心价值:
- 快速定位问题:可视化请求路径,一眼看到瓶颈
- 性能优化:识别慢调用和性能瓶颈
- 错误追踪:追踪错误在系统中的传播
- 依赖分析:了解服务间的调用关系
实施要点:
- 统一使用OpenTelemetry标准
- 确保追踪上下文正确传播
- 合理配置采样策略
- 选择合适的追踪后端
- 持续优化追踪数据质量
延伸阅读
- OpenTelemetry Documentation
- Dapper - Google’s Tracing System
- Jaeger Documentation
- Grafana Tempo
- Zipkin Documentation
- The Three Pillars of Observability
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。