ic

Icomm.RabbitMQ (0.4.6.3)

Published 2025-05-05 22:23:24 +07:00 by Long Nguyễn Đình - Core

Installation

dotnet nuget add source --name ic --username your_username --password your_token https://git.icomm.vn/api/packages/ic/nuget/index.json
dotnet add package --source ic --version 0.4.6.3 Icomm.RabbitMQ

About this package

RabbitMQ ultility

Thư viện hỗ trợ kết nối RabbitMQ

Bao bồm 3 thành phần chính

  • Serializer-Deserializer
  • Producer
  • Consumer

1. DependencyInjection

services
.AddRabbitMQJsonSerializer<PostContent[]>()
// Khởi tạo Producer
// resolve - IProducer<PostContent[]>
.AddRabbitMQConfig<IProducer<PostContent[]>, GenericProducer<PostContent[]>>("posts-offline-1", lifetime: ServiceLifetime.Singleton)
// Khởi tạo consumer
// c1 - resolve FileExtractContentConsumer
.AddRabbitMQConfig<FileExtractContentConsumer>(key: "new-files-extractcontent", lifetime: ServiceLifetime.Singleton)
// c2 - resolve IConsumer<PostContent[]>
.AddRabbitMQConfig<IConsumer<PostContent[]>, FileExtractContentConsumer>(key: "new-files-extractcontent", lifetime: ServiceLifetime.Singleton)
;

2. Serializer-Deserializer

Như đã biết RabbitMQ chỉ lưu message dưới dạng byte[] nên để thuận tiện trong quá trình coding tương tác với thực thể (object) chúng ta sử dụng middleware để biến đổi byte[] thành object

Interface IGenericSerializer<T>

Thư viện hỗ trợ 1 số phương thức biết đổi sau:

  • JsonEncodingSerializer: byte[] <=Encoding=> string <=JsonConvert=> object
  • MPSerializer: sử dụng thư viện MessagePack, cần khai báo MessagePackObject với mỗi thực thể đã định nghĩa
  • MessagePackTypelessSerializer: sử dụng thư viện MessagePack không cần khai báo Metadata MessagePackObject

3. Producer

Interface IProducer<T> với T là thực thể lưu trữ trên queue thông qua Serializer-Deserializer

Sử dụng:

//Thực thể định nghĩa ví dụ
class TModel
{
	public string Name { get; set; }
}
var qc = new QueueConfig(); //Cấu hình queue, có thể lấy từ ConfigManager
//Lớp middleware Serializer, có thể chọn các lớp hỗ trợ sẵn hoặc tự định nghĩa bằng interface `IGenericSerializer`
var serializer = new JsonEncodingSerializer<TModel>(); 
//Khởi tạo cách 1
IProducer<TModel> producer = new GenericProducer<TModel>(qc, serializer);
//Khởi tạo cách 2, Cách này thường sử dụng với DI để khởi tạo Producer sau đó set cấu hình queue
IProducer<TModel> producer = new GenericProducer<TModel>(serializer);
producer.SetQueueConfig(qc);
//Đẩy thực thể lên queue
producer.Publish(new TModel(){ Name = "hello" }); 

4. Consumer

Interface IConsumer<T> với T là thực thể lưu trữ trên queue thông qua Serializer-Deserializer

4.1 Định nghĩa Class Consumer:

class TModelConsumer : GenericConsumer<TModel>
{
	// Khởi tạo ctor cách 1
	public TModelConsumer(QueueConfig qc, IGenericSerializer<TModel> serializer = null) : base(qc, serializer)
	{
	}

	// Khởi tạo ctor cách 2
	public TModelConsumer(IGenericSerializer<TModel> serializer = null) : base(serializer)
	{
	}
	
	// Hàm xử lý các message lấy từ queue về
	public override void Execute(TModel item)
	{
		// Làm gì đó ở đây với TModel item
	}
}

//Tạo thực thể Consumer
IConsumer<TModel> consumer = new TModelConsumer(biến ctor);
// Bắt đầu consumer vào queue để pull message về
consumer.Start();
//Consumer cũng có thể Publish 1 message
consumer.Publish(new TModel(){ Name = "hello" }); 
// Dừng Consumer
consumer.Stop();

4.1 Consumer nhanh với ActionConsumer:

var qc = new QueueConfig(); //Cấu hình queue, có thể lấy từ ConfigManager hoặc tự gán thuộc tính
//Lớp middleware Serializer, có thể chọn các lớp hỗ trợ sẵn hoặc tự định nghĩa bằng interface `IGenericSerializer`
var serializer = new JsonEncodingSerializer<TModel>(); 

IConsumer<TModel> consumer = new ActionConsumer<TModel>(qc,
(item) => 
{
	// Làm gì đó ở đây với TModel item
},
serializer);
// Bắt đầu consumer vào queue để pull message về
consumer.Start();
// Dừng Consumer
consumer.Stop();

Chú ý

Để tránh blocking Task khi sử dụng cấu hình QueueConfig.MaxWorkpool > 1

public override async Task ExecuteAsync(<Object> item, BasicDeliverEventArgs e, CancellationToken cancellationToken)
{
	// Gọi hàm Async
	// Để tránh blocking cần đẩy hàm async này vào Task Pool để xử lý bằng cách
	// Lấy kết quả trả về
	var result = Task.Run(()=> /*...Hàm Async*/).Result;
	// Không lấy kết quả trả về
	Task.Run(()=> /*...Hàm Async*/);
	// Hàm không có tham biến
	Task.Run(/*Hàm Async*/);
}
  • Điều chỉnh rule tắt chương trình khi có lỗi nghiêm trọng

Dependencies

ID Version Target Framework
Icomm.Configs.Abstract 0.1.7 .NETStandard2.1
Icomm.Core 0.0.2.23 .NETStandard2.1
Microsoft.Extensions.Hosting.Abstractions 5.0.0 .NETStandard2.1
RabbitMQ.Client 6.8.1 .NETStandard2.1
System.Diagnostics.DiagnosticSource 7.0.2 .NETStandard2.1
System.Runtime.Loader 4.3.0 .NETStandard2.1
Details
NuGet
2025-05-05 22:23:24 +07:00
33
long.nguyen
44 KiB
Assets (2)
Versions (47) View all
0.6.2.3 2026-05-29
0.6.2.2 2026-05-28
0.6.1 2026-02-06
0.6.0 2026-01-30
0.2.9 2025-05-06