Levelwise
فارسی
پیام‌رسانی

RabbitMQ

در RabbitMQ فرستنده پیام را به یک Exchange می‌دهد، نه مستقیم به صف. Exchange با قانون‌های اتصال تصمیم می‌گیرد پیام به کدام صف‌ها برود. هر پیام بعد از تأیید مصرف‌کننده از صف حذف می‌شود.

بازبینی نشدهبا کمک AI نوشته شدهزمان خواندن: ۱۴ دقیقهمثال فروشگاه اینترنتیکد C# و RabbitMQ.Client

نویسنده: bezzad

مشکل: کارهای کند بعد از ثبت سفارش

در فروشگاه اینترنتی ما، بعد از ثبت سفارش باید یک ایمیل رسید فرستاده شود. سرور ایمیل گاهی چند ثانیه طول می‌کشد و گاهی هم جواب نمی‌دهد.

اگر سرویس سفارش منتظر ایمیل بماند:

  1. مشتری منتظر می‌ماند. صفحه پرداخت چند ثانیه هیچ جوابی نمی‌دهد.
  2. خطای ایمیل، سفارش را خراب می‌کند. سرور ایمیل خراب است، پس ثبت سفارش هم خطا می‌دهد.
  3. روز حراج همه چیز کند می‌شود. هزار سفارش در یک دقیقه یعنی هزار ایمیل در همان دقیقه.

راه حل: سرویس سفارش فقط یک پیام در صف می‌گذارد و سریع جواب می‌دهد. سرویس ایمیل پیام‌ها را با سرعت خودش از صف برمی‌دارد. RabbitMQ یک Message Broker است که همین صف‌ها را نگه می‌دارد.

ایده اصلی: فرستنده صف را نمی‌شناسد

در RabbitMQ فرستنده پیام را مستقیم در صف نمی‌گذارد. پیام را به یک Exchange می‌دهد. سه مفهوم اصلی داریم:

  1. صف (Queue). پیام‌ها را تا وقتی یک مصرف‌کننده آن‌ها را بردارد نگه می‌دارد. هر سرویس معمولاً صف خودش را دارد.
  2. مبادله‌گر (Exchange). پیام را می‌گیرد و تصمیم می‌گیرد به کدام صف‌ها برود. خودش پیامی نگه نمی‌دارد.
  3. اتصال (Binding). یک قانون است که می‌گوید «این صف، پیام‌هایی با این کلید را می‌خواهد».

هر پیام یک Routing Key دارد. مثلاً پیام ثبت سفارش کلید order.placed دارد و پیام لغو سفارش کلید order.cancelled.

سرویس سفارشorder.placedExchangeshop.orderstype: topicorder.*order.placed#email-queuestock-queuereport-queueایمیلانبارگزارشکلید روی هر خط یک قانون اتصال است
صف ایمیل همه پیام‌های سفارش را می‌خواهد. صف انبار فقط ثبت سفارش را. صف گزارش همه چیز را.

چرا این جدایی مفید است؟

  1. فرستنده ساده می‌ماند. سرویس سفارش فقط می‌گوید «سفارش ثبت شد». نمی‌داند چند سرویس گوش می‌دهند.
  2. سرویس جدید بدون تغییر فرستنده. سرویس گزارش فقط یک صف و یک Binding جدید می‌سازد.
  3. هر صف مستقل است. اگر سرویس انبار یک ساعت خاموش باشد، پیام‌هایش در صف خودش می‌ماند. ایمیل‌ها بدون مشکل ادامه دارند.

نوع‌های Exchange

directکلید باید دقیقاً یکی باشدkey: emailemailsmsfanoutکلید مهم نیست، همه می‌گیرندany keyemailreporttopicکلید با یک الگو مقایسه می‌شودorder.cancelledorder.**.placed
در هر سه حالت، فرستنده همان کار را می‌کند. فقط قانون انتخاب صف فرق دارد.
  • نوع direct. کلید پیام باید دقیقاً با کلید Binding یکی باشد. برای وقتی که هر نوع کار صف جدا دارد.
  • نوع fanout. کلید را نگاه نمی‌کند. پیام به همه صف‌های وصل‌شده می‌رود. برای خبرهایی که همه باید بشنوند.
  • نوع topic. کلید با یک الگو مقایسه می‌شود. علامت ستاره یعنی دقیقاً یک کلمه و علامت # یعنی صفر یا چند کلمه. برای رویدادهای کسب‌وکار، معمولاً بهترین انتخاب است.
  • نوع headers. به جای کلید، Header های پیام را مقایسه می‌کند. کمتر استفاده می‌شود.
یک نکته کوچک: یک Exchange پیش‌فرض بدون نام هم هست. اگر پیام را به آن بدهی و Routing Key را اسم صف بگذاری، پیام مستقیم به همان صف می‌رود. برای مثال‌های ساده خوب است، ولی فرستنده را به اسم صف وابسته می‌کند.

چند مصرف‌کننده روی یک صف

در روز حراج، یک نسخه از سرویس ایمیل کافی نیست. پس سه pod از سرویس ایمیل را به همان صف وصل می‌کنیم.

  1. هر پیام فقط به یکی از آن‌ها می‌رسد. RabbitMQ پیام‌ها را بین مصرف‌کننده‌های یک صف پخش می‌کند. به این الگو Competing Consumers می‌گوییم.
  2. تعداد مصرف‌کننده محدودیت ساختاری ندارد. برعکس کافکا، لازم نیست تعداد pod را با تعداد Partition تنظیم کنی.
  3. ولی ترتیب تضمین نمی‌شود. وقتی چند مصرف‌کننده همزمان کار می‌کنند، یا پیامی دوباره به صف برمی‌گردد، ترتیب پردازش با ترتیب ارسال فرق می‌کند.

تأیید پیام (Ack)

پیام چه وقت از صف حذف می‌شود؟ وقتی مصرف‌کننده بگوید «کارم تمام شد». به این کار Ack می‌گوییم.

email-queueپیام‌های منتظرتحویلسرویس ایمیلprefetch = 10ackموفق بود، پیام حذف شودnackصف خطاemail-queue.dlqاگر سرویس بمیرد و جوابی ندهد،پیام دوباره به صف برمی‌گرددپیامی که درست نمی‌شود، در صف خطا منتظر بررسی می‌ماند
تا وقتی Ack نیامده، پیام در اختیار RabbitMQ است. اگر اتصال مصرف‌کننده قطع شود، پیام دوباره تحویل داده می‌شود.
  • حالت 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);
چرا requeue را خاموش کردیم؟ اگر یک پیام خراب باشد (مثلاً JSON غلط)، هر بار که به صف برگردد دوباره شکست می‌خورد. این یک چرخه بی‌پایان است. پیامی که درست نمی‌شود باید به صف خطا برود. برای خطاهای موقت، مثل قطعی کوتاه سرور ایمیل، چند بار با فاصله تلاش کن. کتابخانه‌ای مثل MassTransit این کار را آماده دارد.

پیام گم نشود

برای اینکه پیام با restart شدن سرور RabbitMQ گم نشود، چند چیز با هم لازم است:

  1. صف ماندگار (durable). تعریف صف بعد از restart می‌ماند.
  2. پیام ماندگار (Persistent). خود پیام روی دیسک نوشته می‌شود.
  3. صف Quorum. پیام روی چند سرور کپی می‌شود. اگر یک سرور بمیرد، صف از بین نمی‌رود.
  4. تأیید فرستنده (Publisher Confirms). RabbitMQ به فرستنده خبر می‌دهد که پیام را واقعاً گرفته است. بدون آن، فرستنده نمی‌داند پیام رسیده یا نه.
  5. تأیید دستی مصرف‌کننده. تنظیم autoAck خاموش باشد و Ack بعد از کار فرستاده شود.
مراقب Dual Write باش: ذخیره سفارش در دیتابیس و فرستادن پیام به RabbitMQ دو کار جدا هستند. اگر وسط این دو کار سرویس بمیرد، پیام هیچ وقت نمی‌رود. اینجا هم راه درست الگوی Outbox است.

اشتباه‌های رایج

اشتباه نتیجه راه درست
تنظیم autoAck روشن پیام همان لحظه تحویل حذف می‌شود. اگر سرویس بمیرد، پیام گم شده است. تأیید دستی، بعد از کار.
بدون Prefetch یک مصرف‌کننده همه پیام‌ها را می‌گیرد و بقیه بیکارند. یک عدد معقول برای Prefetch.
Nack با برگشت به صف برای هر خطا پیام خراب در یک چرخه بی‌پایان می‌چرخد. چند تلاش با فاصله، بعد صف خطا.
فرستنده مستقیم به اسم صف هر سرویس جدید یعنی تغییر کد فرستنده. فرستادن به Exchange با Routing Key.
انتظار ترتیب دقیق با چند مصرف‌کننده پیام لغو قبل از ثبت پردازش می‌شود. اگر ترتیب لازم است، یک مصرف‌کننده یا ابزاری مثل کافکا.
مصرف‌کننده غیر Idempotent مشتری دو ایمیل رسید می‌گیرد. چک کردن شناسه پیام قبل از کار.

RabbitMQ یا Kafka؟

RabbitMQ

  • صف کار است. پیام بعد از Ack حذف می‌شود.
  • مسیریابی منعطف با Exchange و Binding.
  • هر پیام جدا تأیید یا رد می‌شود.
  • تعداد مصرف‌کننده آزاد است.
  • مناسب برای کارهای پس‌زمینه و دستورها.

Kafka

  • دفتر ثبت است. پیام بعد از خواندن می‌ماند.
  • می‌شود پیام‌های قدیمی را دوباره خواند.
  • ترتیب داخل هر Partition حفظ می‌شود.
  • مصرف‌کننده مفید حداکثر به تعداد Partition.
  • مناسب برای جریان رویداد با حجم زیاد و چند خواننده.
قانون ساده: اگر سؤال تو «این کار را یک نفر انجام دهد» است، RabbitMQ ساده‌تر است. اگر سؤال «این اتفاق افتاد، هر کس لازم دارد بخواند، حتی بعداً» است، کافکا مناسب‌تر است.

خلاصه در شش خط

  1. فرستنده پیام را به Exchange می‌دهد، نه به صف.
  2. صف‌ها با Binding و Routing Key به Exchange وصل می‌شوند.
  3. نوع topic برای رویدادهای کسب‌وکار معمولاً بهترین انتخاب است.
  4. چند مصرف‌کننده روی یک صف، کار را تقسیم می‌کنند، ولی ترتیب را حفظ نمی‌کنند.
  5. تأیید دستی بعد از کار، Prefetch معقول و صف خطا را فراموش نکن.
  6. پیام ممکن است دو بار برسد. مصرف‌کننده باید Idempotent باشد.