📨 معماری رویدادمحور (Event-Driven Architecture) در NET. با RabbitMQ
⚡️ معماری رویدادمحور (EDA) میتواند برنامهها را منعطفتر و قابلاعتمادتر کند.
به جای اینکه یک بخش سیستم مستقیماً بخش دیگر را فراخوانی کند، اجازه میدهیم رویدادها از طریق پیامرسان (Message Broker) جریان پیدا کنند.
📚 در این راهنمای سریع، یک سیستم ساده رویدادمحور در NET. با استفاده از RabbitMQ پیادهسازی میکنیم.
📌 سناریوی ما
یک تولیدکننده (Producer) رویدادها را ارسال میکند و یک مصرفکننده (Consumer) آنها را دریافت میکند.
برای تست، RabbitMQ را در یک کانتینر Docker اجرا میکنیم (با فعال بودن رابط کاربری مدیریت (Management UI) تا بتوانیم فعالیتها را مشاهده کنیم).
از پکیج رسمی RabbitMQ.Client در یک اپلیکیشن کنسول NET. استفاده خواهیم کرد.
🐳 اجرای RabbitMQ با Docker
اگر RabbitMQ را نصب نکردهاید، میتوانید آن را به سرعت با Docker اجرا کنید:
docker run -it --rm --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:4-management
📍 این دستور یک RabbitMQ Broker روی localhost راهاندازی میکند:
🔌 پورت 5672 → پروتکل AMQP
🌐 پورت 15672 → رابط مدیریت در آدرس:
http://localhost:15672
🛠 مفاهیم پایه RabbitMQ
🏭 Producer → برنامهای که پیامها (رویدادها) را به RabbitMQ ارسال میکند.
📥 Consumer → برنامهای که پیامها را از یک صف دریافت میکند.
📦 Queue → مانند یک صندوق پستی که پیامها را ذخیره میکند. مصرفکنندگان از صفها میخوانند.
🔀 Exchange → مکانیسم مسیردهی که پیامهای دریافتی از تولیدکنندگان را به صفها هدایت میکند.
💡 نکته: در RabbitMQ، تولیدکنندگان هرگز مستقیماً به یک صف ارسال نمیکنند، بلکه به یک Exchange ارسال میکنند. Exchange تعیین میکند که پیام به کدام صف یا صفها برود.🚀 تولیدکننده (Producer) – ارسال رویداد
فرض کنید رویداد OrderPlaced داریم که میتواند سرویسهای پاییندستی مثل انبار، ایمیل اطلاعرسانی و غیره را فعال کند.
var factory = new ConnectionFactory() { HostName = "localhost" };
using var connection = await factory.CreateConnectionAsync();
using var channel = await connection.CreateChannelAsync();
await channel.QueueDeclareAsync(
queue: "orders",
durable: true,
exclusive: false,
autoDelete: false,
arguments: null);
var orderPlaced = new OrderPlaced
{
OrderId = Guid.NewGuid(),
Total = 99.99,
CreatedAt = DateTime.UtcNow
};
var message = JsonSerializer.Serialize(orderPlaced);
var body = Encoding.UTF8.GetBytes(message);
await channel.BasicPublishAsync(
exchange: string.Empty,
routingKey: "orders",
mandatory: true,
basicProperties: new BasicProperties { Persistent = true },
body: body);
Console.WriteLine($"Sent: {message}");📌 نکات مهم:
📂 صف durable است → بعد از ریست RabbitMQ باقی میماند.
💾 پیام Persistent است → روی دیسک ذخیره میشود.
🔤 دادهها به JSON سریالایز شده و به بایت UTF-8 ارسال میشوند.
🎯 مصرفکننده (Consumer) – دریافت رویداد
مصرفکننده به همان صف متصل میشود و پیامها را دریافت میکند:
var factory = new ConnectionFactory() { HostName = "localhost" };
using var connection = await factory.CreateConnectionAsync();
using var channel = await connection.CreateChannelAsync();
await channel.QueueDeclareAsync(
queue: "orders",
durable: true,
exclusive: false,
autoDelete: false,
arguments: null);
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (sender, eventArgs) =>
{
byte[] body = eventArgs.Body.ToArray();
string message = Encoding.UTF8.GetString(body);
var orderPlaced = JsonSerializer.Deserialize<OrderPlaced>(message);
Console.WriteLine($"Received: OrderPlaced - {orderPlaced.OrderId}");
await ((AsyncEventingBasicConsumer)sender)
.Channel.BasicAckAsync(eventArgs.DeliveryTag, multiple: false);
};
await channel.BasicConsumeAsync("orders", autoAck: false, consumer);
Console.WriteLine("Waiting for messages...");📌 نکات مهم:
⏳ autoAck: false → پیام فقط پس از پردازش موفق تأیید میشود.
♻️ اگر پردازش شکست بخورد، میتوان با BasicNack آن را مجدداً در صف قرار داد یا به Dead-Letter Queue فرستاد.