RabbitMQ
در RabbitMQ فرستنده پیام را به یک Exchange میدهد، نه مستقیم به صف. Exchange با قانونهای اتصال تصمیم میگیرد پیام به کدام صفها برود. هر پیام بعد از تأیید مصرفکننده از صف حذف میشود.
نویسنده: bezzad
مشکل: کارهای کند بعد از ثبت سفارش
در فروشگاه اینترنتی ما، بعد از ثبت سفارش باید یک ایمیل رسید فرستاده شود. سرور ایمیل گاهی چند ثانیه طول میکشد و گاهی هم جواب نمیدهد.
اگر سرویس سفارش منتظر ایمیل بماند:
- مشتری منتظر میماند. صفحه پرداخت چند ثانیه هیچ جوابی نمیدهد.
- خطای ایمیل، سفارش را خراب میکند. سرور ایمیل خراب است، پس ثبت سفارش هم خطا میدهد.
- روز حراج همه چیز کند میشود. هزار سفارش در یک دقیقه یعنی هزار ایمیل در همان دقیقه.
راه حل: سرویس سفارش فقط یک پیام در صف میگذارد و سریع جواب میدهد. سرویس ایمیل پیامها را با سرعت خودش از صف برمیدارد. RabbitMQ یک Message Broker است که همین صفها را نگه میدارد.
ایده اصلی: فرستنده صف را نمیشناسد
در RabbitMQ فرستنده پیام را مستقیم در صف نمیگذارد. پیام را به یک Exchange میدهد. سه مفهوم اصلی داریم:
- صف (Queue). پیامها را تا وقتی یک مصرفکننده آنها را بردارد نگه میدارد. هر سرویس معمولاً صف خودش را دارد.
- مبادلهگر (Exchange). پیام را میگیرد و تصمیم میگیرد به کدام صفها برود. خودش پیامی نگه نمیدارد.
- اتصال (Binding). یک قانون است که میگوید «این صف، پیامهایی با این کلید را میخواهد».
هر پیام یک Routing Key دارد. مثلاً پیام ثبت سفارش کلید order.placed دارد و پیام لغو سفارش کلید order.cancelled.
چرا این جدایی مفید است؟
- فرستنده ساده میماند. سرویس سفارش فقط میگوید «سفارش ثبت شد». نمیداند چند سرویس گوش میدهند.
- سرویس جدید بدون تغییر فرستنده. سرویس گزارش فقط یک صف و یک Binding جدید میسازد.
- هر صف مستقل است. اگر سرویس انبار یک ساعت خاموش باشد، پیامهایش در صف خودش میماند. ایمیلها بدون مشکل ادامه دارند.
نوعهای Exchange
- نوع direct. کلید پیام باید دقیقاً با کلید Binding یکی باشد. برای وقتی که هر نوع کار صف جدا دارد.
- نوع fanout. کلید را نگاه نمیکند. پیام به همه صفهای وصلشده میرود. برای خبرهایی که همه باید بشنوند.
- نوع topic. کلید با یک الگو مقایسه میشود. علامت ستاره یعنی دقیقاً یک کلمه و علامت # یعنی صفر یا چند کلمه. برای رویدادهای کسبوکار، معمولاً بهترین انتخاب است.
- نوع headers. به جای کلید، Header های پیام را مقایسه میکند. کمتر استفاده میشود.
چند مصرفکننده روی یک صف
در روز حراج، یک نسخه از سرویس ایمیل کافی نیست. پس سه pod از سرویس ایمیل را به همان صف وصل میکنیم.
- هر پیام فقط به یکی از آنها میرسد. RabbitMQ پیامها را بین مصرفکنندههای یک صف پخش میکند. به این الگو Competing Consumers میگوییم.
- تعداد مصرفکننده محدودیت ساختاری ندارد. برعکس کافکا، لازم نیست تعداد pod را با تعداد Partition تنظیم کنی.
- ولی ترتیب تضمین نمیشود. وقتی چند مصرفکننده همزمان کار میکنند، یا پیامی دوباره به صف برمیگردد، ترتیب پردازش با ترتیب ارسال فرق میکند.
تأیید پیام (Ack)
پیام چه وقت از صف حذف میشود؟ وقتی مصرفکننده بگوید «کارم تمام شد». به این کار Ack میگوییم.
- حالت Ack. کار موفق بود. پیام حذف میشود.
- حالت Nack بدون برگشت به صف. کار ممکن نیست. اگر صف یک Dead Letter Exchange داشته باشد، پیام به صف خطا میرود.
- بدون جواب. اگر مصرفکننده بمیرد، پیام دوباره به صف برمیگردد و به مصرفکننده دیگری میرسد.
نتیجه مهم: پیام ممکن است دو بار برسد. مثلاً ایمیل فرستاده شد، ولی قبل از Ack سرویس مُرد. پس مصرفکننده باید Idempotent باشد، دقیقاً مثل کافکا.
محدودیت پیام همزمان (Prefetch)
اگر محدودیتی نگذاری، RabbitMQ هر تعداد پیام را بدون Ack به یک مصرفکننده میدهد. یک pod صدها پیام میگیرد و pod های دیگر بیکار میمانند. با تنظیم Prefetch میگویی «به هر مصرفکننده حداکثر این تعداد پیام تأییدنشده بده».
کد
نسخههای جدید کتابخانه RabbitMQ.Client متدهای async دارند. اول صف و Exchange را تعریف میکنیم و پیام ثبت سفارش را میفرستیم:
using RabbitMQ.Client;
using System.Text;
using System.Text.Json;
var factory = new ConnectionFactory { HostName = "localhost" };
using var connection = await factory.CreateConnectionAsync();
using var channel = await connection.CreateChannelAsync();
await channel.ExchangeDeclareAsync("shop.orders", ExchangeType.Topic, durable: true);
await channel.QueueDeclareAsync(
queue: "email-queue", durable: true, exclusive: false, autoDelete: false,
arguments: new Dictionary<string, object?>
{
["x-queue-type"] = "quorum", // replicated, safe queue
["x-dead-letter-exchange"] = "shop.orders.dlx" // rejected messages go here
});
await channel.QueueBindAsync("email-queue", "shop.orders", routingKey: "order.*");
var body = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(new OrderPlaced("o-3", 250_000)));
var props = new BasicProperties { Persistent = true, MessageId = "o-3" };
await channel.BasicPublishAsync("shop.orders", "order.placed",
mandatory: true, basicProperties: props, body: body);
public sealed record OrderPlaced(string OrderId, decimal Total);
حالا سرویس ایمیل. Ack را خودمان و بعد از کار میفرستیم:
await channel.BasicQosAsync(prefetchSize: 0, prefetchCount: 10, global: false);
var consumer = new AsyncEventingBasicConsumer(channel);
consumer.ReceivedAsync += async (_, ea) =>
{
try
{
var order = JsonSerializer.Deserialize<OrderPlaced>(ea.Body.Span)!;
await emailSender.SendReceiptOnceAsync(order); // idempotent by OrderId
await channel.BasicAckAsync(ea.DeliveryTag, multiple: false);
}
catch (Exception)
{
// Do not requeue forever: send it to the dead letter exchange.
await channel.BasicNackAsync(ea.DeliveryTag, multiple: false, requeue: false);
}
};
await channel.BasicConsumeAsync("email-queue", autoAck: false, consumer: consumer);
پیام گم نشود
برای اینکه پیام با restart شدن سرور RabbitMQ گم نشود، چند چیز با هم لازم است:
- صف ماندگار (durable). تعریف صف بعد از restart میماند.
- پیام ماندگار (Persistent). خود پیام روی دیسک نوشته میشود.
- صف Quorum. پیام روی چند سرور کپی میشود. اگر یک سرور بمیرد، صف از بین نمیرود.
- تأیید فرستنده (Publisher Confirms). RabbitMQ به فرستنده خبر میدهد که پیام را واقعاً گرفته است. بدون آن، فرستنده نمیداند پیام رسیده یا نه.
- تأیید دستی مصرفکننده. تنظیم autoAck خاموش باشد و Ack بعد از کار فرستاده شود.
اشتباههای رایج
| اشتباه | نتیجه | راه درست |
|---|---|---|
| تنظیم autoAck روشن | پیام همان لحظه تحویل حذف میشود. اگر سرویس بمیرد، پیام گم شده است. | تأیید دستی، بعد از کار. |
| بدون Prefetch | یک مصرفکننده همه پیامها را میگیرد و بقیه بیکارند. | یک عدد معقول برای Prefetch. |
| Nack با برگشت به صف برای هر خطا | پیام خراب در یک چرخه بیپایان میچرخد. | چند تلاش با فاصله، بعد صف خطا. |
| فرستنده مستقیم به اسم صف | هر سرویس جدید یعنی تغییر کد فرستنده. | فرستادن به Exchange با Routing Key. |
| انتظار ترتیب دقیق با چند مصرفکننده | پیام لغو قبل از ثبت پردازش میشود. | اگر ترتیب لازم است، یک مصرفکننده یا ابزاری مثل کافکا. |
| مصرفکننده غیر Idempotent | مشتری دو ایمیل رسید میگیرد. | چک کردن شناسه پیام قبل از کار. |
RabbitMQ یا Kafka؟
RabbitMQ
- صف کار است. پیام بعد از Ack حذف میشود.
- مسیریابی منعطف با Exchange و Binding.
- هر پیام جدا تأیید یا رد میشود.
- تعداد مصرفکننده آزاد است.
- مناسب برای کارهای پسزمینه و دستورها.
Kafka
- دفتر ثبت است. پیام بعد از خواندن میماند.
- میشود پیامهای قدیمی را دوباره خواند.
- ترتیب داخل هر Partition حفظ میشود.
- مصرفکننده مفید حداکثر به تعداد Partition.
- مناسب برای جریان رویداد با حجم زیاد و چند خواننده.
خلاصه در شش خط
- فرستنده پیام را به Exchange میدهد، نه به صف.
- صفها با Binding و Routing Key به Exchange وصل میشوند.
- نوع topic برای رویدادهای کسبوکار معمولاً بهترین انتخاب است.
- چند مصرفکننده روی یک صف، کار را تقسیم میکنند، ولی ترتیب را حفظ نمیکنند.
- تأیید دستی بعد از کار، Prefetch معقول و صف خطا را فراموش نکن.
- پیام ممکن است دو بار برسد. مصرفکننده باید Idempotent باشد.