**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.**
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.
Kiểm tra khả năng hoạt động
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;
}
}
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);
}
}
[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);
}
}