You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 16:32:45