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!"); }
示例项目结构如下图所示:
问题汇总
- 批量读取队列数据一次性写入的代码无法运行,只能逐条写入触发多次数据库提交,如何实现批量写入?
- 队列设置100万容量是否足够?
- 服务停止时等待任务完成应该用
Wait()还是await? - 如何增加消费工作线程数量?如何根据队列积压量动态调整工作线程数?
解决方案
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); }
动态调整线程数
如果需要根据队列积压量动态调整线程数,有两个实现方向:
- 基于
BlockingCollection.Count属性获取当前积压量,设置扩缩容阈值:比如积压超过1万条就新增工作线程,积压低于1000条就主动退出多余工作线程。注意必须设置线程数上下限(比如最少1个,最多8个),避免线程过多把数据库连接打满。 - 更推荐的方案是用.NET内置的
Channel<T>替换BlockingCollection<T>,配合System.Threading.Tasks.Dataflow库中的ActionBlock<T>实现,ActionBlock天然支持并行度控制,也支持运行时动态调整并行度,不需要自己手写线程管理逻辑,性能比BlockingCollection更高,API也更简洁。
内容的提问来源于stack exchange,提问作者Parsa99
相关产品推荐
相关产品推荐

