如何解决System.Threading.Channels多消费者重复处理同记录问题?
问题:System.Threading.Channels重复处理同一条InspectionFileModel记录
使用System.Threading.Channels时,写入者将InspectionFileModel入队后,消费者开始处理该模型,但此时写入者从数据库读取到同一条记录(因处理中状态未及时更新),导致另一消费者重复处理。尝试用ConcurrentDictionary解决但无效,相关代码如下:
public sealed class SendImagesBackgroundService : BackgroundTask { private readonly ILogger<SendImagesBackgroundService> _logger; private readonly IServiceProvider _serviceProvider; private readonly IApiService _apiService; private readonly Channel<InspectionFileModel> _channel; private readonly SendImagesBackgroundServiceOptions _options; private static readonly ConcurrentDictionary<string, bool> s_concurrentDictionary = new(); public SendImagesBackgroundService( ILogger<SendImagesBackgroundService> logger, IServiceProvider serviceProvider, IApiService apiService, IOptions<SendImagesBackgroundServiceOptions> options, IBackgroundTaskLockProvider lockProvider) : base(logger) { UseExclusiveLock(lockProvider); _options = options.Value; _logger = logger; _serviceProvider = serviceProvider; _apiService = apiService; _channel = Channel.CreateBounded<InspectionFileModel>( new BoundedChannelOptions(_options.ChannelQueueSize)); } protected override Task ExecuteAsync(CancellationToken cancellationToken) { Task.Run(() => FetchImages(cancellationToken), cancellationToken); Task.Run(() => StartProcessingInspectionImages(cancellationToken), cancellationToken); return Task.CompletedTask; } private async Task FetchImages(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { using var serviceScope = _serviceProvider.CreateScope(); List<InspectionFileModel>? inspectionImages = default; do { try { var mediatr = serviceScope.ServiceProvider .GetRequiredService<IMediator>(); inspectionImages = await mediatr.Send( new DownloadInspectionFilesQueryNew(),cancellationToken); if (inspectionImages is not {Count: > 0}) { continue; } var nonProcessingImages = inspectionImages.Where( x => !s_concurrentDictionary.ContainsKey(x.Id)); foreach (var image in nonProcessingImages) { s_concurrentDictionary.TryAdd(image.Id, true); while (!(cancellationToken.IsCancellationRequested && await _channel.Writer.WaitToWriteAsync( cancellationToken))) { if (!_channel.Writer.TryWrite(image)) { continue; } break; } } } catch (Exception ex) when (!(cancellationToken .IsCancellationRequested && ex is OperationCanceledException)) { _logger.LogError(ex, "Fetching inspection images Failed"); } } while (!cancellationToken.IsCancellationRequested && inspectionImages is {Count: > 0}); await Task.Delay( _options.DelayBetweenFetchBatchMs, cancellationToken); } _logger.LogInformation("End fetching inspection images"); } private async Task StartProcessingInspectionImages( CancellationToken cancellationToken) { var parallelProcesses = new List<Task>(); for (int i = 0; i < _options.NumberOfParallelTasks; i++) { var task = Task.Run(() => ProcessInspectionImages(cancellationToken), cancellationToken); parallelProcesses.Add(task); } await Task.WhenAll(parallelProcesses); } private async Task ProcessInspectionImages(CancellationToken cancellationToken) { while (!(cancellationToken.IsCancellationRequested && await _channel.Reader.WaitToReadAsync(cancellationToken))) { while (!cancellationToken.IsCancellationRequested && _channel.Reader.TryRead(out var inspectionImage )) { try { await SendInspectionImageToLivo(inspectionImage, cancellationToken); } catch (Exception ex) when (!(cancellationToken .IsCancellationRequested && ex is OperationCanceledException)) { //handle } } } } private async Task SendInspectionImageToLivo(InspectionFileModel image, CancellationToken cancellationToken) { try { //send data over the network } catch (ApiException ex) { //handle } finally { s_concurrentDictionary.TryRemove(image.Id, out bool _); } } public override object? GetTelemetry() => null; }
状态在SendInspectionImageToLivo方法中更新:若返回200则设为成功状态,若返回4**则设为失败状态,后续数据库查询将不再包含这些记录。
问题根源
- 非原子操作导致重复入队:
ContainsKey判断和TryAdd是两步独立操作,多线程环境下可能同时通过过滤条件,导致同一个ID被多次添加到字典并写入通道。 - Channel写入逻辑错误:原代码中
WaitToWriteAsync的条件判断逻辑混乱,取消令牌触发时可能仍执行写入操作,加剧重复问题。 - 数据库层面无状态锁:内存字典仅能在进程内生效,无法阻止数据库层面的重复读取(比如多个实例部署时)。
解决方案
1. 原子化判断与添加操作
将“检查存在+添加”改为原子操作,直接用TryAdd的返回值判断是否允许写入:
// 替换FetchImages方法中的过滤和循环代码 foreach (var image in inspectionImages) { // 原子操作:只有ID不存在时才添加并写入通道 if (s_concurrentDictionary.TryAdd(image.Id, true)) { while (!cancellationToken.IsCancellationRequested) { if (await _channel.Writer.WaitToWriteAsync(cancellationToken)) { if (_channel.Writer.TryWrite(image)) { break; } } else { // 取消触发时移除ID,避免内存泄漏 s_concurrentDictionary.TryRemove(image.Id, out _); break; } } } }
2. 修复Channel写入的取消逻辑
确保取消令牌触发时及时清理字典,避免无效ID残留:
while (!cancellationToken.IsCancellationRequested) { if (await _channel.Writer.WaitToWriteAsync(cancellationToken)) { if (_channel.Writer.TryWrite(image)) break; } else { s_concurrentDictionary.TryRemove(image.Id, out _); break; } }
3. 数据库层面原子标记处理中(终极保障)
从根源避免重复读取,修改DownloadInspectionFilesQueryNew的实现,用数据库事务原子性地将查询到的记录标记为“处理中”:
-- 以MySQL为例,原子查询并更新状态 BEGIN TRANSACTION; SELECT * FROM InspectionFiles WHERE Status = 'Pending' LIMIT @BatchSize FOR UPDATE; UPDATE InspectionFiles SET Status = 'Processing' WHERE Status = 'Pending' LIMIT @BatchSize; COMMIT;
这样后续查询不会再获取到已标记为处理中的记录,彻底解决跨进程/多实例的重复问题。
4. 异常处理中的字典清理优化
根据处理结果决定是否移除字典中的ID,避免可重试场景下的重复入队:
private async Task SendInspectionImageToLivo(InspectionFileModel image, CancellationToken cancellationToken) { bool shouldRemove = true; try { var response = await _apiService.SendImageAsync(image, cancellationToken); if (response.IsSuccessStatusCode) { await UpdateStatus(image.Id, "Success"); } else if ((int)response.StatusCode >= 400 && (int)response.StatusCode < 500) { await UpdateStatus(image.Id, "Failed"); } else { // 5xx等可重试错误,保留ID并重新入队 shouldRemove = false; await _channel.Writer.WriteAsync(image, cancellationToken); } } catch (ApiException ex) when (IsRetryable(ex)) { shouldRemove = false; await _channel.Writer.WriteAsync(image, cancellationToken); } catch (Exception ex) { await UpdateStatus(image.Id, "Failed"); } finally { if (shouldRemove) { s_concurrentDictionary.TryRemove(image.Id, out _); } } }
内容的提问来源于stack exchange,提问作者tabsandze
相关产品推荐
相关产品推荐

