Icomm.RabbitMQ.OpenTelemetry (0.6.2.3)
Published 2026-05-29 17:14:36 +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.jsondotnet add package --source ic --version 0.6.2.3 Icomm.RabbitMQ.OpenTelemetryAbout this package
RabbitMQ
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=>objectMPSerializer: sử dụng thư việnMessagePack, cần khai báoMessagePackObjectvới mỗi thực thể đã định nghĩaMessagePackTypelessSerializer: sử dụng thư việnMessagePackkhông cần khai báo MetadataMessagePackObject
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*/);
}
Dependencies
Details
2026-05-29 17:14:36 +07:00
Assets (2)
Versions (8)
View all
NuGet
25
long.nguyen
18 KiB