Skip to content

乐观锁 - 版本号机制 ​

概述 ​

乐观锁(Optimistic Locking)是一种基于版本号的并发控制机制。它假设冲突很少发生,因此在读取数据时不加锁,只在更新时检查版本号是否发生变化。如果版本号不一致,说明数据已被其他事务修改,更新失败并需要重试。

核心原理 ​

事务 A                          数据库                       事务 B
  |                               |                             |
  |-- 读取数据 (version=1) ------>|                             |
  |                               |                             |-- 读取数据 (version=1)
  |                               |                             |
  |-- 修改数据                     |                             |-- 修改数据
  |-- 更新 (version=2) ---------->|                             |
  |                               |-- 检查 version=1 ✓          |
  |                               |-- 更新成功                   |
  |                               |                             |
  |                               |                             |-- 更新 (version=2)
  |                               |                             |-- 检查 version=1 ✗
  |                               |                             |-- 更新失败(版本冲突)
  |                               |                             |-- 重试:重新读取

PostgreSQL 实现 ​

1. 表设计 ​

sql
CREATE TABLE products (
    id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    name VARCHAR(100) NOT NULL,
    price DECIMAL(10, 2) NOT NULL,
    stock INTEGER NOT NULL DEFAULT 0,
    
    -- 乐观锁版本号
    version INTEGER NOT NULL DEFAULT 0,
    
    created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
    updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);

-- 索引
CREATE INDEX idx_products_name ON products(name);

2. Entity Framework Core 配置 ​

csharp
using Microsoft.EntityFrameworkCore;

public class Product
{
    public Guid Id { get; set; }
    public string Name { get; set; }
    public decimal Price { get; set; }
    public int Stock { get; set; }
    
    // 并发令牌
    [ConcurrencyCheck]
    public int Version { get; set; }
    
    public DateTime CreatedAt { get; set; }
    public DateTime UpdatedAt { get; set; }
}

public class AppDbContext : DbContext
{
    public DbSet<Product> Products { get; set; }
    
    protected override void OnModelCreating(ModelBuilder modelBuilder)
    {
        modelBuilder.Entity<Product>()
            .Property(p => p.Version)
            .IsConcurrencyToken(); // 标记为并发令牌
        
        modelBuilder.Entity<Product>()
            .HasQueryFilter(p => !p.IsDeleted); // 全局查询过滤器
    }
}

3. 库存扣减示例 ​

csharp
public class InventoryService
{
    private readonly AppDbContext _dbContext;
    private readonly ILogger<InventoryService> _logger;
    
    public async Task<Result<bool>> DeductStockAsync(
        Guid productId, 
        int quantity)
    {
        await using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        try
        {
            // 1. 读取产品和版本号
            var product = await _dbContext.Products
                .FirstOrDefaultAsync(p => p.Id == productId);
            
            if (product == null)
            {
                return Result<bool>.Failure("Product not found");
            }
            
            if (product.Stock < quantity)
            {
                return Result<bool>.Failure("Insufficient stock");
            }
            
            // 2. 记录当前版本号
            var originalVersion = product.Version;
            
            // 3. 修改数据
            product.Stock -= quantity;
            product.Version++; // 递增版本号
            product.UpdatedAt = DateTime.UtcNow;
            
            // 4. 保存(EF Core 会自动在 WHERE 子句中添加版本号检查)
            var affectedRows = await _dbContext.SaveChangesAsync();
            
            if (affectedRows == 0)
            {
                // 版本号冲突,说明数据已被其他事务修改
                _logger.LogWarning("Optimistic lock conflict for product {ProductId}", productId);
                return Result<bool>.Failure("Concurrent modification detected, please retry");
            }
            
            await transaction.CommitAsync();
            
            _logger.LogInformation("Stock deducted: Product {ProductId}, Quantity {Quantity}, New Stock {Stock}",
                productId, quantity, product.Stock);
            
            return Result<bool>.Success(true);
        }
        catch (DbUpdateConcurrencyException ex)
        {
            await transaction.RollbackAsync();
            
            _logger.LogWarning(ex, "Concurrency conflict when deducting stock for product {ProductId}", productId);
            
            return Result<bool>.Failure("Concurrent modification detected");
        }
        catch (Exception ex)
        {
            await transaction.RollbackAsync();
            
            _logger.LogError(ex, "Failed to deduct stock for product {ProductId}", productId);
            
            return Result<bool>.Failure("Internal error");
        }
    }
}

4. 手动 SQL 实现(更灵活) ​

csharp
public class ManualOptimisticLockService
{
    private readonly AppDbContext _dbContext;
    
    /// <summary>
    /// 使用原始 SQL 实现乐观锁
    /// </summary>
    public async Task<Result<bool>> UpdateProductPriceAsync(
        Guid productId,
        decimal newPrice,
        int expectedVersion)
    {
        const string sql = @"
            UPDATE products
            SET price = @newPrice,
                version = version + 1,
                updated_at = NOW()
            WHERE id = @productId
              AND version = @expectedVersion";
        
        var rowsAffected = await _dbContext.Database.ExecuteSqlRawAsync(
            sql,
            new NpgsqlParameter("@newPrice", newPrice),
            new NpgsqlParameter("@productId", productId),
            new NpgsqlParameter("@expectedVersion", expectedVersion));
        
        if (rowsAffected == 0)
        {
            return Result<bool>.Failure("Version mismatch, update failed");
        }
        
        return Result<bool>.Success(true);
    }
}

带重试的乐观锁 ​

1. 自动重试机制 ​

csharp
public class RetryableOptimisticLockService
{
    private readonly AppDbContext _dbContext;
    private readonly ILogger<RetryableOptimisticLockService> _logger;
    
    /// <summary>
    /// 带自动重试的库存扣减
    /// </summary>
    public async Task<Result<bool>> DeductStockWithRetryAsync(
        Guid productId,
        int quantity,
        int maxRetries = 3)
    {
        for (int attempt = 1; attempt <= maxRetries; attempt++)
        {
            try
            {
                var result = await DeductStockAsync(productId, quantity);
                
                if (result.IsSuccess)
                {
                    return result;
                }
                
                // 如果是并发冲突,重试
                if (result.Error.Contains("Concurrent modification") && attempt < maxRetries)
                {
                    _logger.LogInformation("Retry attempt {Attempt}/{Max} for product {ProductId}",
                        attempt, maxRetries, productId);
                    
                    // 短暂等待后重试(避免立即冲突)
                    await Task.Delay(TimeSpan.FromMilliseconds(50 * attempt));
                    continue;
                }
                
                // 其他错误,直接返回
                return result;
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Unexpected error on attempt {Attempt}", attempt);
                
                if (attempt == maxRetries)
                {
                    throw;
                }
            }
        }
        
        return Result<bool>.Failure($"Failed after {maxRetries} retries");
    }
    
    private async Task<Result<bool>> DeductStockAsync(Guid productId, int quantity)
    {
        await using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        try
        {
            var product = await _dbContext.Products
                .FirstOrDefaultAsync(p => p.Id == productId);
            
            if (product == null)
            {
                return Result<bool>.Failure("Product not found");
            }
            
            if (product.Stock < quantity)
            {
                return Result<bool>.Failure("Insufficient stock");
            }
            
            product.Stock -= quantity;
            product.Version++;
            product.UpdatedAt = DateTime.UtcNow;
            
            await _dbContext.SaveChangesAsync();
            await transaction.CommitAsync();
            
            return Result<bool>.Success(true);
        }
        catch (DbUpdateConcurrencyException)
        {
            await transaction.RollbackAsync();
            return Result<bool>.Failure("Concurrent modification detected");
        }
    }
}

2. 使用 Polly 实现重试 ​

csharp
using Polly;

public class PollyOptimisticLockService
{
    private readonly AppDbContext _dbContext;
    private readonly IAsyncPolicy _retryPolicy;
    
    public PollyOptimisticLockService(AppDbContext dbContext)
    {
        _dbContext = dbContext;
        
        // 配置重试策略:最多3次,指数退避
        _retryPolicy = Policy
            .Handle<DbUpdateConcurrencyException>()
            .WaitAndRetryAsync(
                retryCount: 3,
                sleepDurationProvider: (attempt, context) => 
                    TimeSpan.FromMilliseconds(100 * Math.Pow(2, attempt)),
                onRetry: (outcome, timespan, retryNumber, context) =>
                {
                    Console.WriteLine($"Retry {retryNumber} after {timespan.TotalMilliseconds}ms");
                });
    }
    
    public async Task DeductStockAsync(Guid productId, int quantity)
    {
        await _retryPolicy.ExecuteAsync(async () =>
        {
            var product = await _dbContext.Products.FindAsync(productId);
            
            if (product.Stock < quantity)
            {
                throw new InvalidOperationException("Insufficient stock");
            }
            
            product.Stock -= quantity;
            product.Version++;
            
            await _dbContext.SaveChangesAsync();
        });
    }
}

状态机流转控制 ​

1. 订单状态机 ​

csharp
public enum OrderStatus
{
    Created = 1,
    Paid = 2,
    Shipped = 3,
    Delivered = 4,
    Cancelled = 5,
    Refunded = 6
}

public class Order
{
    public Guid Id { get; set; }
    public Guid UserId { get; set; }
    public decimal TotalAmount { get; set; }
    
    public OrderStatus Status { get; set; }
    
    [ConcurrencyCheck]
    public int Version { get; set; }
    
    public DateTime CreatedAt { get; set; }
    public DateTime? PaidAt { get; set; }
    public DateTime? ShippedAt { get; set; }
}

public class OrderStateMachine
{
    private readonly AppDbContext _dbContext;
    
    /// <summary>
    /// 定义合法的状态转换
    /// </summary>
    private static readonly Dictionary<OrderStatus, HashSet<OrderStatus>> ValidTransitions = 
        new()
        {
            { OrderStatus.Created, new() { OrderStatus.Paid, OrderStatus.Cancelled } },
            { OrderStatus.Paid, new() { OrderStatus.Shipped, OrderStatus.Refunded } },
            { OrderStatus.Shipped, new() { OrderStatus.Delivered } },
            { OrderStatus.Delivered, new() { OrderStatus.Refunded } },
            { OrderStatus.Cancelled, new() { } }, // 终态
            { OrderStatus.Refunded, new() { } }   // 终态
        };
    
    public async Task<Result<bool>> TransitionStatusAsync(
        Guid orderId,
        OrderStatus newStatus,
        int expectedVersion)
    {
        await using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        try
        {
            var order = await _dbContext.Orders.FindAsync(orderId);
            
            if (order == null)
            {
                return Result<bool>.Failure("Order not found");
            }
            
            // 1. 验证状态转换是否合法
            if (!IsValidTransition(order.Status, newStatus))
            {
                return Result<bool>.Failure(
                    $"Invalid transition from {order.Status} to {newStatus}");
            }
            
            // 2. 记录旧状态
            var oldStatus = order.Status;
            
            // 3. 更新状态和版本号
            order.Status = newStatus;
            order.Version++;
            
            // 4. 更新相关时间戳
            UpdateTimestamps(order, newStatus);
            
            // 5. 保存(带版本号检查)
            await _dbContext.SaveChangesAsync();
            await transaction.CommitAsync();
            
            return Result<bool>.Success(true);
        }
        catch (DbUpdateConcurrencyException)
        {
            await transaction.RollbackAsync();
            return Result<bool>.Failure("Concurrent modification detected");
        }
    }
    
    private bool IsValidTransition(OrderStatus from, OrderStatus to)
    {
        return ValidTransitions.TryGetValue(from, out var validTargets) && 
               validTargets.Contains(to);
    }
    
    private void UpdateTimestamps(Order order, OrderStatus newStatus)
    {
        switch (newStatus)
        {
            case OrderStatus.Paid:
                order.PaidAt = DateTime.UtcNow;
                break;
            case OrderStatus.Shipped:
                order.ShippedAt = DateTime.UtcNow;
                break;
        }
    }
}

// 使用示例
public class OrderService
{
    private readonly OrderStateMachine _stateMachine;
    
    public async Task PayOrderAsync(Guid orderId)
    {
        var order = await _dbContext.Orders.FindAsync(orderId);
        
        // 尝试转换状态:Created -> Paid
        var result = await _stateMachine.TransitionStatusAsync(
            orderId, 
            OrderStatus.Paid,
            order.Version);
        
        if (!result.IsSuccess)
        {
            throw new InvalidOperationException(result.Error);
        }
    }
}

性能优化 ​

1. 批量更新 ​

csharp
public class BatchOptimisticLockService
{
    private readonly AppDbContext _dbContext;
    
    /// <summary>
    /// 批量扣减库存(带乐观锁)
    /// </summary>
    public async Task<Result<int>> BatchDeductStockAsync(
        List<StockDeductionRequest> requests)
    {
        await using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        try
        {
            int successCount = 0;
            
            foreach (var request in requests)
            {
                var sql = @"
                    UPDATE products
                    SET stock = stock - @quantity,
                        version = version + 1,
                        updated_at = NOW()
                    WHERE id = @productId
                      AND version = @version
                      AND stock >= @quantity";
                
                var rowsAffected = await _dbContext.Database.ExecuteSqlRawAsync(
                    sql,
                    new NpgsqlParameter("@quantity", request.Quantity),
                    new NpgsqlParameter("@productId", request.ProductId),
                    new NpgsqlParameter("@version", request.Version));
                
                if (rowsAffected > 0)
                {
                    successCount++;
                }
            }
            
            await transaction.CommitAsync();
            
            return Result<int>.Success(successCount);
        }
        catch (Exception ex)
        {
            await transaction.RollbackAsync();
            return Result<int>.Failure($"Batch update failed: {ex.Message}");
        }
    }
}

public class StockDeductionRequest
{
    public Guid ProductId { get; set; }
    public int Quantity { get; set; }
    public int Version { get; set; }
}

2. 减少版本号冲突 ​

csharp
public class OptimizedInventoryService
{
    /// <summary>
    /// 使用分区策略减少热点商品的版本冲突
    /// </summary>
    public async Task<Result<bool>> DeductStockPartitionedAsync(
        Guid productId,
        int quantity)
    {
        // 为热门商品使用多个库存分片
        var partitionKey = GetPartitionKey(productId);
        
        var sql = @"
            UPDATE inventory_partitions
            SET stock = stock - @quantity,
                version = version + 1
            WHERE product_id = @productId
              AND partition_key = @partitionKey
              AND version = @version
              AND stock >= @quantity";
        
        // ... 执行更新
        
        return Result<bool>.Success(true);
    }
    
    private int GetPartitionKey(Guid productId)
    {
        // 根据产品 ID 和用户 ID 计算分区键
        // 使同一产品的不同用户请求分散到不同分区
        return Math.Abs(productId.GetHashCode()) % 4; // 4 个分区
    }
}

监控与诊断 ​

1. 统计版本冲突率 ​

csharp
public class OptimisticLockMetrics
{
    private readonly Counter<long> _updateAttempts;
    private readonly Counter<long> _versionConflicts;
    private readonly Histogram<double> _retryCount;
    
    public void RecordUpdateAttempt(bool success, int retryCount)
    {
        _updateAttempts.Add(1);
        
        if (!success)
        {
            _versionConflicts.Add(1);
        }
        
        _retryCount.Record(retryCount);
    }
    
    public double GetConflictRate()
    {
        // 从指标系统获取
        var attempts = _updateAttempts.GetCount();
        var conflicts = _versionConflicts.GetCount();
        
        return attempts > 0 ? (double)conflicts / attempts : 0;
    }
}

2. PostgreSQL 监控查询 ​

sql
-- 查看频繁发生版本冲突的记录
SELECT 
    id,
    name,
    version,
    updated_at,
    age(NOW(), updated_at) as last_updated_ago
FROM products
WHERE version > 100  -- 版本号很高的记录
ORDER BY version DESC
LIMIT 20;

-- 统计每个表的并发更新次数
SELECT 
    schemaname,
    relname as table_name,
    n_tup_upd as total_updates,
    n_dead_tup as dead_tuples
FROM pg_stat_user_tables
ORDER BY n_tup_upd DESC;

最佳实践总结 ​

✅ DO ​

  1. 始终使用 ConcurrencyCheck:让 EF Core 自动处理版本号
  2. 实现自动重试:版本冲突时自动重试 2-3 次
  3. 设置合理的重试间隔:使用指数退避
  4. 监控冲突率:如果超过 5%,考虑优化方案
  5. 记录详细日志:便于排查问题

❌ DON'T ​

  1. 不要忽略 DbUpdateConcurrencyException:必须处理
  2. 不要无限重试:设置最大重试次数
  3. 不要在长事务中使用:增加冲突概率
  4. 不要对只读数据使用:浪费资源

对比:乐观锁 vs 悲观锁 ​

特性乐观锁悲观锁
适用场景读多写少写多读少
性能高(无锁开销)低(锁竞争)
并发度高低
实现复杂度中低
死锁风险无有
重试需求需要不需要

总结 ​

乐观锁是处理并发更新的优雅方案:

✅ 优点:

  • 无锁开销,性能高
  • 不会死锁
  • 适合读多写少场景

⚠️ 注意事项:

  • 需要实现重试机制
  • 高并发写场景可能频繁冲突
  • 监控冲突率,及时调整策略

在实际应用中,乐观锁常与唯一索引、分布式锁结合使用,形成多层并发控制体系。

Released under the MIT License.