Skip to content

诊断监听器与 OpenTelemetry ​

概述 ​

EF Core 通过 DiagnosticSource 和 DiagnosticListener 提供了一套完整的事件诊断系统,允许开发者订阅数据库操作事件,实现自定义监控、日志记录和性能分析。结合 OpenTelemetry,可以实现标准化的分布式追踪。

核心组件 ​

组件说明
DiagnosticSource.NET 标准的事件发布/订阅机制
DiagnosticListener监听特定来源的诊断事件
Activity表示一个操作单元,支持分布式追踪
OpenTelemetryCNCF 标准的可观测性框架
Application InsightsAzure 监控服务(可选)

可用事件列表 ​

csharp
// EF Core 发布的主要事件
Microsoft.EntityFrameworkCore.Database.Command.CommandExecuting
Microsoft.EntityFrameworkCore.Database.Command.CommandExecuted
Microsoft.EntityFrameworkCore.Database.Command.CommandError
Microsoft.EntityFrameworkCore.Database.Connection.ConnectionOpening
Microsoft.EntityFrameworkCore.Database.Connection.ConnectionOpened
Microsoft.EntityFrameworkCore.Database.Connection.ConnectionClosing
Microsoft.EntityFrameworkCore.Database.Connection.ConnectionClosed
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionStarting
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionStarted
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionCommitting
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionCommitted
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionRollingBack
Microsoft.EntityFrameworkCore.Database.Transaction.TransactionRolledBack
Microsoft.EntityFrameworkCore.Query.QueryExecutionPlanned
Microsoft.EntityFrameworkCore.ChangeTracker.DetectChangesStarting
Microsoft.EntityFrameworkCore.ChangeTracker.DetectChangesCompleted

基础用法 ​

创建诊断监听器 ​

csharp
// EfCoreDiagnosticObserver.cs
using System.Diagnostics;

public class EfCoreDiagnosticObserver : IObserver<KeyValuePair<string, object?>>
{
    private readonly ILogger<EfCoreDiagnosticObserver> _logger;

    public EfCoreDiagnosticObserver(ILogger<EfCoreDiagnosticObserver> logger)
    {
        _logger = logger;
    }

    public void OnCompleted()
    {
        _logger.LogInformation("Diagnostic observer completed");
    }

    public void OnError(Exception error)
    {
        _logger.LogError(error, "Diagnostic observer error");
    }

    public void OnNext(KeyValuePair<string, object?> value)
    {
        switch (value.Key)
        {
            case "Microsoft.EntityFrameworkCore.Database.Command.CommandExecuting":
                HandleCommandExecuting(value.Value);
                break;
                
            case "Microsoft.EntityFrameworkCore.Database.Command.CommandExecuted":
                HandleCommandExecuted(value.Value);
                break;
                
            case "Microsoft.EntityFrameworkCore.Database.Command.CommandError":
                HandleCommandError(value.Value);
                break;
                
            case "Microsoft.EntityFrameworkCore.Query.QueryExecutionPlanned":
                HandleQueryPlanned(value.Value);
                break;
        }
    }

    private void HandleCommandExecuting(object? payload)
    {
        if (payload == null) return;
        
        var command = GetProperty<DbCommand>(payload, "Command");
        var connection = GetProperty<DbConnection>(payload, "Connection");
        var startTime = GetProperty<DateTime>(payload, "StartTime");
        
        if (command != null)
        {
            _logger.LogInformation("""
                SQL Command Executing:
                Connection: {Connection}
                Command: {CommandText}
                Parameters: {Parameters}
                StartTime: {StartTime}
                """,
                connection?.ConnectionString,
                command.CommandText,
                string.Join(", ", command.Parameters.Cast<DbParameter>()
                    .Select(p => $"{p.ParameterName}={p.Value}")),
                startTime);
        }
    }

    private void HandleCommandExecuted(object? payload)
    {
        if (payload == null) return;
        
        var command = GetProperty<DbCommand>(payload, "Command");
        var duration = GetProperty<TimeSpan>(payload, "Duration");
        
        if (command != null && duration.HasValue)
        {
            _logger.LogInformation("""
                SQL Command Executed:
                Duration: {Duration}ms
                Command: {CommandText}
                """,
                duration.Value.TotalMilliseconds,
                command.CommandText);
            
            // 检测慢查询
            if (duration.Value > TimeSpan.FromSeconds(1))
            {
                _logger.LogWarning("SLOW QUERY: {Duration}ms - {Command}", 
                    duration.Value.TotalMilliseconds, 
                    command.CommandText);
            }
        }
    }

    private void HandleCommandError(object? payload)
    {
        if (payload == null) return;
        
        var command = GetProperty<DbCommand>(payload, "Command");
        var exception = GetProperty<Exception>(payload, "Exception");
        
        _logger.LogError(exception, """
            SQL Command Error:
            Command: {CommandText}
            """,
            command?.CommandText);
    }

    private void HandleQueryPlanned(object? payload)
    {
        if (payload == null) return;
        
        var queryContext = GetProperty<object>(payload, "QueryContext");
        var queryPlan = GetProperty<string>(payload, "QueryPlan");
        
        _logger.LogDebug("Query Execution Plan: {Plan}", queryPlan);
    }

    private static T? GetProperty<T>(object obj, string propertyName)
    {
        var type = obj.GetType();
        var property = type.GetProperty(propertyName);
        return property != null ? (T?)property.GetValue(obj) : default;
    }
}

注册监听器 ​

csharp
// Program.cs
using System.Diagnostics;

var builder = WebApplication.CreateBuilder(args);

// 注册诊断监听器
builder.Services.AddSingleton<EfCoreDiagnosticObserver>();

builder.Services.AddDbContext<AppDbContext>((sp, options) =>
{
    options.UseSqlServer(builder.Configuration.GetConnectionString("DefaultConnection"));
    
    // 订阅诊断事件
    var observer = sp.GetRequiredService<EfCoreDiagnosticObserver>();
    var subscription = DiagnosticListener.AllListeners.Subscribe(new DiagnosticObserverWrapper(observer));
});

var app = builder.Build();
app.Run();

// 包装器(处理订阅生命周期)
public class DiagnosticObserverWrapper : IObserver<DiagnosticListener>
{
    private readonly IObserver<KeyValuePair<string, object?>> _observer;
    private readonly List<IDisposable> _subscriptions = new();

    public DiagnosticObserverWrapper(IObserver<KeyValuePair<string, object?>> observer)
    {
        _observer = observer;
    }

    public void OnCompleted()
    {
        foreach (var sub in _subscriptions)
        {
            sub.Dispose();
        }
    }

    public void OnError(Exception error)
    {
    }

    public void OnNext(DiagnosticListener value)
    {
        if (value.Name == "Microsoft.EntityFrameworkCore")
        {
            var subscription = value.Subscribe(_observer);
            _subscriptions.Add(subscription);
        }
    }
}

Activity 与分布式追踪 ​

启用 Activity 支持 ​

csharp
// Program.cs
using System.Diagnostics;

// 配置 ActivitySource
AppContext.SetSwitch("System.Diagnostics.DiagnosticSource.ActivityTrackingOptions", 
    (int)ActivityTrackingOptions.SpanId | 
    (int)ActivityTrackingOptions.TraceId | 
    (int)ActivityTrackingOptions.ParentId);

builder.Services.AddOpenTelemetry()
    .WithTracing(tracing =>
    {
        tracing.AddSource("MyApp")
               .AddEntityFrameworkCoreInstrumentation()
               .AddAspNetCoreInstrumentation()
               .AddConsoleExporter();  // 开发环境输出到控制台
    });

自定义 Activity ​

csharp
// ProductService.cs
public class ProductService
{
    private static readonly ActivitySource ActivitySource = new("MyApp.ProductService");
    private readonly AppDbContext _context;

    public async Task<List<Product>> GetProductsAsync()
    {
        using var activity = ActivitySource.StartActivity("GetProducts");
        activity?.SetTag("operation", "database_query");
        activity?.SetTag("entity", "Product");
        
        try
        {
            var products = await _context.Products.ToListAsync();
            
            activity?.SetTag("result_count", products.Count);
            activity?.SetStatus(ActivityStatusCode.Ok);
            
            return products;
        }
        catch (Exception ex)
        {
            activity?.SetStatus(ActivityStatusCode.Error, ex.Message);
            activity?.RecordException(ex);
            throw;
        }
    }

    public async Task<int> CreateProductAsync(Product product)
    {
        using var activity = ActivitySource.StartActivity("CreateProduct");
        activity?.SetTag("product_name", product.Name);
        activity?.SetTag("product_price", product.Price);
        
        try
        {
            _context.Products.Add(product);
            await _context.SaveChangesAsync();
            
            activity?.SetTag("product_id", product.Id);
            activity?.SetStatus(ActivityStatusCode.Ok);
            
            return product.Id;
        }
        catch (Exception ex)
        {
            activity?.SetStatus(ActivityStatusCode.Error, ex.Message);
            activity?.RecordException(ex);
            throw;
        }
    }
}

查看 Trace 输出 ​

// 控制台输出示例:
Activity.TraceId:          0af7651916cd43dd8448eb211c80319c
Activity.SpanId:           b7ad6b7169203331
Activity.TraceFlags:       Recorded
Activity.ActivitySourceName: MyApp.ProductService
Activity.DisplayName:      GetProducts
Activity.Kind:             Internal
Activity.StartTime:        2026-04-06T10:30:45.1234567Z
Activity.Duration:         00:00:00.0450000
Activity.Tags:
    operation: database_query
    entity: Product
    result_count: 100
Associated Instruments:
    System.Data.SqlClient.SqlCommand.ExecuteAsync

-- EF Core SQL:
SELECT [p].[Id], [p].[Name], [p].[Price]
FROM [Products] AS [p]
Duration: 42ms

OpenTelemetry 集成 ​

安装 NuGet 包 ​

bash
dotnet add package OpenTelemetry.Extensions.Hosting
dotnet add package OpenTelemetry.Instrumentation.EntityFrameworkCore
dotnet add package OpenTelemetry.Exporter.Console
dotnet add package OpenTelemetry.Exporter.OpenTelemetryProtocol

完整配置 ​

csharp
// Program.cs
using OpenTelemetry.Resources;
using OpenTelemetry.Trace;

var builder = WebApplication.CreateBuilder(args);

// 配置 OpenTelemetry
builder.Services.AddOpenTelemetry()
    .ConfigureResource(resource => resource
        .AddService(
            serviceName: "MyApp",
            serviceVersion: "1.0.0",
            serviceInstanceId: Environment.MachineName))
    .WithTracing(tracing =>
    {
        tracing
            // 添加 EF Core  instrumentation
            .AddEntityFrameworkCoreInstrumentation(options =>
            {
                options.SetDbStatementForText = true;  // 记录 SQL 文本
                options.EnrichWithIDbCommand = (activity, command) =>
                {
                    activity.SetTag("db.command_type", command.CommandType.ToString());
                    activity.SetTag("db.command_timeout", command.CommandTimeout);
                };
            })
            
            // 添加 ASP.NET Core instrumentation
            .AddAspNetCoreInstrumentation(options =>
            {
                options.RecordException = true;
                options.EnrichWithHttpResponse = (activity, response) =>
                {
                    activity.SetTag("http.response.status_code", response.StatusCode);
                };
            })
            
            // 添加 HttpClient instrumentation
            .AddHttpClientInstrumentation()
            
            // 控制台导出器(开发环境)
            .AddConsoleExporter()
            
            // OTLP 导出器(生产环境 - 发送到 Jaeger/Zipkin/Prometheus)
            .AddOtlpExporter(options =>
            {
                options.Endpoint = new Uri(builder.Configuration["OTEL_ENDPOINT"] 
                    ?? "http://localhost:4317");
            });
    });

builder.Services.AddDbContext<AppDbContext>(options =>
    options.UseSqlServer(builder.Configuration.GetConnectionString("DefaultConnection")));

var app = builder.Build();
app.Run();

发送到 Jaeger ​

bash
# Docker 运行 Jaeger
docker run -d --name jaeger \
  -e COLLECTOR_ZIPKIN_HOST_PORT=:9411 \
  -p 5775:5775/udp \
  -p 6831:6831/udp \
  -p 6832:6832/udp \
  -p 5778:5778 \
  -p 16686:16686 \
  -p 14268:14268 \
  -p 14250:14250 \
  -p 9411:9411 \
  jaegertracing/all-in-one:latest

# 访问 UI: http://localhost:16686
csharp
// appsettings.json
{
  "OTEL_ENDPOINT": "http://localhost:4317"
}

// Program.cs
.AddOtlpExporter(options =>
{
    options.Endpoint = new Uri(builder.Configuration["OTEL_ENDPOINT"]);
})

Application Insights 集成 ​

安装 NuGet 包 ​

bash
dotnet add package Microsoft.ApplicationInsights.AspNetCore
dotnet add package Microsoft.ApplicationInsights.DependencyCollector

配置 ​

csharp
// Program.cs
builder.Services.AddApplicationInsightsTelemetry(options =>
{
    options.ConnectionString = builder.Configuration["APPLICATIONINSIGHTS_CONNECTION_STRING"];
});

builder.Services.AddOpenTelemetry()
    .WithTracing(tracing =>
    {
        tracing.AddEntityFrameworkCoreInstrumentation()
               .AddApplicationInsightsExporter();  // 发送到 AI
    });

自定义遥测 ​

csharp
// TelemetryService.cs
public class TelemetryService
{
    private readonly TelemetryClient _telemetryClient;

    public TelemetryService(TelemetryClient telemetryClient)
    {
        _telemetryClient = telemetryClient;
    }

    public void TrackDatabaseOperation(string operation, TimeSpan duration, bool success)
    {
        var properties = new Dictionary<string, string>
        {
            ["operation"] = operation,
            ["success"] = success.ToString(),
            ["duration_ms"] = duration.TotalMilliseconds.ToString()
        };

        var metrics = new Dictionary<string, double>
        {
            ["duration"] = duration.TotalMilliseconds
        };

        _telemetryClient.TrackEvent("DatabaseOperation", properties, metrics);
    }

    public void TrackSlowQuery(string sql, TimeSpan duration)
    {
        _telemetryClient.TrackTrace($"SLOW QUERY ({duration.TotalMilliseconds}ms): {sql}", 
            SeverityLevel.Warning);
    }
}

实际应用场景 ​

场景 1: 性能监控仪表板 ​

csharp
// MetricsCollector.cs
public class MetricsCollector : IObserver<KeyValuePair<string, object?>>
{
    private readonly ConcurrentDictionary<string, long> _queryCounts = new();
    private readonly ConcurrentDictionary<string, long> _totalDurations = new();

    public void OnNext(KeyValuePair<string, object?> value)
    {
        if (value.Key == "Microsoft.EntityFrameworkCore.Database.Command.CommandExecuted")
        {
            var command = GetProperty<DbCommand>(value.Value, "Command");
            var duration = GetProperty<TimeSpan>(value.Value, "Duration");
            
            if (command != null && duration.HasValue)
            {
                var operation = ExtractOperationName(command.CommandText);
                
                _queryCounts.AddOrUpdate(operation, 1, (k, v) => v + 1);
                _totalDurations.AddOrUpdate(operation, 
                    duration.Value.Ticks, 
                    (k, v) => v + duration.Value.Ticks);
            }
        }
    }

    public Dictionary<string, QueryStats> GetStatistics()
    {
        return _queryCounts.Keys.Select(key => new
        {
            Operation = key,
            Count = _queryCounts[key],
            AvgDuration = TimeSpan.FromTicks(_totalDurations[key] / _queryCounts[key])
        })
        .ToDictionary(
            x => x.Operation,
            x => new QueryStats 
            { 
                Count = x.Count, 
                AvgDuration = x.AvgDuration 
            });
    }

    private string ExtractOperationName(string sql)
    {
        if (sql.StartsWith("SELECT", StringComparison.OrdinalIgnoreCase))
            return "SELECT";
        if (sql.StartsWith("INSERT", StringComparison.OrdinalIgnoreCase))
            return "INSERT";
        if (sql.StartsWith("UPDATE", StringComparison.OrdinalIgnoreCase))
            return "UPDATE";
        if (sql.StartsWith("DELETE", StringComparison.OrdinalIgnoreCase))
            return "DELETE";
        return "OTHER";
    }
}

public record QueryStats
{
    public long Count { get; init; }
    public TimeSpan AvgDuration { get; init; }
}

// API 端点: 获取统计数据
app.MapGet("/api/metrics/database", (MetricsCollector collector) =>
{
    return Results.Ok(collector.GetStatistics());
});

场景 2: 审计日志 ​

csharp
// AuditLogger.cs
public class AuditLogger : IObserver<KeyValuePair<string, object?>>
{
    private readonly ILogger<AuditLogger> _logger;
    private readonly IHttpContextAccessor _httpContextAccessor;

    public AuditLogger(
        ILogger<AuditLogger> logger,
        IHttpContextAccessor httpContextAccessor)
    {
        _logger = logger;
        _httpContextAccessor = httpContextAccessor;
    }

    public void OnNext(KeyValuePair<string, object?> value)
    {
        if (value.Key == "Microsoft.EntityFrameworkCore.Database.Command.CommandExecuted")
        {
            var command = GetProperty<DbCommand>(value.Value, "Command");
            var duration = GetProperty<TimeSpan>(value.Value, "Duration");
            
            if (command != null)
            {
                var userId = _httpContextAccessor.HttpContext?.User?.FindFirst("sub")?.Value 
                    ?? "anonymous";
                
                _logger.LogInformation("""
                    AUDIT: Database Operation
                    User: {UserId}
                    Operation: {Operation}
                    Duration: {Duration}ms
                    SQL: {Sql}
                    Timestamp: {Timestamp}
                    """,
                    userId,
                    command.CommandType,
                    duration?.TotalMilliseconds ?? 0,
                    SanitizeSql(command.CommandText),
                    DateTime.UtcNow);
            }
        }
    }

    private string SanitizeSql(string sql)
    {
        // 移除敏感数据(密码、令牌等)
        return Regex.Replace(sql, @"'(password|token|secret)'[^']*'", "'***'", 
            RegexOptions.IgnoreCase);
    }
}

最佳实践 ​

✅ 推荐做法 ​

1. 生产环境使用采样 ​

csharp
.AddOtlpExporter(options =>
{
    options.Endpoint = new Uri(configuration["OTEL_ENDPOINT"]);
})
.SetSampler(new ParentBasedSampler(
    new TraceIdRatioBasedSampler(0.1)));  // 10% 采样

2. 避免记录敏感数据 ​

csharp
// ❌ 错误: 记录完整 SQL(可能包含敏感数据)
activity.SetTag("db.statement", command.CommandText);

// ✅ 正确: 仅记录操作类型
activity.SetTag("db.operation", ExtractOperationType(command.CommandText));

3. 设置合理的超时 ​

csharp
.AddOtlpExporter(options =>
{
    options.TimeoutMilliseconds = 5000;  // 5 秒超时
})

❌ 避免的错误 ​

1. 不要在监听器中执行耗时操作 ​

csharp
// ❌ 错误
public void OnNext(KeyValuePair<string, object?> value)
{
    // ⚠️ 同步 HTTP 调用
    httpClient.PostAsync("http://logging-service", content).Wait();
}

// ✅ 正确: 异步批量发送
private readonly Channel<LogEntry> _logChannel = Channel.CreateBounded<LogEntry>(1000);

public void OnNext(KeyValuePair<string, object?> value)
{
    _logChannel.Writer.TryWrite(new LogEntry(...));
}

2. 不要忘记取消订阅 ​

csharp
// ❌ 错误: 内存泄漏
var subscription = DiagnosticListener.AllListeners.Subscribe(observer);
// 忘记 Dispose

// ✅ 正确
public class ManagedObserver : IDisposable
{
    private readonly IDisposable _subscription;

    public ManagedObserver()
    {
        _subscription = DiagnosticListener.AllListeners.Subscribe(this);
    }

    public void Dispose()
    {
        _subscription.Dispose();
    }
}

总结 ​

技术选型对比 ​

方案复杂度功能适用场景
DiagnosticListener⭐ 简单基础监控小型应用
Activity + Console⭐⭐ 中等本地调试开发环境
OpenTelemetry⭐⭐⭐ 复杂标准化追踪生产环境
Application Insights⭐⭐ 中等Azure 生态Azure 部署

核心要点 ​

  1. DiagnosticListener: 底层事件订阅机制
  2. Activity: 分布式追踪标准
  3. OpenTelemetry: CNCF 标准,厂商中立
  4. 采样策略: 生产环境必须采样
  5. 敏感数据: 避免记录密码、令牌等

基于 MIT 许可发布