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

MassTransit

کتابخانه MassTransit کارهای تکراری پیام‌رسانی را در .NET انجام می‌دهد. تو فقط پیام و مصرف‌کننده را می‌نویسی. ساخت صف‌ها، تلاش دوباره، صف خطا، Outbox و Saga را MassTransit آماده دارد.

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

نویسنده: bezzad

مشکل: کد تکراری دور هر پیام

در درس RabbitMQ دیدیم که برای یک پیام ساده چقدر کد لازم است. در فروشگاه اینترنتی ما ده‌ها نوع پیام داریم. برای هر کدام باید این کارها را بنویسیم:

  1. ساخت Exchange، صف و Binding.
  2. تبدیل شیء به JSON و برعکس.
  3. تلاش دوباره وقتی سرور ایمیل یک لحظه جواب نمی‌دهد.
  4. فرستادن پیام خراب به صف خطا.
  5. ثبت پیام در Outbox تا با دیتابیس هماهنگ باشد.
  6. ساختن Scope برای Dependency Injection برای هر پیام.

هر تیم این‌ها را کمی متفاوت می‌نویسد و هر کدام چند باگ کوچک دارد. MassTransit همه این‌ها را یک بار و درست نوشته است.

ایده اصلی: کد تو فقط کسب‌وکار است

کد مافقط کسب‌وکارIConsumer<T>IPublishEndpointrecord OrderPlacedMassTransitکارهای تکراریتلاش دوبارهصف خطاساخت صف‌هاتبدیل پیامOutboxSagaانتقال‌دهندهقابل تعویضRabbitMQAzure Service BusAmazon SQSIn-Memory
کد ما فقط پیام را منتشر می‌کند یا مصرف می‌کند. بقیه کارها در لایه MassTransit است.
  1. پیام یک record ساده است. هیچ وابستگی به RabbitMQ ندارد.
  2. مصرف‌کننده یک کلاس با یک متد است. رابط IConsumer را پیاده می‌کند.
  3. انتقال‌دهنده قابل تعویض است. امروز RabbitMQ، فردا Azure Service Bus. کد مصرف‌کننده تغییر نمی‌کند. برای تست هم یک انتقال‌دهنده داخل حافظه هست.
درباره Kafka: در MassTransit، کافکا یک «Rider» است، نه یک انتقال‌دهنده کامل. یعنی کنار یک انتقال‌دهنده اصلی کار می‌کند و همه امکانات صف را ندارد.

پیام و مصرف‌کننده

پیام‌ها را در یک پروژه مشترک (مثلاً Shop.Contracts) نگه دار. MassTransit نوع پیام را با namespace و نام کلاس می‌شناسد. پس namespace باید در فرستنده و گیرنده یکی باشد.

namespace Shop.Contracts;

public sealed record OrderPlaced(Guid OrderId, string CustomerEmail, decimal Total);

مصرف‌کننده سرویس ایمیل:

using MassTransit;
using Shop.Contracts;

public sealed class SendReceiptConsumer(IEmailSender email) : IConsumer<OrderPlaced>
{
    public async Task Consume(ConsumeContext<OrderPlaced> context)
    {
        var order = context.Message;
        // Idempotent: the sender skips orders that already got a receipt.
        await email.SendReceiptOnceAsync(order.OrderId, order.CustomerEmail, context.CancellationToken);
    }
}

اگر متد Consume بدون خطا تمام شود، MassTransit خودش پیام را Ack می‌کند. اگر خطا بدهد، مسیر تلاش دوباره شروع می‌شود.

راه‌اندازی

using MassTransit;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddMassTransit(x =>
{
    x.SetKebabCaseEndpointNameFormatter();   // queue names like "send-receipt"
    x.AddConsumer<SendReceiptConsumer>();

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("localhost", "/", h =>
        {
            h.Username("guest");
            h.Password("guest");
        });

        cfg.UseMessageRetry(r =>
        {
            r.Intervals(TimeSpan.FromSeconds(1), TimeSpan.FromSeconds(5)); // two retries
            r.Ignore<ArgumentException>();   // a bad message will not get better
        });

        cfg.ConfigureEndpoints(context);     // one queue per consumer, bound to its message types
    });
});

var app = builder.Build();
app.Run();

متد ConfigureEndpoints کارهای زیادی را انجام می‌دهد:

سرویس سفارشPublishShop.Contracts:OrderPlacedیک مبادله‌گر برای هر نوع پیامsend-receiptصف مصرف‌کننده ایمیلreserve-stockصف مصرف‌کننده انبارsend-receipt_errorهمه این‌ها خودکار ساخته می‌شوند
یک Exchange برای نوع پیام، یک صف برای هر مصرف‌کننده، و یک صف خطا کنار هر صف.
  1. برای هر نوع پیام یک Exchange می‌سازد. اسمش از namespace و نام کلاس ساخته می‌شود.
  2. برای هر مصرف‌کننده یک صف می‌سازد. اسم صف از نام کلاس مصرف‌کننده ساخته می‌شود.
  3. صف را به Exchange پیام‌هایش وصل می‌کند. پس هر سرویس جدیدی که OrderPlaced را مصرف کند، خودکار نسخه خودش را می‌گیرد.

انتشار پیام: Publish یا Send؟

  • متد Publish برای رویداد. «سفارش ثبت شد». فرستنده نمی‌داند چه کسی گوش می‌دهد. صفر، یک یا چند مصرف‌کننده می‌گیرند.
  • متد Send برای دستور. «این ایمیل را بفرست». یک مقصد مشخص دارد و فقط یک مصرف‌کننده آن را انجام می‌دهد.
app.MapPost("/orders", async (
    PlaceOrder cmd, ShopDbContext db, IPublishEndpoint publisher, CancellationToken ct) =>
{
    var order = Order.Place(cmd.CustomerEmail, cmd.Items);
    db.Orders.Add(order);

    await publisher.Publish(new OrderPlaced(order.Id, order.CustomerEmail, order.Total), ct);
    await db.SaveChangesAsync(ct); // order and message: one transaction

    return Results.Created($"/orders/{order.Id}", order.Id);
});

تلاش دوباره و صف خطا

پیام رسیدOrderPlacedاجرای اولخطاتلاش دومخطاتلاش سومخطاصف خطا_error1s5sموفقخطای موقت معمولاً در تلاش دوم یا سوم حل می‌شود
فاصله بین تلاش‌ها به سرویس خراب فرصت می‌دهد تا دوباره بالا بیاید.
  1. مصرف‌کننده خطا داد. MassTransit پیام را طبق تنظیم چند بار دیگر به همان مصرف‌کننده می‌دهد.
  2. خطای موقت معمولاً حل می‌شود. مثلاً سرور ایمیل بعد از چند ثانیه جواب می‌دهد.
  3. اگر همه تلاش‌ها شکست خورد، پیام به صفی با پسوند error می‌رود. جزئیات خطا هم در Header های پیام ذخیره می‌شود.
  4. پیام گم نمی‌شود. یک نفر صف خطا را بررسی می‌کند. بعد از رفع مشکل، پیام را به صف اصلی برمی‌گرداند.
خطای دائمی را تکرار نکن: اگر پیام از اول خراب است (مثلاً مبلغ منفی)، تلاش دوباره فقط وقت هدر می‌دهد. این نوع خطاها را با Ignore از تلاش دوباره خارج کن تا مستقیم به صف خطا بروند.

پیام و دیتابیس با هم: Outbox

مشکل Dual Write را در درس کافکا دیدیم: سفارش در دیتابیس ذخیره شد، ولی پیام هیچ وقت نرفت. MassTransit یک Outbox آماده با EF Core دارد:

builder.Services.AddMassTransit(x =>
{
    x.AddEntityFrameworkOutbox<ShopDbContext>(o =>
    {
        o.UseSqlServer();
        o.UseBusOutbox();   // Publish writes to the outbox table, not to the broker
    });

    // ... consumers and UsingRabbitMq as before ...
});

public sealed class ShopDbContext(DbContextOptions<ShopDbContext> options) : DbContext(options)
{
    public DbSet<Order> Orders => Set<Order>();

    protected override void OnModelCreating(ModelBuilder modelBuilder)
    {
        modelBuilder.AddInboxStateEntity();
        modelBuilder.AddOutboxMessageEntity();
        modelBuilder.AddOutboxStateEntity();
    }
}

حالا قدم‌ها این شکلی است:

  1. متد Publish پیام را فقط در جدول Outbox می‌گذارد. هنوز چیزی به RabbitMQ نرفته است.
  2. متد SaveChangesAsync سفارش و پیام را در یک تراکنش ذخیره می‌کند. یا هر دو، یا هیچ کدام.
  3. یک سرویس پس‌زمینه پیام‌های جدول را به RabbitMQ می‌فرستد. اگر RabbitMQ قطع باشد، بعداً می‌فرستد.

همین بسته یک Inbox هم برای سمت مصرف‌کننده دارد. این Inbox شناسه پیام‌های دیده‌شده را نگه می‌دارد و پیام تکراری را کنار می‌گذارد. باید آن را جدا برای صف مصرف‌کننده فعال کنی، با همان جدول‌هایی که بالا ساختیم.

Saga State Machine

برای کارهای طولانی که چند سرویس را درگیر می‌کنند، MassTransit یک ماشین وضعیت دارد. وضعیت هر سفارش در دیتابیس ذخیره می‌شود و هر پیام آن را یک قدم جلو می‌برد. این موضوع در درس Saga با جزئیات و کد کامل آمده است.

لایسنس را چک کن: نسخه ۸ MassTransit متن‌باز است. نسخه ۹ با لایسنس تجاری منتشر شده است. قبل از شروع یک پروژه جدید، شرایط لایسنس نسخه‌ای را که انتخاب می‌کنی بخوان.

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

اشتباه نتیجه راه درست
namespace متفاوت برای یک پیام در دو سرویس پیام فرستاده می‌شود ولی هیچ مصرف‌کننده‌ای آن را نمی‌گیرد. پیام‌ها در یک پروژه یا بسته مشترک.
Publish بدون Outbox بعد از ذخیره دیتابیس سفارش ذخیره شد ولی پیام گم شد. Outbox با EF Core.
تلاش دوباره برای خطای دائمی پیام خراب چند بار بی‌دلیل اجرا می‌شود. خطاهای دائمی را Ignore کن.
صف خطا را کسی نگاه نمی‌کند سفارش‌ها بی‌صدا گیر می‌کنند. هشدار روی تعداد پیام صف خطا.
Send برای رویداد مصرف‌کننده‌های جدید پیام را نمی‌گیرند. رویداد با Publish، دستور با Send.
فکر کنی MassTransit پیام تکراری را کاملاً حذف می‌کند بدون Inbox، مشتری دو ایمیل می‌گیرد. Inbox یا مصرف‌کننده Idempotent.

چه وقت MassTransit؟

مناسب

  • چند سرویس .NET با پیام با هم حرف می‌زنند.
  • تلاش دوباره، صف خطا و Outbox لازم داری و نمی‌خواهی خودت بنویسی.
  • کار طولانی با Saga لازم داری.
  • می‌خواهی بعداً انتقال‌دهنده را عوض کنی.

نامناسب

  • فقط یک پیام ساده بین دو سرویس داری. کتابخانه رسمی RabbitMQ کافی است.
  • سیستم تو اصلاً .NET نیست، یا بیشتر سرویس‌ها به زبان دیگری هستند.
  • کار اصلی تو جریان رویداد سنگین روی کافکا است. کتابخانه Confluent.Kafka مستقیم‌تر است.

خلاصه در شش خط

  1. کتابخانه MassTransit کارهای تکراری پیام‌رسانی را برای تو انجام می‌دهد.
  2. پیام یک record ساده است و مصرف‌کننده یک کلاس با متد Consume.
  3. متد ConfigureEndpoints صف‌ها و Exchange ها را خودش می‌سازد.
  4. رویداد را Publish کن و دستور را Send.
  5. تلاش دوباره با فاصله، و بعد صف خطا. خطای دائمی را تکرار نکن.
  6. با Outbox و Inbox، پیام نه گم می‌شود و نه دو بار اثر می‌گذارد.