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

如何解决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**则设为失败状态,后续数据库查询将不再包含这些记录。


问题根源

  1. 非原子操作导致重复入队:ContainsKey判断和TryAdd是两步独立操作,多线程环境下可能同时通过过滤条件,导致同一个ID被多次添加到字典并写入通道。
  2. Channel写入逻辑错误:原代码中WaitToWriteAsync的条件判断逻辑混乱,取消令牌触发时可能仍执行写入操作,加剧重复问题。
  3. 数据库层面无状态锁:内存字典仅能在进程内生效,无法阻止数据库层面的重复读取(比如多个实例部署时)。

解决方案

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 17:11:10