📖 数据密集型设计

分布式追踪与链路监控

深入探讨分布式追踪与链路监控技术

一、分布式追踪概述

分布式追踪是在分布式系统中追踪请求完整执行路径的技术,能够帮助开发者理解系统行为、定位性能瓶颈、排查故障。在数据密集型应用中,分布式追踪是保障系统可靠性和性能的关键工具。

二、分布式追踪原理

2.1 追踪数据模型

graph TD A[Trace] --> B[Span1: API Gateway] A --> C[Span2: Service A] A --> D[Span3: Service B] A --> E[Span4: Database] B --> C C --> D D --> E B --> B1[Duration: 100ms] C --> C1[Duration: 200ms] D --> D1[Duration: 300ms] E --> E1[Duration: 400ms] B --> B2[Tags: http.method=GET] C --> C2[Tags: db.type=mysql]

2.2 追踪数据模型表

概念 描述 作用
Trace 一次完整请求的追踪 关联所有相关Span
Span 单个操作的追踪 记录操作详情
Trace ID 唯一标识一次追踪 跨服务关联
Span ID 唯一标识一个Span 定位具体操作
Parent Span ID 父Span的ID 构建调用关系
Tags 键值对元数据 添加上下文信息
Logs 时间点日志 记录关键事件

三、分布式追踪实现

3.1 追踪上下文传播

public class TraceContext
{
    public string TraceId { get; set; }
    public string SpanId { get; set; }
    public string ParentSpanId { get; set; }
    public Dictionary<string, string> Tags { get; set; } = new Dictionary<string, string>();
    public List<LogEntry> Logs { get; set; } = new List<LogEntry>();
    
    public static TraceContext FromHttpHeaders(HttpRequest headers)
    {
        return new TraceContext
        {
            TraceId = headers["trace-id"] ?? GenerateTraceId(),
            SpanId = headers["span-id"] ?? GenerateSpanId(),
            ParentSpanId = headers["parent-span-id"]
        };
    }
    
    public void InjectToHttpHeaders(HttpRequestMessage request)
    {
        request.Headers.Add("trace-id", TraceId);
        request.Headers.Add("span-id", SpanId);
        
        if (!string.IsNullOrEmpty(ParentSpanId))
        {
            request.Headers.Add("parent-span-id", ParentSpanId);
        }
    }
    
    private static string GenerateTraceId()
    {
        return Guid.NewGuid().ToString("N");
    }
    
    private static string GenerateSpanId()
    {
        return Guid.NewGuid().ToString("N").Substring(0, 16);
    }
}

3.2 追踪拦截器

public class TracingMiddleware
{
    private readonly RequestDelegate _next;
    private readonly ITracer _tracer;
    
    public async Task InvokeAsync(HttpContext context)
    {
        var traceContext = TraceContext.FromHttpHeaders(context.Request);
        
        using var span = _tracer.StartSpan(context.Request.Path, traceContext);
        
        span.SetTag("http.method", context.Request.Method);
        span.SetTag("http.path", context.Request.Path);
        span.SetTag("http.query", context.Request.QueryString.ToString());
        span.SetTag("client.ip", context.Connection.RemoteIpAddress?.ToString());
        
        try
        {
            await _next(context);
            
            span.SetTag("http.status_code", context.Response.StatusCode);
        }
        catch (Exception ex)
        {
            span.SetTag("error", true);
            span.Log(ex.Message);
            throw;
        }
    }
}

3.3 分布式追踪服务

public class DistributedTracingService
{
    public async Task<TResult> ExecuteWithTracingAsync<TResult>(
        string operationName,
        Func<TraceContext, Task<TResult>> operation)
    {
        var traceContext = new TraceContext
        {
            TraceId = TraceContext.GenerateTraceId(),
            SpanId = TraceContext.GenerateSpanId()
        };
        
        using var span = _tracer.StartSpan(operationName, traceContext);
        
        try
        {
            return await operation(traceContext);
        }
        catch (Exception ex)
        {
            span.SetTag("error", true);
            span.Log(ex.Message);
            throw;
        }
    }
    
    public async Task ExecuteWithTracingAsync(
        string operationName,
        Func<TraceContext, Task> operation)
    {
        var traceContext = new TraceContext
        {
            TraceId = TraceContext.GenerateTraceId(),
            SpanId = TraceContext.GenerateSpanId()
        };
        
        using var span = _tracer.StartSpan(operationName, traceContext);
        
        try
        {
            await operation(traceContext);
        }
        catch (Exception ex)
        {
            span.SetTag("error", true);
            span.Log(ex.Message);
            throw;
        }
    }
    
    public async Task<TraceResult> GetTraceAsync(string traceId)
    {
        return await _traceRepository.GetTraceAsync(traceId);
    }
    
    public async Task<List<TraceResult>> SearchTracesAsync(TraceQuery query)
    {
        return await _traceRepository.SearchTracesAsync(query);
    }
}

四、Jaeger分布式追踪

4.1 Jaeger架构

graph TD A[客户端] --> B[Jaeger Agent] B --> C[Jaeger Collector] C --> D[Jaeger Query] D --> E[Jaeger UI] C --> F[存储: Cassandra/Elasticsearch] D --> F G[服务A] --> B H[服务B] --> B I[服务C] --> B

4.2 Jaeger配置

public class JaegerConfigurationService
{
    public ITracer ConfigureJaeger(string serviceName)
    {
        var sampler = new ConstSampler(true);
        
        var reporter = new RemoteReporter.Builder()
            .WithSender(new UdpSender("localhost", 6831))
            .Build();
        
        var tracer = new Tracer.Builder(serviceName)
            .WithSampler(sampler)
            .WithReporter(reporter)
            .Build();
        
        GlobalTracer.Register(tracer);
        
        return tracer;
    }
    
    public async Task SendTraceAsync(TraceContext traceContext)
    {
        var span = _tracer.BuildSpan(traceContext.SpanId)
            .AsChildOf(new SpanContext(traceContext.TraceId, traceContext.ParentSpanId))
            .Start();
        
        foreach (var tag in traceContext.Tags)
        {
            span.SetTag(tag.Key, tag.Value);
        }
        
        foreach (var log in traceContext.Logs)
        {
            span.Log(log.Timestamp, log.Fields);
        }
        
        span.Finish();
    }
}

4.3 Jaeger追踪示例

public class OrderService
{
    public async Task<Order> CreateOrderAsync(OrderRequest request)
    {
        using var span = _tracer.BuildSpan("CreateOrder").Start();
        
        span.SetTag("order.items_count", request.Items.Count);
        span.SetTag("order.total_amount", request.TotalAmount);
        
        try
        {
            var inventorySpan = _tracer.BuildSpan("CheckInventory")
                .AsChildOf(span.Context)
                .Start();
            
            await _inventoryService.CheckInventoryAsync(request.Items);
            inventorySpan.Finish();
            
            var paymentSpan = _tracer.BuildSpan("ProcessPayment")
                .AsChildOf(span.Context)
                .Start();
            
            await _paymentService.ProcessPaymentAsync(request);
            paymentSpan.Finish();
            
            var order = await _orderRepository.CreateAsync(request);
            span.SetTag("order.id", order.Id);
            
            return order;
        }
        catch (Exception ex)
        {
            span.SetTag("error", true);
            span.Log(ex.Message);
            throw;
        }
    }
}

五、Zipkin分布式追踪

5.1 Zipkin架构

graph TD A[客户端] --> B[Zipkin Collector] B --> C[Zipkin Storage] C --> D[Zipkin Query API] D --> E[Zipkin UI] F[服务A] --> B G[服务B] --> B H[服务C] --> B

5.2 Zipkin配置

public class ZipkinConfigurationService
{
    public void ConfigureZipkin(IApplicationBuilder app)
    {
        app.UseZipkinTracing(new ZipkinOptions
        {
            ServiceName = "OrderService",
            ZipkinEndpoint = new Uri("http://localhost:9411/api/v2/spans"),
            SampleRate = 1.0
        });
    }
    
    public async Task SendSpanAsync(Span span)
    {
        using var httpClient = new HttpClient();
        
        var spanData = new ZipkinSpan
        {
            TraceId = span.TraceId,
            SpanId = span.SpanId,
            ParentSpanId = span.ParentSpanId,
            Name = span.Name,
            Timestamp = span.Timestamp.ToUnixTimeMilliseconds(),
            Duration = span.Duration.TotalMilliseconds,
            Tags = span.Tags,
            Logs = span.Logs.Select(l => new ZipkinLog
            {
                Timestamp = l.Timestamp.ToUnixTimeMilliseconds(),
                Fields = l.Fields
            }).ToList()
        };
        
        await httpClient.PostAsJsonAsync("http://localhost:9411/api/v2/spans", spanData);
    }
}

六、链路监控

6.1 链路监控指标

public class TraceMetrics
{
    public string TraceId { get; set; }
    public long TotalDurationMs { get; set; }
    public int SpanCount { get; set; }
    public int ServiceCount { get; set; }
    public List<SpanMetrics> SpanMetrics { get; set; } = new List<SpanMetrics>();
    public bool HasError { get; set; }
}

public class SpanMetrics
{
    public string SpanId { get; set; }
    public string ServiceName { get; set; }
    public string OperationName { get; set; }
    public long DurationMs { get; set; }
    public double DurationPercentage { get; set; }
    public bool HasError { get; set; }
}

public class TraceMonitor
{
    public async Task<TraceMetrics> GetTraceMetricsAsync(string traceId)
    {
        var trace = await _traceRepository.GetTraceAsync(traceId);
        
        return new TraceMetrics
        {
            TraceId = traceId,
            TotalDurationMs = trace.Spans.Max(s => s.EndTime) - trace.Spans.Min(s => s.StartTime),
            SpanCount = trace.Spans.Count,
            ServiceCount = trace.Spans.Select(s => s.ServiceName).Distinct().Count(),
            SpanMetrics = trace.Spans.Select(s => new SpanMetrics
            {
                SpanId = s.SpanId,
                ServiceName = s.ServiceName,
                OperationName = s.OperationName,
                DurationMs = s.EndTime - s.StartTime,
                HasError = s.Tags.ContainsKey("error")
            }).ToList(),
            HasError = trace.Spans.Any(s => s.Tags.ContainsKey("error"))
        };
    }
    
    public async Task MonitorAsync()
    {
        var slowTraces = await _traceRepository.GetSlowTracesAsync(1000);
        
        foreach (var trace in slowTraces)
        {
            await _alertService.SendAlert("慢追踪", 
                $"TraceID: {trace.TraceId}, 耗时: {trace.Duration}ms");
        }
        
        var errorTraces = await _traceRepository.GetErrorTracesAsync();
        
        foreach (var trace in errorTraces)
        {
            await _alertService.SendAlert("错误追踪", 
                $"TraceID: {trace.TraceId}, 错误: {trace.Error}");
        }
    }
}

6.2 链路可视化

public class TraceVisualizationService
{
    public string GenerateMermaidGraph(TraceResult trace)
    {
        var sb = new StringBuilder();
        sb.AppendLine("graph TD");
        
        var spans = trace.Spans.OrderBy(s => s.StartTime).ToList();
        
        foreach (var span in spans)
        {
            var nodeId = $"span_{span.SpanId.Replace("-", "_")}";
            sb.AppendLine($"    {nodeId}[{span.ServiceName}: {span.OperationName} ({span.Duration}ms)]");
            
            if (!string.IsNullOrEmpty(span.ParentSpanId))
            {
                var parentId = $"span_{span.ParentSpanId.Replace("-", "_")}";
                sb.AppendLine($"    {parentId} --> {nodeId}");
            }
            
            if (span.Tags.ContainsKey("error"))
            {
                sb.AppendLine($"    style {nodeId} fill:#ff4d4f");
            }
        }
        
        return sb.ToString();
    }
    
    public async Task<TraceResult> GetTraceAsync(string traceId)
    {
        return await _traceRepository.GetTraceAsync(traceId);
    }
}

七、采样策略

7.1 采样策略对比

策略 描述 优点 缺点 适用场景
全量采样 采样所有请求 完整追踪 开销大 开发测试
固定比例 按比例采样 可控开销 可能错过重要请求 生产环境
基于速率 按速率限制 控制总量 突发流量可能被丢弃 高流量场景
基于规则 按规则采样 精准采样 规则复杂 特定场景

7.2 采样策略实现

public class TraceSampler
{
    private readonly SamplingStrategy _strategy;
    private readonly double _rate;
    private readonly int _maxTracesPerSecond;
    private int _tracesThisSecond;
    private DateTime _lastSecond;
    
    public bool ShouldSample(TraceContext context)
    {
        return _strategy switch
        {
            SamplingStrategy.Always => true,
            SamplingStrategy.Never => false,
            SamplingStrategy.FixedRate => ShouldSampleByRate(),
            SamplingStrategy.RateLimit => ShouldSampleByRateLimit(),
            SamplingStrategy.RuleBased => ShouldSampleByRules(context),
            _ => false
        };
    }
    
    private bool ShouldSampleByRate()
    {
        return new Random().NextDouble() < _rate;
    }
    
    private bool ShouldSampleByRateLimit()
    {
        var now = DateTime.Now;
        
        if (now.Second != _lastSecond.Second)
        {
            _tracesThisSecond = 0;
            _lastSecond = now;
        }
        
        if (_tracesThisSecond < _maxTracesPerSecond)
        {
            _tracesThisSecond++;
            return true;
        }
        
        return false;
    }
    
    private bool ShouldSampleByRules(TraceContext context)
    {
        if (context.Tags.ContainsKey("error"))
            return true;
        
        if (context.Tags.TryGetValue("http.status_code", out var statusCode))
        {
            if (statusCode.StartsWith("5"))
                return true;
        }
        
        if (context.Tags.TryGetValue("http.path", out var path))
        {
            if (path.Contains("/api/v1/orders"))
                return true;
        }
        
        return ShouldSampleByRate();
    }
}

public enum SamplingStrategy
{
    Always,
    Never,
    FixedRate,
    RateLimit,
    RuleBased
}

八、分布式追踪最佳实践

8.1 全链路追踪

确保请求从入口到出口都有完整的追踪信息。

8.2 合理采样

根据业务需求和系统负载选择合适的采样策略。

8.3 添加有意义的Tags

在关键Span中添加有意义的Tags,便于问题排查。

8.4 日志关联

public class TracingLogger
{
    public void LogInformation(string message, TraceContext context)
    {
        _logger.LogInformation("[Trace={TraceId}, Span={SpanId}] {Message}", 
            context.TraceId, context.SpanId, message);
    }
    
    public void LogError(Exception ex, TraceContext context)
    {
        _logger.LogError(ex, "[Trace={TraceId}, Span={SpanId}] {Message}", 
            context.TraceId, context.SpanId, ex.Message);
    }
    
    public void LogWarning(string message, TraceContext context)
    {
        _logger.LogWarning("[Trace={TraceId}, Span={SpanId}] {Message}", 
            context.TraceId, context.SpanId, message);
    }
}

8.5 监控与告警

建立追踪监控和告警机制,及时发现异常。

九、总结

分布式追踪是数据密集型应用中保障系统可靠性和性能的关键技术。Jaeger和Zipkin是主流的分布式追踪系统,能够帮助开发者理解系统行为、定位性能瓶颈、排查故障。合理选择采样策略、添加有意义的Tags、建立监控告警机制,能够充分发挥分布式追踪的价值。