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

C#中类AWS SQS/Service Bus的内存消息队列实现优化问询

云队列消息处理的并发控制与数据结构优化问题

问题描述

我需要从云队列中出队消息并处理,要求每个Worker实例可自定义并发处理消息数(如100或10条),以充分利用Worker容量。目前使用C# Channels连接生产者与消费者,通过ConcurrentDictionary跟踪处理任务避免超出并发上限,但未找到具备类似AWS SQS/Azure Service Bus的peek&lock特性的内存数据结构(消息被消费后锁定,处理完成再删除,期间不可被其他消费者获取),想了解是否有更优的数据结构替代当前方案。

同时针对现有代码有两点疑问:

  1. ProduceAsync函数需在并发处理数低于配置值时从云队列拉取消息写入Channel,否则等待Channel有可用空间;
  2. ConsumeAsync当前以fire-and-forget方式调用ProcessMessageAsync,导致所有消息并发处理,希望利用现有ConcurrentDictionary优化,避免额外维护Task列表来await Task.WhenAll。

现有实现代码

internal class Program
{
    static readonly Random random = new Random();
    private ConcurrentDictionary<string, Task> messageProcessingTasksMap =
        new ConcurrentDictionary<string, Task>();

    static async Task Main(string[] args)
    {
        CancellationTokenSource cts = new CancellationTokenSource();
        Console.CancelKeyPress += (s, e) =>
        {
            e.Cancel = true;
            cts.Cancel();
        };

        var channel = Channel.CreateBounded<string>(new BoundedChannelOptions(100)
        {
            SingleReader = true,
            SingleWriter = true,
            AllowSynchronousContinuations = false
        });

        var program = new Program();
        var statsTask = program.TaskStatsAsync(channel.Reader, cts.Token);
        var producerTask = program.ProduceAsyc(channel.Writer, cts.Token);
        var consumerTask = program.ConsumeAsync(channel.Reader, cts.Token);

        await Task.WhenAll(statsTask, producerTask, consumerTask);
        Console.WriteLine("Press any key to exit...");
        Console.Read();
    }

    private async Task ProduceAsyc(ChannelWriter<string> channelWriter,
        CancellationToken cancellationToken)
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            // 100 将从配置读取,不同Worker实例值不同
            if (messageProcessingTasksMap.Count < 100)
            {
                int messagesToDequeue = 100 - messageProcessingTasksMap.Count;
                if (messagesToDequeue > 10)
                {
                    messagesToDequeue = 10;
                }
        
                for (int i = 0; i < messagesToDequeue; i++)
                {
                    // 实际场景中这里会从云队列出队消息并写入Channel
                    await channelWriter.WriteAsync($"Message-{i}", cancellationToken);
                }
            }
            // 原代码未处理并发满时的等待,添加短暂延迟避免空转
            await Task.Delay(100, cancellationToken);
        }

        Console.WriteLine("Stopping producer. Marking channel complete");
        channelWriter.Complete();
        Console.WriteLine("Producer stopped...");
    }

    private async Task ConsumeAsync(ChannelReader<string> channelReader,
        CancellationToken cancellationToken)
    {
        while (await channelReader.WaitToReadAsync(cancellationToken))
        {
            while (messageProcessingTasksMap.Count < 100 &&
                channelReader.TryRead(out var message))
            {
                var taskId = Guid.NewGuid().ToString();
                var processingTask = ProcessMessageAsync(message, taskId, cancellationToken);
                messageProcessingTasksMap.TryAdd(taskId, processingTask);
                // 任务完成后自动从字典移除,避免内存膨胀
                _ = processingTask.ContinueWith(t => 
                {
                    messageProcessingTasksMap.TryRemove(taskId, out _);
                }, cancellationToken);
            }
        }

        // 等待所有在处理的任务完成
        await Task.WhenAll(messageProcessingTasksMap.Values);
        Console.WriteLine("Stopping consumer...");
        Console.WriteLine($"Channel Count: {channelReader.Count}");
    }

    private async Task ProcessMessageAsync(string message, string taskId, CancellationToken cancellationToken)
    {
        // 模拟业务处理延迟
        await Task.Delay(random.Next(100, 300), cancellationToken);
        // 移除操作放在ContinueWith中更可靠,避免异常时未清理
        // messageProcessingTasksMap.Remove(taskId, out var _);
    }

    private async Task TaskStatsAsync(ChannelReader<string> channelReader,
        CancellationToken cancellationToken)
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            Console.WriteLine(
                $"Number of tasks being processed: {messageProcessingTasksMap.Count}" +
                $", Channel Count: {channelReader.Count}");
            await Task.Delay(1000, cancellationToken);
        }

        Console.WriteLine("Exiting Stats Task");
    }
}

优化方案与疑问解答

1. 具备Peek&Lock特性的内存数据结构替代方案

.NET内置数据结构没有直接支持peek&lock的,但可以基于ConcurrentQueue+ConcurrentDictionary自定义实现:

  • 用ConcurrentQueue<T>存储待处理消息(建议用带唯一ID的消息对象,而非纯字符串)
  • 用ConcurrentDictionary<string, object>存储已锁定的消息ID,标记为"处理中"状态
  • 消费者逻辑:
    1. 从队列Dequeue一条消息
    2. 尝试将消息ID加入字典,成功则锁定完成,开始处理
    3. 如果加入失败(说明被其他消费者锁定),则将消息重新放回队列尾部
    4. 处理完成后从字典移除ID;若处理失败,可选择将消息重新入队并移除锁定

这种实现完全可控,贴合peek&lock的核心逻辑,无需依赖第三方组件。

2. ProduceAsync函数优化

原代码在并发满时会空循环浪费CPU,优化思路:

  • 用SemaphoreSlim替代ConcurrentDictionary跟踪并发量,通过WaitAsync等待可用槽位
  • 结合Channel的WriteAsync阻塞特性,当Channel满时自动等待,无需额外判断

优化后的示例:

private readonly SemaphoreSlim _concurrencySemaphore;
private readonly int _maxConcurrency;
private readonly int _maxBatchSize = 10;

public Program(int maxConcurrency = 100)
{
    _maxConcurrency = maxConcurrency;
    _concurrencySemaphore = new SemaphoreSlim(maxConcurrency);
}

private async Task ProduceAsyc(ChannelWriter<string> channelWriter, CancellationToken cancellationToken)
{
    while (!cancellationToken.IsCancellationRequested)
    {
        int availableSlots = _concurrencySemaphore.CurrentCount;
        if (availableSlots <= 0)
        {
            // 等待有任务完成释放槽位
            await _concurrencySemaphore.WaitAsync(cancellationToken);
            _concurrencySemaphore.Release(); // 释放仅为触发后续计算
            continue;
        }

        int messagesToDequeue = Math.Min(availableSlots, _maxBatchSize);
        // 实际场景替换为云队列批量拉取API
        var messages = Enumerable.Range(0, messagesToDequeue).Select(i => $"Message-{Guid.NewGuid()}");
        
        foreach (var msg in messages)
        {
            // Channel满时自动等待,无需手动判断
            await channelWriter.WriteAsync(msg, cancellationToken);
        }
    }

    channelWriter.Complete();
    Console.WriteLine("Producer stopped...");
}

3. ConsumeAsync函数优化

原fire-and-forget的问题可通过SemaphoreSlim结合异步等待解决,无需额外维护Task列表:

  • 处理消息前先获取信号量,完成后释放,直接控制并发上限
  • 用await foreach遍历Channel消息,逻辑更简洁
  • 无需ConcurrentDictionary跟踪任务,信号量直接管控并发

优化后的示例:

private async Task ConsumeAsync(ChannelReader<string> channelReader, CancellationToken cancellationToken)
{
    var processingTasks = new List<Task>();
    await foreach (var message in channelReader.ReadAllAsync(cancellationToken))
    {
        await _concurrencySemaphore.WaitAsync(cancellationToken);
        // 任务完成后自动释放信号量
        var task = ProcessMessageAsync(message, cancellationToken)
            .ContinueWith(t => _concurrencySemaphore.Release(), cancellationToken);
        processingTasks.Add(task);
    }

    // 等待所有处理任务完成
    await Task.WhenAll(processingTasks);
    Console.WriteLine("Stopping consumer...");
}

private async Task ProcessMessageAsync(string message, CancellationToken cancellationToken)
{
    await Task.Delay(random.Next(100, 300), cancellationToken);
    Console.WriteLine($"Processed message: {message}");
}

这种方式既避免了fire-and-forget的失控问题,又简化了代码逻辑,可靠性更高。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 08:25:58