乐观锁 - 并发冲突处理
目录
1. 概述
1.1 什么是并发冲突?
在使用乐观锁时,多个请求同时修改同一条记录会导致并发冲突:
时间线:
T1: 请求A读取商品库存(stock=10, version=1)
T2: 请求B读取商品库存(stock=10, version=1)
T3: 请求A扣减库存(stock=9, version=2)✓ 成功
T4: 请求B扣减库存(WHERE stock=10 AND version=1)✗ 失败!version已变为2结果:请求B抛出 DbUpdateConcurrencyException
1.2 为什么需要处理?
- 数据一致性:防止丢失更新
- 用户体验:优雅处理冲突,避免直接报错
- 系统稳定性:防止异常导致服务中断
2. 并发冲突场景
2.1 典型场景
| 场景 | 冲突原因 | 解决方案 |
|---|---|---|
| 库存扣减 | 多人同时购买同一商品 | 重试或返回缺货 |
| 余额更新 | 同时充值/消费 | 重试确保最终一致 |
| 状态修改 | 多人同时审批 | 业务规则判断 |
| 计数器 | 高频累加操作 | 数据库原子操作 |
2.2 冲突检测原理
csharp
// EF Core 生成的 SQL
UPDATE products
SET stock = @p0, version = @p1
WHERE id = @p2 AND version = @p3; -- 关键:版本号检查
-- 如果受影响行数为 0,说明版本号不匹配,抛出并发异常3. EF Core 并发异常处理
3.1 基础处理
csharp
public class ProductService
{
private readonly OrderDbContext _dbContext;
public async Task<Result> DeductStockAsync(Guid productId, int quantity)
{
try
{
var product = await _dbContext.Products.FindAsync(productId);
if (product.Stock < quantity)
{
return Result.Fail("Insufficient stock");
}
product.Stock -= quantity;
product.Version++; // 乐观锁版本号
await _dbContext.SaveChangesAsync(); // 可能抛出并发异常
return Result.Success();
}
catch (DbUpdateConcurrencyException ex)
{
// 处理并发冲突
return HandleConcurrencyConflict(ex, productId);
}
}
private Result HandleConcurrencyConflict(
DbUpdateConcurrencyException ex,
Guid productId)
{
// 获取冲突的实体
var entry = ex.Entries.Single();
var databaseValues = entry.GetDatabaseValues();
var currentValues = entry.CurrentValues;
// 重新加载最新数据
entry.Reload();
// 根据业务逻辑决定如何处理
var updatedProduct = entry.Entity as Product;
if (updatedProduct.Stock < ((Product)currentValues.ToObject()).Stock)
{
return Result.Fail("Stock was reduced by another transaction");
}
// 可以选择重试
return RetryDeduct(productId, updatedProduct);
}
}3.2 提取冲突信息
csharp
public class ConcurrencyConflictInfo
{
public string TableName { get; set; } = string.Empty;
public Dictionary<string, object> OriginalValues { get; set; } = new();
public Dictionary<string, object> CurrentValues { get; set; } = new();
public Dictionary<string, object> DatabaseValues { get; set; } = new();
}
public static class ConcurrencyHelper
{
public static ConcurrencyConflictInfo ExtractConflictInfo(
DbUpdateConcurrencyException ex)
{
var entry = ex.Entries.Single();
var info = new ConcurrencyConflictInfo
{
TableName = entry.Entity.GetType().Name
};
// 原始值(客户端读取时的值)
foreach (var prop in entry.OriginalValues.Properties)
{
info.OriginalValues[prop.Name] = entry.OriginalValues[prop];
}
// 当前值(客户端尝试设置的值)
foreach (var prop in entry.CurrentValues.Properties)
{
info.CurrentValues[prop.Name] = entry.CurrentValues[prop];
}
// 数据库值(当前的实际值)
var databaseValues = entry.GetDatabaseValues();
foreach (var prop in databaseValues.Properties)
{
info.DatabaseValues[prop.Name] = databaseValues[prop];
}
return info;
}
}4. 重试策略
4.1 Polly 重试策略
csharp
using Polly;
public class ResilientProductService
{
private readonly OrderDbContext _dbContext;
private readonly IAsyncPolicy<Result> _retryPolicy;
public ResilientProductService(OrderDbContext dbContext)
{
_dbContext = dbContext;
// 配置重试策略
_retryPolicy = Policy<Result>
.Handle<DbUpdateConcurrencyException>()
.WaitAndRetryAsync(
retryCount: 3,
sleepDurationProvider: (retryAttempt, context) =>
{
// 指数退避:100ms, 200ms, 400ms
return TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100);
},
onRetry: (result, timeSpan, retryCount, context) =>
{
Console.WriteLine($"Retry {retryCount} after concurrency conflict");
});
}
public async Task<Result> DeductStockWithRetryAsync(
Guid productId, int quantity)
{
return await _retryPolicy.ExecuteAsync(async () =>
{
var product = await _dbContext.Products.FindAsync(productId);
if (product.Stock < quantity)
{
return Result.Fail("Insufficient stock");
}
product.Stock -= quantity;
product.Version++;
await _dbContext.SaveChangesAsync();
return Result.Success();
});
}
}4.2 自定义重试上下文
csharp
public class RetryContext
{
public Guid ProductId { get; set; }
public int RequestedQuantity { get; set; }
public int RetryCount { get; set; }
}
public async Task<Result> DeductStockWithContextAsync(
Guid productId, int quantity)
{
var context = new Context();
context["productId"] = productId;
context["quantity"] = quantity;
return await _retryPolicy.ExecuteAsync(
async ctx =>
{
var pid = (Guid)ctx["productId"];
var qty = (int)ctx["quantity"];
return await DeductStockInternalAsync(pid, qty);
},
context);
}4.3 手动重试实现
csharp
public class ManualRetryService
{
private readonly OrderDbContext _dbContext;
private readonly int _maxRetries = 3;
public async Task<Result> ExecuteWithManualRetryAsync<T>(
Func<Task<Result>> operation)
{
Exception? lastException = null;
for (int i = 0; i < _maxRetries; i++)
{
try
{
return await operation();
}
catch (DbUpdateConcurrencyException ex)
{
lastException = ex;
// 指数退避
await Task.Delay(TimeSpan.FromMilliseconds(Math.Pow(2, i) * 100));
// 重新加载上下文
await ReloadDbContextAsync();
}
}
return Result.Fail(
$"Operation failed after {_maxRetries} retries: {lastException?.Message}");
}
private async Task ReloadDbContextAsync()
{
// 清除缓存,强制从数据库重新加载
_dbContext.ChangeTracker.Clear();
}
}5. 业务层面的冲突解决
5.1 合并策略
csharp
public class MergeStrategyService
{
/// <summary>
/// 合并并发修改(以计数器为例)
/// </summary>
public async Task<Result> IncrementCounterAsync(string key, int delta)
{
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// 使用数据库原子操作,避免并发冲突
var sql = @"
UPDATE counters
SET value = value + @Delta, version = version + 1
WHERE key = @Key";
var rowsAffected = await _dbContext.Database.ExecuteSqlRawAsync(
sql,
new SqlParameter("@Key", key),
new SqlParameter("@Delta", delta));
if (rowsAffected == 0)
{
// 记录不存在,插入新记录
await _dbContext.Counters.AddAsync(new Counter
{
Key = key,
Value = delta,
Version = 1
});
await _dbContext.SaveChangesAsync();
}
await transaction.CommitAsync();
return Result.Success();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
}5.2 最后写入胜出(LWW)
csharp
public class LastWriteWinsService
{
public async Task<Result> UpdateWithLWWAsync(
Guid entityId,
Action<Entity> updateAction)
{
try
{
var entity = await _dbContext.Entities.FindAsync(entityId);
// 应用更新
updateAction(entity);
// 强制覆盖(忽略版本号)
_dbContext.Entry(entity).Property(e => e.Version).IsModified = false;
await _dbContext.SaveChangesAsync();
return Result.Success();
}
catch (DbUpdateConcurrencyException)
{
// LWW 策略:忽略并发冲突,最后一次写入生效
return Result.Success("Updated (concurrent modification ignored)");
}
}
}5.3 业务规则判断
csharp
public class ApprovalService
{
/// <summary>
/// 审批流程(多人可能同时审批)
/// </summary>
public async Task<Result> ApproveAsync(Guid applicationId, string approverId)
{
try
{
var application = await _dbContext.Applications.FindAsync(applicationId);
// 业务规则:检查是否已被其他人审批
if (application.Status != ApplicationStatus.Pending)
{
return Result.Fail($"Application already {application.Status}");
}
application.Status = ApplicationStatus.Approved;
application.ApprovedBy = approverId;
application.ApprovedAt = DateTime.UtcNow;
application.Version++;
await _dbContext.SaveChangesAsync();
return Result.Success();
}
catch (DbUpdateConcurrencyException)
{
// 重新加载并检查状态
var application = await _dbContext.Applications.FindAsync(applicationId);
if (application.Status != ApplicationStatus.Pending)
{
return Result.Fail(
$"Application was already {application.Status} by another approver");
}
// 状态仍为 Pending,可能是其他并发问题,建议重试
return Result.Fail("Concurrent approval detected, please retry");
}
}
}6. 监控与告警
6.1 并发冲突指标
csharp
public class ConcurrencyMetrics
{
private readonly Counter<long> _conflictsTotal;
private readonly Counter<long> _conflictsByEntity;
private readonly Histogram<int> _retriesPerOperation;
private readonly Counter<long> _operationsFailedAfterRetries;
public ConcurrencyMetrics(IMeterFactory meterFactory)
{
var meter = meterFactory.Create("Idempotency.OptimisticLocking");
_conflictsTotal = meter.CreateCounter<long>(
"optimistic_lock.conflicts.total",
description: "Total number of concurrency conflicts");
_conflictsByEntity = meter.CreateCounter<long>(
"optimistic_lock.conflicts.by_entity",
description: "Concurrency conflicts by entity type");
_retriesPerOperation = meter.CreateHistogram<int>(
"optimistic_lock.retries.per_operation",
description: "Number of retries per operation");
_operationsFailedAfterRetries = meter.CreateCounter<long>(
"optimistic_lock.failures.after_retries",
description: "Operations that failed after all retries");
}
public void RecordConflict(string entityType)
{
_conflictsTotal.Add(1);
_conflictsByEntity.Add(1,
new KeyValuePair<string, object?>("entity_type", entityType));
}
public void RecordRetries(int retryCount)
{
_retriesPerOperation.Record(retryCount);
}
public void RecordFailure()
{
_operationsFailedAfterRetries.Add(1);
}
}6.2 Prometheus 告警规则
yaml
groups:
- name: optimistic_locking
rules:
# 高并发冲突率
- alert: HighConcurrencyConflictRate
expr: rate(optimistic_lock_conflicts_total[5m]) > 10
for: 5m
labels:
severity: warning
annotations:
summary: "High concurrency conflict rate"
description: "More than 10 conflicts per second in the last 5 minutes"
# 重试失败率高
- alert: HighRetryFailureRate
expr: rate(optimistic_lock_failures_after_retries_total[5m]) /
(rate(optimistic_lock_conflicts_total[5m]) + 1) > 0.5
for: 5m
labels:
severity: critical
annotations:
summary: "High retry failure rate"
description: "More than 50% of operations fail after retries"
# 平均重试次数过高
- alert: HighAverageRetries
expr: histogram_quantile(0.95,
rate(optimistic_lock_retries_per_operation_bucket[5m])) > 3
for: 10m
labels:
severity: warning
annotations:
summary: "High average retry count"
description: "95th percentile of retries is above 3"6.3 日志记录
csharp
public class ConcurrencyLoggingService
{
private readonly ILogger<ConcurrencyLoggingService> _logger;
public void LogConflict(ConcurrencyConflictInfo conflictInfo)
{
_logger.LogWarning(
"Concurrency conflict detected on {Table}. " +
"Original: {@OriginalValues}, " +
"Current: {@CurrentValues}, " +
"Database: {@DatabaseValues}",
conflictInfo.TableName,
conflictInfo.OriginalValues,
conflictInfo.CurrentValues,
conflictInfo.DatabaseValues);
}
public void LogRetry(int retryCount, string operation)
{
_logger.LogInformation(
"Retry {RetryCount} for operation {Operation}",
retryCount, operation);
}
public void LogExhaustedRetries(string operation, Exception ex)
{
_logger.LogError(ex,
"All retries exhausted for operation {Operation}",
operation);
}
}总结
乐观锁的并发冲突处理是保证数据一致性的关键:
核心要点
- 精确捕获:区分
DbUpdateConcurrencyException和其他异常 - 重试策略:使用 Polly 实现指数退避重试
- 业务判断:根据业务规则决定冲突解决策略
- 监控告警:及时发现高并发冲突场景
最佳实践
- 默认重试 3 次,指数退避
- 记录详细的冲突信息用于排查
- 监控冲突率,超过阈值时优化业务逻辑
- 对于高频冲突场景,考虑使用分布式锁
通过完善的并发冲突处理,可以让乐观锁方案在高并发环境下依然保持稳定可靠。