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

.NET Core Kafka消费者插入SQL Server遇deadlock死锁问题求助

.NET Core微服务Kafka消费者高流量下SQL Server死锁问题

我是.NET新手,用.NET Core搭建了一个微服务作为Kafka消费者,消费函数中包含插入SQL Server数据库的逻辑。流量较低时服务可正常运行,但流量较高时频繁出现死锁错误。我已尝试使用with no lock,现在正在寻找死锁时自动重试消费者的方法,恳请有相关经验的人士提供帮助。

错误信息

数据库保存上下文类型为'Entities.RepositoryContext'的更改时发生异常。
      Microsoft.EntityFrameworkCore.DbUpdateException: 更新条目时出错。请查看内部异常了解详细信息。
       ---> Microsoft.Data.SqlClient.SqlException (0x80131904): 事务(进程 ID 414)与另一个进程在锁资源上发生死锁,已被选为死锁牺牲品。请重新运行该事务。

insertData函数代码

private async Task insertData(Dto data)
{
    var dataFromDB = _repository.myRepo.FirstOrDefaultWithNoLock(p => p.ID.Equals(data.ID));

    if (dataFromDB == null)
    {
        dataFromDB = _mapper.Map<Model>(data);
        _repository.myRepo.Create(dataFromDB);
        _repository.myRepo.SaveChangesWithIdentityInsert();
        return;
    }

    return;
}

SaveChangesWithIdentityInsert实现代码

public static void EnableIdentityInsert<T>(this DbContext context) => SetIdentityInsert<T>(context, true);
public static void DisableIdentityInsert<T>(this DbContext context) => SetIdentityInsert<T>(context, false);

private static void SetIdentityInsert<T>([NotNull] DbContext context, bool enable)
{
    if (context == null) throw new ArgumentNullException(nameof(context));
    var entityType = context.Model.FindEntityType(typeof(T));
    var value = enable ? "ON" : "OFF";
    context.Database.ExecuteSqlRaw($"SET IDENTITY_INSERT dbo.{entityType.GetTableName()} {value}");
}

public static void SaveChangesWithIdentityInsert<T>([NotNull] this DbContext context)
{
    if (context == null) throw new ArgumentNullException(nameof(context));
    var strategy = context.Database.CreateExecutionStrategy();
    strategy.Execute(
    () =>
    {
        using var transaction = context.Database.BeginTransaction(isolationLevel: System.Data.IsolationLevel.ReadUncommitted);
        context.EnableIdentityInsert<T>();
        context.SaveChanges();
        context.DisableIdentityInsert<T>();
        transaction.Commit();
    });
    context.ChangeTracker.Clear();
}

解决方案建议

1. 调整事务与隔离级别

  • 当前使用的ReadUncommitted隔离级别可能读取脏数据,加剧并发插入冲突。建议开启数据库的READ_COMMITTED_SNAPSHOT模式(执行SQL:ALTER DATABASE YourDatabase SET READ_COMMITTED_SNAPSHOT ON;),该级别可避免脏读,同时降低锁竞争。
  • 将查询和插入操作合并到同一个事务内,避免“查询-插入”的间隙导致多请求同时插入同一ID。修改后的代码示例:
private async Task insertData(Dto data)
{
    var strategy = _repository.myRepo.Database.CreateExecutionStrategy();
    await strategy.ExecuteAsync(async () =>
    {
        using var transaction = await _repository.myRepo.Database.BeginTransactionAsync(System.Data.IsolationLevel.ReadCommitted);
        try
        {
            var dataFromDB = await _repository.myRepo.FirstOrDefaultAsync(p => p.ID.Equals(data.ID));
            if (dataFromDB == null)
            {
                dataFromDB = _mapper.Map<Model>(data);
                _repository.myRepo.Create(dataFromDB);
                _repository.myRepo.EnableIdentityInsert<Model>();
                await _repository.myRepo.SaveChangesAsync();
                _repository.myRepo.DisableIdentityInsert<Model>();
            }
            await transaction.CommitAsync();
        }
        catch
        {
            await transaction.RollbackAsync();
            throw;
        }
    });
}

2. 完善EF Core重试策略

EF Core的默认执行策略会自动重试死锁等可恢复异常,确保事务重新执行。建议使用异步版本ExecuteAsync适配Kafka消费的异步场景,提升并发处理能力。

3. 数据库层面优化

  • 确保ID字段为主键并建立聚集索引,精准的索引范围能缩小锁的持有范围,减少冲突概率。
  • 缩短事务执行时间,只在事务内执行必要的插入操作,避免无关操作延长锁持有时间。

4. Kafka消费者重试配置

  • 在Kafka消费者配置中设置RetryBackoffMs,增加重试间隔,避免短时间内大量重试加重数据库负载。
  • 配置死信队列(DLQ),将多次消费失败的消息转移到DLQ,避免阻塞正常消息消费,后续再手动处理DLQ中的消息。

5. 优化IDENTITY_INSERT使用

仅在需要手动指定ID值时开启IDENTITY_INSERT,如果大部分场景下ID由数据库自动生成,避免每次插入都执行该操作,减少额外数据库开销。

内容的提问来源于stack exchange,提问作者Firman Alamsyah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 10:39:52