Skip to content

乐观并发控制原理 ​

EF Core 的并发冲突检测与处理机制

📖 目录 ​


什么是并发控制 ​

并发问题场景 ​

时间线:
T1: 用户A 读取产品 (Price = 100)
T2: 用户B 读取产品 (Price = 100)
T3: 用户A 修改价格为 120,保存
T4: 用户B 修改价格为 80,保存

结果: 用户A的修改被覆盖!(丢失更新)

为什么需要并发控制 ​

典型场景:

  1. 电商库存扣减

    • 用户A和用户B同时购买最后一件商品
    • 可能导致超卖
  2. 订单状态更新

    • 客服A标记为"已发货"
    • 客服B同时标记为"已取消"
    • 状态冲突
  3. 账户余额修改

    • 同时存取款
    • 余额计算错误

乐观 vs 悲观 ​

悲观并发控制 ​

csharp
// 锁定记录,其他事务等待
await using var transaction = await context.Database.BeginTransactionAsync(
    IsolationLevel.Serializable);

var product = await context.Products
    .Where(p => p.Id == id)
    .FirstOrDefaultAsync();

product.Price = newPrice;
await context.SaveChangesAsync();

await transaction.CommitAsync();

特点:

  • ✅ 保证数据一致性
  • ❌ 性能差(锁等待)
  • ❌ 可能死锁
  • ❌ 可扩展性差

适用场景:

  • 高冲突率
  • 金融系统
  • 关键业务

乐观并发控制(推荐) ​

csharp
// 不锁定,保存时检查是否被修改
var product = await context.Products.FindAsync(id);

product.Price = newPrice;

try
{
    await context.SaveChangesAsync();
}
catch (DbUpdateConcurrencyException)
{
    // 处理冲突
}

特点:

  • ✅ 性能好(无锁)
  • ✅ 可扩展性强
  • ✅ 适合大多数场景
  • ⚠️ 需要处理冲突

适用场景:

  • 低冲突率(大部分场景)
  • Web 应用
  • API 服务

实现方式 ​

方式 1: RowVersion / Timestamp(推荐) ​

实体配置 ​

csharp
public class Product
{
    public int Id { get; set; }
    public string Name { get; set; }
    public decimal Price { get; set; }
    
    // 行版本(并发令牌)
    [Timestamp]
    public byte[] RowVersion { get; set; }
}

或 Fluent API:

csharp
modelBuilder.Entity<Product>(entity =>
{
    entity.Property(e => e.RowVersion)
          .IsRowVersion();
});

数据库结构 ​

sql
CREATE TABLE Products (
    Id INT PRIMARY KEY IDENTITY(1,1),
    Name NVARCHAR(100),
    Price DECIMAL(18,2),
    RowVersion ROWVERSION  -- SQL Server
    -- 或 timestamp (PostgreSQL: xmin)
);

工作原理 ​

T1: 读取产品 (RowVersion = 0x00000000000007D1)
T2: 读取产品 (RowVersion = 0x00000000000007D1)
T3: A 修改价格,保存
    UPDATE Products SET Price = 120, RowVersion = NEW 
    WHERE Id = 1 AND RowVersion = 0x00000000000007D1
    ✅ 成功 (RowVersion 匹配)
    
T4: B 修改价格,保存
    UPDATE Products SET Price = 80, RowVersion = NEW
    WHERE Id = 1 AND RowVersion = 0x00000000000007D1
    ❌ 失败 (RowVersion 已改变,0 行受影响)
    → 抛出 DbUpdateConcurrencyException

方式 2: ConcurrencyCheck 特性 ​

csharp
public class Product
{
    public int Id { get; set; }
    
    [ConcurrencyCheck]
    public decimal Price { get; set; }
    
    [ConcurrencyCheck]
    public int StockQuantity { get; set; }
}

工作原理:

sql
UPDATE Products 
SET Price = @p0, StockQuantity = @p1
WHERE Id = @p2 
  AND Price = @original_Price      -- 检查原始值
  AND StockQuantity = @original_StockQuantity

缺点:

  • ❌ 只能检查特定字段
  • ❌ WHERE 子句变长
  • ✅ 适合少量字段

方式 3: 手动跟踪 ​

csharp
public class Product
{
    public int Id { get; set; }
    public decimal Price { get; set; }
    public DateTime LastModified { get; set; }
}

// 更新时检查
var product = await context.Products.FindAsync(id);
var originalLastModified = product.LastModified;

product.Price = newPrice;
product.LastModified = DateTime.UtcNow;

var affectedRows = await context.SaveChangesAsync();

if (affectedRows == 0)
{
    throw new ConcurrencyException("Record was modified by another user");
}

缺点:

  • ❌ 需要手动实现
  • ❌ 容易出错

并发冲突处理 ​

1. 客户端获胜(Client Wins) ​

csharp
try
{
    await context.SaveChangesAsync();
}
catch (DbUpdateConcurrencyException ex)
{
    foreach (var entry in ex.Entries)
    {
        // 使用客户端的值,覆盖数据库
        entry.OriginalValues.SetValues(entry.GetDatabaseValues());
    }
    
    // 重试保存
    await context.SaveChangesAsync();
}

适用场景:

  • 用户的修改更重要
  • 最后写入获胜

2. 数据库获胜(Store Wins) ​

csharp
try
{
    await context.SaveChangesAsync();
}
catch (DbUpdateConcurrencyException ex)
{
    foreach (var entry in ex.Entries)
    {
        // 刷新为数据库的最新值
        var databaseValues = entry.GetDatabaseValues();
        entry.CurrentValues.SetValues(databaseValues);
    }
    
    // 通知用户
    Console.WriteLine("Data was modified by another user. Refreshing...");
}

适用场景:

  • 数据完整性优先
  • 需要人工审查

3. 合并变更(Merge) ​

csharp
try
{
    await context.SaveChangesAsync();
}
catch (DbUpdateConcurrencyException ex)
{
    foreach (var entry in ex.Entries)
    {
        var proposedValues = entry.CurrentValues;
        var databaseValues = entry.GetDatabaseValues();
        var originalValues = entry.OriginalValues;
        
        // 自定义合并逻辑
        foreach (var property in proposedValues.Properties)
        {
            var proposedValue = proposedValues[property];
            var databaseValue = databaseValues[property];
            var originalValue = originalValues[property];
            
            // 如果客户端和数据库都修改了同一字段
            if (!Equals(proposedValue, originalValue) && 
                !Equals(databaseValue, originalValue))
            {
                // 冲突!需要特殊处理
                throw new ConcurrencyConflictException(
                    $"Conflict on {property.Name}: " +
                    $"Client={proposedValue}, DB={databaseValue}");
            }
            else if (!Equals(proposedValue, originalValue))
            {
                // 只有客户端修改,使用客户端值
                entry.Property(property.Name).CurrentValue = proposedValue;
            }
            // 否则使用数据库值(不变)
        }
    }
    
    await context.SaveChangesAsync();
}

适用场景:

  • 不同字段可以独立修改
  • 智能合并

4. 重试策略 ​

csharp
public async Task UpdateWithRetryAsync(int productId, decimal newPrice, int maxRetries = 3)
{
    for (int i = 0; i < maxRetries; i++)
    {
        try
        {
            var product = await context.Products.FindAsync(productId);
            product.Price = newPrice;
            await context.SaveChangesAsync();
            return; // 成功
        }
        catch (DbUpdateConcurrencyException)
        {
            if (i == maxRetries - 1)
                throw; // 最后一次仍失败,抛出异常
            
            // 等待后重试
            await Task.Delay(TimeSpan.FromMilliseconds(100 * (i + 1)));
        }
    }
}

适用场景:

  • 临时冲突
  • 高并发场景

API 中的并发控制 ​

Minimal API 示例 ​

csharp
app.Put("/api/products/{id}", async (int id, ProductDto dto, AppDbContext context) =>
{
    var product = await context.Products.FindAsync(id);
    if (product == null)
        return Results.NotFound();
    
    product.Name = dto.Name;
    product.Price = dto.Price;
    
    try
    {
        await context.SaveChangesAsync();
        return Results.NoContent();
    }
    catch (DbUpdateConcurrencyException)
    {
        if (!ProductExists(id))
            return Results.NotFound();
        else
            return Results.Conflict("Product was modified by another user");
    }
});

返回并发信息 ​

csharp
[HttpPut("{id}")]
public async Task<ActionResult> UpdateProduct(int id, ProductDto dto)
{
    var product = await context.Products.FindAsync(id);
    if (product == null)
        return NotFound();
    
    // 检查客户端的 RowVersion
    if (dto.RowVersion != null && 
        !dto.RowVersion.SequenceEqual(product.RowVersion))
    {
        return Conflict(new
        {
            message = "Product was modified by another user",
            currentData = new
            {
                product.Name,
                product.Price,
                product.RowVersion
            }
        });
    }
    
    product.Name = dto.Name;
    product.Price = dto.Price;
    
    try
    {
        await context.SaveChangesAsync();
        return NoContent();
    }
    catch (DbUpdateConcurrencyException)
    {
        return Conflict("Concurrency conflict");
    }
}

批量操作中的并发 ​

ExecuteUpdate 不支持并发检查 ​

csharp
// ⚠️ ExecuteUpdate 不会检查并发
await context.Products
    .Where(p => p.CategoryId == 1)
    .ExecuteUpdateAsync(setters => setters
        .SetProperty(p => p.Price, p => p.Price * 1.1m));

// 如果需要并发控制,使用传统方式
var products = await context.Products
    .Where(p => p.CategoryId == 1)
    .ToListAsync();

foreach (var product in products)
{
    product.Price *= 1.1m;
}

try
{
    await context.SaveChangesAsync();
}
catch (DbUpdateConcurrencyException ex)
{
    // 处理冲突
}

监控和日志 ​

记录并发冲突 ​

csharp
public class ConcurrencyLogger : SaveChangesInterceptor
{
    private readonly ILogger<ConcurrencyLogger> _logger;

    public override ValueTask<InterceptionResult<int>> SavingChangesAsync(
        DbContextEventData eventData,
        InterceptionResult<int> result,
        CancellationToken cancellationToken = default)
    {
        var context = eventData.Context;
        if (context != null)
        {
            var concurrentEntries = context.ChangeTracker.Entries()
                .Count(e => e.State == EntityState.Modified);
            
            _logger.LogInformation("Saving {Count} modified entities", concurrentEntries);
        }
        
        return base.SavingChangesAsync(eventData, result, cancellationToken);
    }

    public override void ThrowConcurrentOperationException(
        ConcurrentUpdateExceptionEventData eventData)
    {
        _logger.LogWarning(eventData.Exception, 
            "Concurrency conflict detected for {EntityType}", 
            eventData.Entry.Entity.GetType().Name);
        
        base.ThrowConcurrentOperationException(eventData);
    }
}

// 注册
builder.Services.AddDbContext<AppDbContext>(options =>
{
    options.UseSqlServer(connectionString);
    options.AddInterceptors(new ConcurrencyLogger(logger));
});

最佳实践 ​

1. 始终使用 RowVersion ​

csharp
// ✅ 推荐
[Timestamp]
public byte[] RowVersion { get; set; }

// ❌ 避免
[ConcurrencyCheck]
public decimal Price { get; set; }

2. API 返回 RowVersion ​

csharp
public class ProductDto
{
    public int Id { get; set; }
    public string Name { get; set; }
    public decimal Price { get; set; }
    public byte[] RowVersion { get; set; } // 返回给客户端
}

// 客户端下次请求时带回 RowVersion

3. 友好的错误提示 ​

csharp
catch (DbUpdateConcurrencyException)
{
    return BadRequest(new
    {
        error = "concurrency_conflict",
        message = "This record has been modified by another user. Please refresh and try again.",
        action = "refresh_and_retry"
    });
}

4. 前端处理 ​

javascript
// JavaScript 示例
async function updateProduct(id, data) {
    try {
        const response = await fetch(`/api/products/${id}`, {
            method: 'PUT',
            headers: { 'Content-Type': 'application/json' },
            body: JSON.stringify(data)
        });
        
        if (response.status === 409) {
            // 并发冲突
            const error = await response.json();
            alert(error.message);
            // 重新加载最新数据
            await loadProduct(id);
        }
    } catch (error) {
        console.error('Update failed:', error);
    }
}

5. 测试并发 ​

csharp
[Fact]
public async Task ConcurrentUpdates_ShouldDetectConflict()
{
    var product = new Product { Name = "Test", Price = 100 };
    context.Products.Add(product);
    await context.SaveChangesAsync();
    
    // 模拟两个上下文同时修改
    var context1 = CreateNewContext();
    var context2 = CreateNewContext();
    
    var product1 = await context1.Products.FindAsync(product.Id);
    var product2 = await context2.Products.FindAsync(product.Id);
    
    product1.Price = 120;
    await context1.SaveChangesAsync(); // 成功
    
    product2.Price = 80;
    await Assert.ThrowsAsync<DbUpdateConcurrencyException>(
        () => context2.SaveChangesAsync()); // 抛出异常
}

📊 性能对比 ​

方式性能复杂度适用场景
无并发控制⭐⭐⭐⭐⭐低单用户系统
乐观并发(RowVersion)⭐⭐⭐⭐中Web 应用(推荐)
乐观并发(ConcurrencyCheck)⭐⭐⭐中少量字段检查
悲观并发(锁)⭐⭐高高冲突/金融系统

💡 小结 ​

核心要点:

  • ✅ 乐观并发适合大多数 Web 应用
  • ✅ 使用 RowVersion/Timestamp 实现
  • ✅ 捕获 DbUpdateConcurrencyException 处理冲突
  • ✅ API 返回 RowVersion 给客户端
  • ✅ 提供友好的错误提示
  • ❌ ExecuteUpdate 不支持并发检查

实现步骤:

1. 添加 RowVersion 属性
2. 配置为 [Timestamp] 或 IsRowVersion()
3. 捕获 DbUpdateConcurrencyException
4. 选择冲突解决策略
5. 返回友好错误给客户端

下一步:

  1. 学习 处理并发异常
  2. 掌握 悲观并发控制
  3. 理解 事务隔离级别

基于 MIT 许可发布