面试题咨询:基于.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
相关产品推荐
相关产品推荐

