Appearance
诊断监听器与 OpenTelemetry
概述
EF Core 通过 DiagnosticSource 和 DiagnosticListener 提供了一套完整的事件诊断系统,允许开发者订阅数据库操作事件,实现自定义监控、日志记录和性能分析。结合 OpenTelemetry,可以实现标准化的分布式追踪。
核心组件
| 组件 | 说明 |
|---|---|
| DiagnosticSource | .NET 标准的事件发布/订阅机制 |
| DiagnosticListener | 监听特定来源的诊断事件 |
| Activity | 表示一个操作单元,支持分布式追踪 |
| OpenTelemetry | CNCF 标准的可观测性框架 |
| Application Insights | Azure 监控服务(可选) |
可用事件列表
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: 42msOpenTelemetry 集成
安装 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:16686csharp
// 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 部署 |
核心要点
- DiagnosticListener: 底层事件订阅机制
- Activity: 分布式追踪标准
- OpenTelemetry: CNCF 标准,厂商中立
- 采样策略: 生产环境必须采样
- 敏感数据: 避免记录密码、令牌等