使用Quartz .NET和EF Core实现水平扩展时的并发数据更新问题
解决方案:Quartz集群下EF Core避免竞态与负载拆分
一、核心问题分析
当前代码的问题在于:多个Quartz实例同时执行作业时,都会查询Created < date && IsHidden == false的日志数据,但查询与更新操作之间没有加锁机制,导致不同实例会拿到同一批未处理数据,最终出现重复更新的竞态问题。
二、数据库行锁实现(EF Core层面)
解决竞态的核心是在查询阶段就锁定目标数据行,防止其他实例读取;同时让实例自动跳过已被锁定的行,实现负载自动拆分。以下针对不同数据库给出实现方案:
1. SQL Server 版本修改代码
public async Task<int> UpdateNameAndMarkHiddenLogsAfterReachingDate(DateTime date) { const int batchSize = 10; int hidden = 0; int hiddenInThisLoopCycle; do { using (var transaction = await _context.Database.BeginTransactionAsync()) { try { // 使用UPDLOCK加更新锁,READPAST跳过已被锁定的行,避免实例间等待 var logs = await _context.Logs .FromSqlRaw(@"SELECT TOP {0} * FROM Logs WHERE Created < @date AND IsHidden = 0 WITH (UPDLOCK, READPAST)", batchSize) .AddParameter(new SqlParameter("@date", date)) .ToListAsync(); if (!logs.Any()) { hiddenInThisLoopCycle = 0; await transaction.CommitAsync(); break; } foreach (var log in logs) { log.Message = "UpdatedLog"; log.IsHidden = true; log.LogDetails.Add(new LogDetail() { ChangedByThread = Environment.CurrentManagedThreadId.ToString(), LastChanged = DateTime.UtcNow, Name = Thread.CurrentThread.Name }); } hiddenInThisLoopCycle = await _context.SaveChangesAsync(); hidden += hiddenInThisLoopCycle; Console.WriteLine($"处理 {hiddenInThisLoopCycle} 条日志,累计处理 {hidden} 条"); await transaction.CommitAsync(); } catch (Exception ex) { await transaction.RollbackAsync(); Console.WriteLine($"事务失败: {ex.Message}"); throw; } } } while (hiddenInThisLoopCycle > 0); return hidden; }
2. PostgreSQL 版本修改代码
public async Task<int> UpdateNameAndMarkHiddenLogsAfterReachingDate(DateTime date) { const int batchSize = 10; int hidden = 0; int hiddenInThisLoopCycle; do { using (var transaction = await _context.Database.BeginTransactionAsync()) { try { // PostgreSQL使用FOR UPDATE SKIP LOCKED实现锁定+跳过已锁行 var logs = await _context.Logs .FromSqlRaw(@"SELECT * FROM Logs WHERE Created < @date AND IsHidden = false LIMIT {0} FOR UPDATE SKIP LOCKED", batchSize) .AddParameter(new NpgsqlParameter("@date", date)) .ToListAsync(); if (!logs.Any()) { hiddenInThisLoopCycle = 0; await transaction.CommitAsync(); break; } foreach (var log in logs) { log.Message = "UpdatedLog"; log.IsHidden = true; log.LogDetails.Add(new LogDetail() { ChangedByThread = Environment.CurrentManagedThreadId.ToString(), LastChanged = DateTime.UtcNow, Name = Thread.CurrentThread.Name }); } hiddenInThisLoopCycle = await _context.SaveChangesAsync(); hidden += hiddenInThisLoopCycle; Console.WriteLine($"处理 {hiddenInThisLoopCycle} 条日志,累计处理 {hidden} 条"); await transaction.CommitAsync(); } catch (Exception ex) { await transaction.RollbackAsync(); Console.WriteLine($"事务失败: {ex.Message}"); throw; } } } while (hiddenInThisLoopCycle > 0); return hidden; }
锁策略说明
UPDLOCK(SQL Server)/FOR UPDATE(PostgreSQL):对读取的行加更新锁,其他事务仅能读取但无法修改或加同类锁,直到当前事务提交。READPAST(SQL Server)/SKIP LOCKED(PostgreSQL):自动跳过已被其他事务锁定的行,让不同实例直接获取未处理的批次数据,实现无冲突的负载拆分。
三、Quartz集群优化配置
除了数据库锁,Quartz集群本身的配置也需要调整,确保作业合理分配:
- 所有Quartz实例必须使用同一个数据库作为JobStore存储(配置ADO.NET JobStore),这是集群协同的基础。
- 保留
[DisallowConcurrentExecution]特性,确保单个实例不会并发执行同一个作业;同时给JobDetail设置RequestsRecovery = true,防止实例宕机后作业丢失。 - 调整
org.quartz.jobStore.clusterCheckinInterval参数(默认15000ms),根据实例数量适当缩短,让集群更快感知实例状态变化。 - 若作业触发时间固定,可给不同实例配置带随机偏移的Cron表达式,减少多实例同时触发作业的概率。
四、额外优化建议
- 移除原代码中的外层大事务,改为每批次处理一个小事务,避免长时间占用锁资源,提升整体并发性能。
- 若日志数据量极大,可按
Created字段分片,比如让不同实例处理不同时间范围的日志(如实例1处理7天前的数据,实例2处理14天前的数据),进一步降低锁竞争。 - 根据数据库性能调整
batchSize大小,平衡单批次处理效率与锁竞争频率。
内容的提问来源于stack exchange,提问作者Ihor Arkhypenko
相关产品推荐
相关产品推荐

