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

基于Channel的ASP.NET Core托管服务任务并发处理问题

问题描述

我按照微软的托管服务教程创建了托管服务,程序运行正常,但需要实现队列任务的并发处理。我的应用有4个客户端,每个客户端可同时处理4个任务,因此任意时刻需要并行处理16个任务。

我已经做了部分代码修改:
在MonitorLoop类中:

private int count = 0;
private async ValueTask MonitorAsync()
{
    while (!_cancellationToken.IsCancellationRequested)
    {
        await _taskQueue.QueueAsync(BuildWorkItem);
        Interlocked.Increment(ref count);
        Console.WriteLine($"Count: {count}");
    }
}

同一类中还有:

if (delayLoop == 3)
{
    _logger.LogInformation("Queued Background Task {Guid} is complete.", guid);
    Interlocked.Decrement(ref count);
}

测试发现,当设置Capacity为4时,count值不会超过5,队列满时会等待空闲位置再添加,但当前任务是逐个串行处理的。

以下是QueuedHostedService类中BackgroundProcessing方法的代码:

private async Task BackgroundProcessing(CancellationToken stoppingToken)
{
    while (!stoppingToken.IsCancellationRequested)
    {
        var workItem = await TaskQueue.DequeueAsync(stoppingToken);

        try
        {
            //instead of getting a single item from the queue, somehow, here
            //we should be able to process them in parallel for 4 clients
            //with a limit for maximum items each client can process
            await workItem(stoppingToken);
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Error occurred executing {WorkItem}.", nameof(workItem));
        }
    }
}

我希望实现任务并行处理,不确定使用Channel作为队列是否是最优方案,也考虑过改用ConcurrentQueue,但不清楚如何搭建支持4个客户端各4线程的可靠架构。

解决方案

一、队列选型结论:继续用Channel更合适

Channel是.NET专为异步生产-消费场景设计的高效队列,自带背压(Backpressure)机制,无需自己实现队列满时的阻塞逻辑,比ConcurrentQueue更适配你的异步托管服务场景,完全不需要替换。

二、核心实现思路

要满足4个客户端各4个并发任务(总计16个并行度),核心是:

  1. 给每个任务标记所属客户端ID
  2. 用信号量(SemaphoreSlim)精准控制单个客户端的并发上限
  3. 启动足够的消费线程(总计16个)来并行处理队列任务

三、具体代码修改

1. 定义带客户端标识的任务结构

让每个待处理任务携带客户端ID,方便跟踪单客户端并发数:

public record ClientWorkItem(int ClientId, Func<CancellationToken, Task> Work);

2. 重构队列接口与实现

将原有队列改为存储ClientWorkItem,保持Channel的背压特性:

public interface IClientTaskQueue
{
    ValueTask QueueAsync(ClientWorkItem workItem);
    ValueTask<ClientWorkItem> DequeueAsync(CancellationToken cancellationToken);
}

public class ClientTaskQueue : IClientTaskQueue
{
    private readonly Channel<ClientWorkItem> _channel;

    public ClientTaskQueue(int capacity)
    {
        var options = new BoundedChannelOptions(capacity)
        {
            FullMode = BoundedChannelFullMode.Wait
        };
        _channel = Channel.CreateBounded<ClientWorkItem>(options);
    }

    public async ValueTask QueueAsync(ClientWorkItem workItem)
    {
        await _channel.Writer.WriteAsync(workItem);
    }

    public async ValueTask<ClientWorkItem> DequeueAsync(CancellationToken cancellationToken)
    {
        return await _channel.Reader.ReadAsync(cancellationToken);
    }
}

3. 修改托管服务的消费逻辑

用每个客户端专属的信号量控制并发,同时启动16个并行消费任务:

public class QueuedHostedService : BackgroundService
{
    private readonly ILogger<QueuedHostedService> _logger;
    private readonly IClientTaskQueue _taskQueue;
    private readonly Dictionary<int, SemaphoreSlim> _clientSemaphores;

    public QueuedHostedService(ILogger<QueuedHostedService> logger, IClientTaskQueue taskQueue)
    {
        _logger = logger;
        _taskQueue = taskQueue;
        // 初始化4个客户端的信号量,每个最大并发4
        _clientSemaphores = Enumerable.Range(1, 4)
            .ToDictionary(id => id, _ => new SemaphoreSlim(4, 4));
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Queued Hosted Service is running.");

        // 启动16个并行消费任务(4客户端×4并发)
        var processingTasks = Enumerable.Range(0, 16)
            .Select(_ => ProcessQueueItemsAsync(stoppingToken))
            .ToList();

        await Task.WhenAll(processingTasks);
    }

    private async Task ProcessQueueItemsAsync(CancellationToken stoppingToken)
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            var workItem = await _taskQueue.DequeueAsync(stoppingToken);
            var semaphore = _clientSemaphores[workItem.ClientId];

            try
            {
                await semaphore.WaitAsync(stoppingToken);
                await workItem.Work(stoppingToken);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Error executing task for client {ClientId}", workItem.ClientId);
            }
            finally
            {
                semaphore.Release();
            }
        }
    }

    public override async Task StopAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Queued Hosted Service is stopping.");

        foreach (var semaphore in _clientSemaphores.Values)
        {
            semaphore.Dispose();
        }

        await base.StopAsync(stoppingToken);
    }
}

4. 调整任务生成与入队逻辑

在MonitorLoop中生成带客户端ID的任务,可通过轮询方式分配客户端:

private int _currentClientId = 0;

private ClientWorkItem BuildWorkItem()
{
    // 轮询分配客户端ID(1-4)
    var clientId = Interlocked.Increment(ref _currentClientId) % 4;
    clientId = clientId == 0 ? 4 : clientId;

    return new ClientWorkItem(clientId, async token =>
    {
        var guid = Guid.NewGuid();
        var delayLoop = 0;
        _logger.LogInformation("Queued Background Task {Guid} started for client {ClientId}.", guid, clientId);

        while (!token.IsCancellationRequested && delayLoop < 3)
        {
            await Task.Delay(1000, token);
            delayLoop++;
            _logger.LogInformation("Queued Background Task {Guid} is running for client {ClientId} ({DelayLoop}/3).", guid, clientId, delayLoop);
        }

        if (delayLoop == 3)
        {
            _logger.LogInformation("Queued Background Task {Guid} is complete for client {ClientId}.", guid, clientId);
            Interlocked.Decrement(ref count);
        }
    });
}

// 修改MonitorAsync的入队代码
private async ValueTask MonitorAsync()
{
    while (!_cancellationToken.IsCancellationRequested)
    {
        await _taskQueue.QueueAsync(BuildWorkItem());
        Interlocked.Increment(ref count);
        Console.WriteLine($"Count: {count}");
        // 可根据业务需求添加入队间隔,避免瞬间填满队列
        // await Task.Delay(100, _cancellationToken);
    }
}

四、关键说明

  • SemaphoreSlim确保每个客户端最多同时处理4个任务,16个消费线程保证整体并行度上限为16
  • Channel的背压机制会在队列满时自动阻塞生产者,避免内存溢出,无需手动实现队列容量控制
  • 若需要更精细化的客户端任务分配,可调整BuildWorkItem中的客户端ID分配逻辑,比如按任务类型绑定客户端

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 03:35:20