Skip to content

乐观锁 - 并发冲突处理 ​

目录 ​


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);
    }
}

总结 ​

乐观锁的并发冲突处理是保证数据一致性的关键:

核心要点 ​

  1. 精确捕获:区分 DbUpdateConcurrencyException 和其他异常
  2. 重试策略:使用 Polly 实现指数退避重试
  3. 业务判断:根据业务规则决定冲突解决策略
  4. 监控告警:及时发现高并发冲突场景

最佳实践 ​

  • 默认重试 3 次,指数退避
  • 记录详细的冲突信息用于排查
  • 监控冲突率,超过阈值时优化业务逻辑
  • 对于高频冲突场景,考虑使用分布式锁

通过完善的并发冲突处理,可以让乐观锁方案在高并发环境下依然保持稳定可靠。

Released under the MIT License.