Sử dụng System.Threading.Channels cho mô hình xuất bản/đăng ký trong tiến trình

**Trước tiên, Channel về bản chất là một kiểu tập hợp mới trong .NET, tương tự như `Queue` nhưng có nhiều điểm khác biệt.** **System.Threading.Channels được giới thiệu từ .NET Core 3.0, sở hữu API bất đồng bộ, hiệu năng cao và an toàn đa luồng. Nó có thể được sử dụng làm hàng đợi tin nhắn để sản xuất và tiêu thụ dữ liệu thông qua các API công khai `Writer` và `Reader`, tương ứng với vai trò nhà sản xuất và người tiêu dùng. Khác với các hàng đợi khác như RabbitMQ, Channel hoạt động hoàn toàn trong cùng một tiến trình.**

namespace TestWebApplication.Channels
{
    public class QueueManager<TItem>
    {
        private static readonly Lazy lazy = new(() => new QueueManager<TItem>());
        
        private ChannelReader<TItem> _queueReader;
        private ChannelWriter<TItem> _queueWriter;

        public ChannelReader<TItem> QueueReader => _queueReader;
        public ChannelWriter<TItem> QueueWriter => _queueWriter;

        private QueueManager()
        {
            var channel = Channel.CreateBounded<TItem>(new BoundedChannelOptions(100)
            {
                FullMode = BoundedChannelFullMode.Wait
            });

            _queueReader = channel.Reader;
            _queueWriter = channel.Writer;
        }

        public static QueueManager<TItem> GetInstance => lazy.Value;
    }
}
Phương thức `WaitToReadAsync` sẽ chờ đến khi kênh có thể đọc được. Phương thức `ReadAsync(stoppingToken)` sẽ đọc nội dung từ kênh.

public class QueueProcessingService : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await Task.Factory.StartNew(async () =>
        {
            while (await QueueManager<string>.GetInstance.QueueReader.WaitToReadAsync(stoppingToken))
            {
                string message = await QueueManager<string>.GetInstance.QueueReader.ReadAsync(stoppingToken);
                await Console.Out.WriteLineAsync($"Tin nhắn nhận được: {message}");
            }
        }, TaskCreationOptions.LongRunning);
    }
}
Kiểm tra khả năng hoạt động

[Route("api/[controller]")]
[ApiController]
public class QueueController : ControllerBase
{
    [HttpGet]
    public async Task<IActionResult> Send(string message)
    {
        await QueueManager<string>.GetInstance.QueueWriter.WriteAsync(message);
        return Ok(message);
    }

    [HttpGet("Receive")]
    public async Task<IActionResult> Receive()
    {
        string message = await QueueManager<string>.GetInstance.QueueReader.ReadAsync();
        return Ok(message);
    }
}

Thẻ: System.Threading.Channels .NET-Core C#-Concurrency

Đăng vào ngày 7 tháng 8 lúc 05:26