当前位置:   article > 正文

一个缓存泛型自动处理队列,留有处理事件接口

一个缓存泛型自动处理队列,留有处理事件接口

在C#中,创建一个通用泛型缓存队列工具通常涉及到使用泛型(T)来支持任意类型的队列项,并且实现先进先出(FIFO)的队列逻辑。为了提供自动处理和事件通知的功能,我们可以使用后台线程来处理队列中的项,并通过事件或委托来通知监听者队列的状态变化。

以下是一个简单的泛型缓存队列工具的示例,它使用了System.Collections.Concurrent命名空间下的BlockingCollection作为线程安全的队列,并通过Task在后台自动处理队列中的项。它还定义了一个事件,当队列中的项被处理时可以触发该事件。

csharp
using System;  
using System.Collections.Concurrent;  
using System.Threading;  
using System.Threading.Tasks;  
  
public class GenericCacheQueue<T>  
{  
    private readonly BlockingCollection<T> _queue = new BlockingCollection<T>();  
    private Task _processingTask;  
    private CancellationTokenSource _cancellationTokenSource;  
  
    // 定义事件,当项被处理时触发  
    public event Action<T> ItemProcessed;  
  
    // 添加项到队列中  
    public void Enqueue(T item)  
    {  
        _queue.Add(item);  
    }  
  
    // 开始处理队列中的项  
    public void StartProcessing()  
    {  
        _cancellationTokenSource = new CancellationTokenSource();  
        _processingTask = Task.Run(() => ProcessQueue(_cancellationTokenSource.Token), _cancellationTokenSource.Token);  
    }  
  
    // 停止处理队列中的项并清理资源  
    public void StopProcessing()  
    {  
        _cancellationTokenSource.Cancel();  
        try  
        {  
            _processingTask.Wait();  
        }  
        catch (AggregateException ae)  
        {  
            // 处理任务取消的异常  
            foreach (var innerException in ae.InnerExceptions)  
            {  
                if (innerException is OperationCanceledException)  
                {  
                    // 忽略由于取消而产生的异常  
                }  
                else  
                {  
                    // 处理其他异常  
                    throw;  
                }  
            }  
        }  
        _queue.CompleteAdding(); // 标记添加完成,以便消费者知道何时退出循环  
    }  
  
    // 后台处理队列中的项  
    private void ProcessQueue(CancellationToken cancellationToken)  
    {  
        foreach (var item in _queue.GetConsumingEnumerable(cancellationToken))  
        {  
            // 处理项  
            ProcessItem(item);  
        }  
    }  
  
    // 处理单个项的逻辑  
    private void ProcessItem(T item)  
    {  
        // 在这里添加处理项的逻辑  
        Console.WriteLine($"Processing item: {item}");  
        // 触发事件通知项已被处理  
        ItemProcessed?.Invoke(item);  
    }  


}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
  • 31
  • 32
  • 33
  • 34
  • 35
  • 36
  • 37
  • 38
  • 39
  • 40
  • 41
  • 42
  • 43
  • 44
  • 45
  • 46
  • 47
  • 48
  • 49
  • 50
  • 51
  • 52
  • 53
  • 54
  • 55
  • 56
  • 57
  • 58
  • 59
  • 60
  • 61
  • 62
  • 63
  • 64
  • 65
  • 66
  • 67
  • 68
  • 69
  • 70
  • 71
  • 72
  • 73
  • 74
  • 75
  • 76

使用上述GenericCacheQueue类时,你需要实现自己的处理逻辑,并在需要的时候触发ItemProcessed事件。下面是如何使用这个队列的示例:

csharp
class Program  
{  
    static void Main(string[] args)  
    {  
        var cacheQueue = new GenericCacheQueue<int>();  
        cacheQueue.ItemProcessed += CacheQueue_ItemProcessed; // 订阅事件  
  
        cacheQueue.StartProcessing(); // 开始处理队列中的项  
  
        // 添加一些项到队列中  
        for (int i = 0; i < 10; i++)  
        {  
            cacheQueue.Enqueue(i);  
        }  
  
        // 等待一段时间,让后台线程处理队列中的项  
        Thread.Sleep(2000);  
  
        cacheQueue.StopProcessing(); // 停止处理并清理资源  
    }  
  
    private static void CacheQueue_ItemProcessed(int item)  
    {  
        // 当项被处理时,这里的代码将被执行  
        Console.WriteLine($"Item {item} processed.");  
    }  
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28

请注意,上述代码示例是一个简化的版本,用于演示如何构建这样的工具。在实际应用中,你可能需要添加更多的错误处理逻辑、日志记录、性能优化以及更复杂的处理逻辑。此外,对于大型或生产环境的应用,你可能需要考虑使用更高级的并发控制机制或队列服务(如RabbitMQ、Azure Service Bus等)。

声明:本文内容由网友自发贡献,不代表【wpsshop博客】立场,版权归原作者所有,本站不承担相应法律责任。如您发现有侵权的内容,请联系我们。转载请注明出处:https://www.wpsshop.cn/w/知新_RL/article/detail/502018
推荐阅读
相关标签
  

闽ICP备14008679号