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

面试题咨询:基于.NET标准API的内存型拉取队列库C#实现方案

设计思路

本实现完全基于标准.NET API开发,未使用PriorityQueue等受限类型,核心设计逻辑匹配所有需求点:

  • 多独立队列维护:采用ConcurrentDictionary做全局队列容器,键为队列唯一标识,天然支持线程安全的队列动态创建、查询
  • 全量消息投递:放弃传统队列消费即删除的逻辑,队列存储全量消息快照,每个消费者独立维护消费偏移量,确保所有消费者都能拉取到完整的未过期消息
  • 过期清理机制:利用消息按生产时间顺序追加的特性,无需优先队列即可实现高效清理:每次消费/生产时触发懒清理,同时搭配后台定时任务批量扫全量过期消息,判断过期的规则为取队列最大保留周期和消息自定义TTL的最小值作为生效过期时间
  • 多生产多消费安全:所有公共操作加锁保证原子性,消费者偏移量存储在线程安全字典中,支持任意数量的生产者、消费者并发接入
核心代码实现

消息实体定义

public class QueueMessage<T>
{
    public T Data { get; set; }
    public DateTime CreatedAt { get; set; }
    public TimeSpan? Ttl { get; set; }

    public bool IsExpired(TimeSpan queueMaxRetention)
    {
        // 消息TTL优先级低于队列全局最大保留周期,取较小值作为生效过期时长
        var effectiveTtl = Ttl.HasValue 
            ? TimeSpan.FromTicks(Math.Min(Ttl.Value.Ticks, queueMaxRetention.Ticks)) 
            : queueMaxRetention;
        return DateTime.Now - CreatedAt > effectiveTtl;
    }
}

队列实现

public interface IMemoryQueue<T>
{
    /// <summary>
    /// 生产消息
    /// </summary>
    /// <param name="data">消息内容</param>
    /// <param name="ttl">消息自定义TTL,不填则使用队列全局最大保留周期</param>
    void Produce(T data, TimeSpan? ttl = null);
    /// <summary>
    /// 消费者拉取消息
    /// </summary>
    /// <param name="consumerId">消费者唯一标识</param>
    /// <param name="data">输出拉取到的消息</param>
    /// <returns>是否拉取到未过期消息</returns>
    bool TryConsume(string consumerId, out T data);
}

public class MemoryQueue<T> : IMemoryQueue<T>
{
    private readonly List<QueueMessage<T>> _messages = new List<QueueMessage<T>>();
    private readonly ConcurrentDictionary<string, int> _consumerOffsets = new ConcurrentDictionary<string, int>();
    private readonly TimeSpan _maxRetention;
    private readonly object _lock = new object();

    public MemoryQueue(TimeSpan maxRetention)
    {
        _maxRetention = maxRetention;
        // 启动后台定时清理任务,不需要可关闭改为纯懒清理模式
        _ = RunBackgroundCleanupAsync();
    }

    public void Produce(T data, TimeSpan? ttl = null)
    {
        lock (_lock)
        {
            _messages.Add(new QueueMessage<T>
            {
                Data = data,
                CreatedAt = DateTime.Now,
                Ttl = ttl
            });
        }
    }

    public bool TryConsume(string consumerId, out T data)
    {
        data = default;
        // 新消费者首次接入默认从最早未过期消息开始消费
        var currentOffset = _consumerOffsets.GetOrAdd(consumerId, 0);

        lock (_lock)
        {
            // 先清理一轮过期消息
            CleanupExpiredMessages();
            if (currentOffset >= _messages.Count) return false;

            var targetMsg = _messages[currentOffset];
            if (targetMsg.IsExpired(_maxRetention)) return false;

            // 更新消费者偏移量
            _consumerOffsets[consumerId] = currentOffset + 1;
            data = targetMsg.Data;
            return true;
        }
    }

    private void CleanupExpiredMessages()
    {
        // 消息按生成时间顺序排列,仅需删除开头连续的过期消息即可
        var removeCount = 0;
        while (removeCount < _messages.Count && _messages[removeCount].IsExpired(_maxRetention))
        {
            removeCount++;
        }
        if (removeCount == 0) return;

        _messages.RemoveRange(0, removeCount);
        // 同步更新所有消费者的偏移量,避免越界
        foreach (var consumerId in _consumerOffsets.Keys)
        {
            _consumerOffsets[consumerId] = Math.Max(0, _consumerOffsets[consumerId] - removeCount);
        }
    }

    private async Task RunBackgroundCleanupAsync()
    {
        while (true)
        {
            await Task.Delay(TimeSpan.FromSeconds(10));
            lock (_lock)
            {
                CleanupExpiredMessages();
            }
        }
    }
}

全局队列管理器

public static class MemoryQueueManager
{
    private static readonly ConcurrentDictionary<string, object> _queueContainer = new ConcurrentDictionary<string, object>();

    /// <summary>
    /// 获取或创建独立队列
    /// </summary>
    /// <typeparam name="T">队列消息类型</typeparam>
    /// <param name="queueName">队列唯一名称</param>
    /// <param name="maxRetention">队列全局最大消息保留周期</param>
    public static IMemoryQueue<T> GetOrCreateQueue<T>(string queueName, TimeSpan maxRetention)
    {
        return (IMemoryQueue<T>)_queueContainer.GetOrAdd(queueName, _ => new MemoryQueue<T>(maxRetention));
    }
}
使用示例
// 创建名称为order_queue的队列,全局最大保留周期1小时
var orderQueue = MemoryQueueManager.GetOrCreateQueue<string>("order_queue", TimeSpan.FromHours(1));

// 生产者发送消息,自定义TTL为30分钟
orderQueue.Produce("order_001", TimeSpan.FromMinutes(30));
orderQueue.Produce("order_002"); // 不填TTL则默认用队列1小时保留周期

// 消费者1拉取消息
while (orderQueue.TryConsume("consumer_01", out var order))
{
    Console.WriteLine($"消费者1处理:{order}");
}

// 消费者2拉取消息,可获取全量未过期消息,和消费者1的偏移量完全独立
while (orderQueue.TryConsume("consumer_02", out var order))
{
    Console.WriteLine($"消费者2处理:{order}");
}
可选优化点
  • 并发量要求高的场景可将lock替换为ReaderWriterLockSlim,区分读写场景提升吞吐量
  • 需要阻塞拉取消息的场景可搭配AutoResetEvent实现无消息时消费者阻塞,有新消息生产时自动唤醒
  • 大流量场景可将消息分段存储,清理时直接删除整个过期段,进一步提升清理效率
  • 可增加消息确认机制,消费失败支持重置偏移量重新消费

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 12:48:04