乐观锁 - 版本号机制
概述
乐观锁(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
- 始终使用 ConcurrencyCheck:让 EF Core 自动处理版本号
- 实现自动重试:版本冲突时自动重试 2-3 次
- 设置合理的重试间隔:使用指数退避
- 监控冲突率:如果超过 5%,考虑优化方案
- 记录详细日志:便于排查问题
❌ DON'T
- 不要忽略 DbUpdateConcurrencyException:必须处理
- 不要无限重试:设置最大重试次数
- 不要在长事务中使用:增加冲突概率
- 不要对只读数据使用:浪费资源
对比:乐观锁 vs 悲观锁
| 特性 | 乐观锁 | 悲观锁 |
|---|---|---|
| 适用场景 | 读多写少 | 写多读少 |
| 性能 | 高(无锁开销) | 低(锁竞争) |
| 并发度 | 高 | 低 |
| 实现复杂度 | 中 | 低 |
| 死锁风险 | 无 | 有 |
| 重试需求 | 需要 | 不需要 |
总结
乐观锁是处理并发更新的优雅方案:
✅ 优点:
- 无锁开销,性能高
- 不会死锁
- 适合读多写少场景
⚠️ 注意事项:
- 需要实现重试机制
- 高并发写场景可能频繁冲突
- 监控冲突率,及时调整策略
在实际应用中,乐观锁常与唯一索引、分布式锁结合使用,形成多层并发控制体系。