MassTransit
کتابخانه MassTransit کارهای تکراری پیامرسانی را در .NET انجام میدهد. تو فقط پیام و مصرفکننده را مینویسی. ساخت صفها، تلاش دوباره، صف خطا، Outbox و Saga را MassTransit آماده دارد.
نویسنده: bezzad
مشکل: کد تکراری دور هر پیام
در درس RabbitMQ دیدیم که برای یک پیام ساده چقدر کد لازم است. در فروشگاه اینترنتی ما دهها نوع پیام داریم. برای هر کدام باید این کارها را بنویسیم:
- ساخت Exchange، صف و Binding.
- تبدیل شیء به JSON و برعکس.
- تلاش دوباره وقتی سرور ایمیل یک لحظه جواب نمیدهد.
- فرستادن پیام خراب به صف خطا.
- ثبت پیام در Outbox تا با دیتابیس هماهنگ باشد.
- ساختن Scope برای Dependency Injection برای هر پیام.
هر تیم اینها را کمی متفاوت مینویسد و هر کدام چند باگ کوچک دارد. MassTransit همه اینها را یک بار و درست نوشته است.
ایده اصلی: کد تو فقط کسبوکار است
- پیام یک record ساده است. هیچ وابستگی به RabbitMQ ندارد.
- مصرفکننده یک کلاس با یک متد است. رابط IConsumer را پیاده میکند.
- انتقالدهنده قابل تعویض است. امروز RabbitMQ، فردا Azure Service Bus. کد مصرفکننده تغییر نمیکند. برای تست هم یک انتقالدهنده داخل حافظه هست.
پیام و مصرفکننده
پیامها را در یک پروژه مشترک (مثلاً 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 کارهای زیادی را انجام میدهد:
- برای هر نوع پیام یک Exchange میسازد. اسمش از namespace و نام کلاس ساخته میشود.
- برای هر مصرفکننده یک صف میسازد. اسم صف از نام کلاس مصرفکننده ساخته میشود.
- صف را به 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);
});
تلاش دوباره و صف خطا
- مصرفکننده خطا داد. MassTransit پیام را طبق تنظیم چند بار دیگر به همان مصرفکننده میدهد.
- خطای موقت معمولاً حل میشود. مثلاً سرور ایمیل بعد از چند ثانیه جواب میدهد.
- اگر همه تلاشها شکست خورد، پیام به صفی با پسوند error میرود. جزئیات خطا هم در Header های پیام ذخیره میشود.
- پیام گم نمیشود. یک نفر صف خطا را بررسی میکند. بعد از رفع مشکل، پیام را به صف اصلی برمیگرداند.
پیام و دیتابیس با هم: 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();
}
}
حالا قدمها این شکلی است:
- متد Publish پیام را فقط در جدول Outbox میگذارد. هنوز چیزی به RabbitMQ نرفته است.
- متد SaveChangesAsync سفارش و پیام را در یک تراکنش ذخیره میکند. یا هر دو، یا هیچ کدام.
- یک سرویس پسزمینه پیامهای جدول را به RabbitMQ میفرستد. اگر RabbitMQ قطع باشد، بعداً میفرستد.
همین بسته یک Inbox هم برای سمت مصرفکننده دارد. این Inbox شناسه پیامهای دیدهشده را نگه میدارد و پیام تکراری را کنار میگذارد. باید آن را جدا برای صف مصرفکننده فعال کنی، با همان جدولهایی که بالا ساختیم.
Saga State Machine
برای کارهای طولانی که چند سرویس را درگیر میکنند، MassTransit یک ماشین وضعیت دارد. وضعیت هر سفارش در دیتابیس ذخیره میشود و هر پیام آن را یک قدم جلو میبرد. این موضوع در درس Saga با جزئیات و کد کامل آمده است.
اشتباههای رایج
| اشتباه | نتیجه | راه درست |
|---|---|---|
| namespace متفاوت برای یک پیام در دو سرویس | پیام فرستاده میشود ولی هیچ مصرفکنندهای آن را نمیگیرد. | پیامها در یک پروژه یا بسته مشترک. |
| Publish بدون Outbox بعد از ذخیره دیتابیس | سفارش ذخیره شد ولی پیام گم شد. | Outbox با EF Core. |
| تلاش دوباره برای خطای دائمی | پیام خراب چند بار بیدلیل اجرا میشود. | خطاهای دائمی را Ignore کن. |
| صف خطا را کسی نگاه نمیکند | سفارشها بیصدا گیر میکنند. | هشدار روی تعداد پیام صف خطا. |
| Send برای رویداد | مصرفکنندههای جدید پیام را نمیگیرند. | رویداد با Publish، دستور با Send. |
| فکر کنی MassTransit پیام تکراری را کاملاً حذف میکند | بدون Inbox، مشتری دو ایمیل میگیرد. | Inbox یا مصرفکننده Idempotent. |
چه وقت MassTransit؟
مناسب
- چند سرویس .NET با پیام با هم حرف میزنند.
- تلاش دوباره، صف خطا و Outbox لازم داری و نمیخواهی خودت بنویسی.
- کار طولانی با Saga لازم داری.
- میخواهی بعداً انتقالدهنده را عوض کنی.
نامناسب
- فقط یک پیام ساده بین دو سرویس داری. کتابخانه رسمی RabbitMQ کافی است.
- سیستم تو اصلاً .NET نیست، یا بیشتر سرویسها به زبان دیگری هستند.
- کار اصلی تو جریان رویداد سنگین روی کافکا است. کتابخانه Confluent.Kafka مستقیمتر است.
خلاصه در شش خط
- کتابخانه MassTransit کارهای تکراری پیامرسانی را برای تو انجام میدهد.
- پیام یک record ساده است و مصرفکننده یک کلاس با متد Consume.
- متد ConfigureEndpoints صفها و Exchange ها را خودش میسازد.
- رویداد را Publish کن و دستور را Send.
- تلاش دوباره با فاصله، و بعد صف خطا. خطای دائمی را تکرار نکن.
- با Outbox و Inbox، پیام نه گم میشود و نه دو بار اثر میگذارد.