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

C# BlockingCollection实现异步日志队列批量存库问题求解

.NET 异步日志批量落库实现方案

问题背景

在应用程序中,大量操作日志需要写入数据库,但日志写入流程不能拖慢正常请求的响应速度,因此需要通过异步队列类机制实现日志的异步落库。
整体设计架构如下图所示:
整体架构
现有实现代码如下:

日志模型

public class ActivityLog
{
    [DatabaseGenerated(DatabaseGeneratedOption.Identity)]
    public string Id { get; set; }
    public IPAddress IPAddress { get; set; }
    public string Action { get; set; }
    public string? Metadata { get; set; }
    public string? UserId { get; set; }
    public DateTime CreationTime { get; set; }
}

队列实现

public class LogQueue
{
    private const int QueueCapacity = 1_000_000; // 疑问:100万的队列容量是否足够?
    private readonly BlockingCollection<ActivityLog> logs = new(QueueCapacity);

    public bool IsCompleted => logs.IsCompleted;
    public void Add(ActivityLog log) => logs.Add(log);
    public IEnumerable<ActivityLog> GetConsumingEnumerable() => logs.GetConsumingEnumerable();
    public void Complete() => logs.CompleteAdding();
}

后台工作者(存在问题的部分)

public class DbLogWorker : IHostedService
{
    private readonly LogQueue queue;
    private readonly IServiceScopeFactory scf;
    private Task jobTask;

    public DbLogWorker(LogQueue queue, IServiceScopeFactory scf)
    {
        this.queue = queue;
        this.scf = scf;
        jobTask = new Task(Job, TaskCreationOptions.LongRunning);
    }
  
    private void Job()
    {
        using var scope = scf.CreateScope();
        var dbContext = scope.ServiceProvider.GetRequiredService<ApplicationDbContext>();

        // 以下代码无法正常运行,设计初衷是减少数据库往返次数
        //while(!queue.IsCompleted)
        //{
        //    var items = queue.GetConsumingEnumerable();
        //    dbContext.AddRange(items);
        //    dbContext.SaveChanges();
        //}

        // 以下写法可以正常运行,但如果队列中有10条待处理日志,就会产生10次数据库交互(性能不佳)
        foreach (var item in queue.GetConsumingEnumerable())
        {
            dbContext.Add(item);
            dbContext.SaveChanges();
        }
    }

    public Task StartAsync(CancellationToken cancellationToken)
    {
        jobTask.Start();
        return Task.CompletedTask;
    }

    public Task StopAsync(CancellationToken cancellationToken)
    {
        queue.Complete();
        jobTask.Wait(); // 疑问:应该使用'await jobTask'?
        return Task.CompletedTask;
    }
}

依赖注入配置

builder.Services.AddSingleton<LogQueue>();
builder.Services.AddHostedService<DbLogWorker>();

控制器使用示例

[HttpGet("/")]
public IActionResult Get(string? name = "N/A")
{
    var log = new ActivityLog()
    {
        CreationTime = DateTime.UtcNow,
        Action = "Home page visit",
        IPAddress = HttpContext.Connection.RemoteIpAddress ?? IPAddress.Any,
        Metadata = $"{{ name: {name} }}",
        UserId = User.FindFirstValue(ClaimTypes.NameIdentifier)
    };
    queue.Add(log);
    return Ok("Welcome!");
}

示例项目结构如下图所示:
示例项目结构

问题汇总

  1. 批量读取队列数据一次性写入的代码无法运行,只能逐条写入触发多次数据库提交,如何实现批量写入?
  2. 队列设置100万容量是否足够?
  3. 服务停止时等待任务完成应该用Wait()还是await?
  4. 如何增加消费工作线程数量?如何根据队列积压量动态调整工作线程数?

解决方案

1. 批量写入代码失效原因与修复

批量代码无法运行的核心原因是GetConsumingEnumerable()是阻塞式无限枚举,只要队列没有标记为完成,它就会一直等待新元素加入,永远不会返回,所以AddRange会一直卡着等新数据,根本走不到SaveChanges步骤。
正确的批量消费逻辑要加两个阈值:批量大小阈值和超时阈值,避免一直等待数据攒批:

  • 攒够N条(比如100条)就立刻写入
  • 哪怕没攒够N条,等待超过M毫秒(比如500ms)也把现有数据写入,防止低流量下日志长时间不入库
    修复后的消费逻辑参考:
// 把BlockingCollection改为可访问,或者给LogQueue加TryTake方法包装
public class LogQueue
{
    private const int QueueCapacity = 1_000_000;
    private readonly BlockingCollection<ActivityLog> logs = new(QueueCapacity);
    // 新增TryTake包装
    public bool TryTake(out ActivityLog log, int millisecondsTimeout) 
        => logs.TryTake(out log, millisecondsTimeout);
    // 其余原有方法保持不变
}

private void Job()
{
    // 批大小和超时时间建议放到配置文件中
    const int batchSize = 100;
    const int timeoutMs = 500;
    var batch = new List<ActivityLog>(batchSize);

    while (!queue.IsCompleted)
    {
        // 先尝试取第一条,阻塞等待直到有数据或队列完成
        if (queue.TryTake(out var log, timeoutMs))
        {
            batch.Add(log);
        }

        // 循环取当前队列里所有可用数据,不阻塞,取完立刻退出
        while (batch.Count < batchSize && queue.TryTake(out var nextLog, 0))
        {
            batch.Add(nextLog);
        }

        // 批次有数据就写入
        if (batch.Count > 0)
        {
            try
            {
                // DbContext是轻量级对象,不要长期持有,每次批量写入新建Scope用完即释放
                using var scope = scf.CreateScope();
                var dbContext = scope.ServiceProvider.GetRequiredService<ApplicationDbContext>();
                dbContext.AddRange(batch);
                dbContext.SaveChanges();
            }
            catch (Exception ex)
            {
                // 此处必须加失败处理逻辑:比如重试、落本地文件兜底,禁止直接吞异常丢日志
            }
            finally
            {
                batch.Clear();
            }
        }
    }

    // 队列标记完成后,把剩余未写入的最后一批数据写完
    if (batch.Count > 0)
    {
        using var scope = scf.CreateScope();
        var dbContext = scope.ServiceProvider.GetRequiredService<ApplicationDbContext>();
        dbContext.AddRange(batch);
        dbContext.SaveChanges();
    }
}

不要长期持有同一个DbContext,DbContext的ChangeTracker跟踪过多实体会导致性能持续下降,每次批量操作新建Scope是EF Core的最佳实践。

2. 100万队列容量是否足够

这个值没有通用答案,需要结合业务情况判断:

  • 按单条日志1KB估算,100万条日志大概占用1GB内存,大部分应用服务器都可以承受
  • 核心判断逻辑是峰值日志产生速度和消费速度的差值:比如峰值每秒产生1万条日志,数据库消费速度每秒3000条,每秒积压7000条,100万容量可以扛大概140秒的流量峰值,超过这个时间队列打满,Add操作会阻塞,反过来影响正常业务请求
    如果业务峰值很高,不建议设置过大的固定容量,队列满了之后要做降级处理,比如丢弃非核心日志、写本地文件兜底,绝对不能让队列阻塞主请求流程。

3. StopAsync里用Wait还是await

必须用await,Wait()是同步阻塞调用,在ASP.NET Core环境中很容易触发线程池死锁。
直接把StopAsync改成异步实现即可:

public async Task StopAsync(CancellationToken cancellationToken)
{
    queue.Complete();
    await jobTask;
}

另外手动new Task加LongRunning选项的写法没有必要,直接在StartAsync中用Task.Factory.StartNew指定LongRunning选项创建任务即可。

4. 多消费线程与动态扩缩容实现

固定数量消费线程

BlockingCollection本身是线程安全的,多个线程同时消费不会出现数据竞争,直接启动N个消费任务即可:

private Task[] jobTasks;
private const int FixedWorkerCount = 4; // 根据数据库写入能力配置固定线程数

public Task StartAsync(CancellationToken cancellationToken)
{
    jobTasks = Enumerable.Range(0, FixedWorkerCount)
        .Select(_ => Task.Factory.StartNew(Job, TaskCreationOptions.LongRunning))
        .ToArray();
    return Task.CompletedTask;
}

public async Task StopAsync(CancellationToken cancellationToken)
{
    queue.Complete();
    await Task.WhenAll(jobTasks);
}

动态调整线程数

如果需要根据队列积压量动态调整线程数,有两个实现方向:

  1. 基于BlockingCollection.Count属性获取当前积压量,设置扩缩容阈值:比如积压超过1万条就新增工作线程,积压低于1000条就主动退出多余工作线程。注意必须设置线程数上下限(比如最少1个,最多8个),避免线程过多把数据库连接打满。
  2. 更推荐的方案是用.NET内置的Channel<T>替换BlockingCollection<T>,配合System.Threading.Tasks.Dataflow库中的ActionBlock<T>实现,ActionBlock天然支持并行度控制,也支持运行时动态调整并行度,不需要自己手写线程管理逻辑,性能比BlockingCollection更高,API也更简洁。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 08:12:12