الگوی Outbox و Inbox
تغییر دیتابیس و فرستادن پیام دو کار جدا هستند و ممکن است یکی انجام شود و دیگری نه. با Outbox، پیام را همراه با داده در یک تراکنش ذخیره کن و بعداً بفرست. با Inbox، شناسه پیامهای دریافتی را ذخیره کن تا پیام تکراری دو بار اجرا نشود.
نویسنده: bezzad
مشکل: دو نوشتن، بدون تراکنش مشترک
مشتری سفارش ثبت میکند. سرویس سفارش دو کار انجام میدهد:
- سفارش را در دیتابیس خودش ذخیره میکند.
- پیام «سفارش ثبت شد» (OrderPlaced) را به Kafka یا RabbitMQ میفرستد. سرویس انبار با این پیام کالا را رزرو میکند.
کد ساده اینطور است:
db.Orders.Add(order);
await db.SaveChangesAsync(ct);
await bus.PublishAsync(new OrderPlaced(order.Id, order.Lines), ct);
ظاهراً درست است. ولی دیتابیس و Broker دو سیستم جدا هستند. هیچ تراکنشی هر دو را با هم پوشش نمیدهد. به این مشکل Dual Write میگویند.
چه چیزی ممکن است خراب شود؟
- برنامه بین دو خط میافتد. مثلاً هنگام انتشار نسخه جدید، Pod بسته میشود. سفارش هست، پیام نیست. انبار هیچ وقت کالا را رزرو نمیکند.
- سرویس Broker چند ثانیه در دسترس نیست. فرستادن پیام خطا میدهد. سفارش قبلاً ذخیره شده و برنمیگردد.
- ترتیب را عوض کنیم، مشکل برعکس میشود. اگر اول پیام را بفرستیم و بعد ذخیره شکست بخورد، انبار کالای سفارشی را رزرو میکند که اصلاً وجود ندارد.
ایده Outbox: پیام را هم در دیتابیس بنویس
دیتابیس فقط یک چیز را تضمین میکند: همه نوشتنهای داخل یک تراکنش با هم انجام میشوند. پس پیام را هم در همان دیتابیس مینویسیم:
- در یک تراکنش، هم سفارش را ذخیره میکنیم و هم یک ردیف در جدول OutboxMessages. این ردیف متن پیام است.
- یک کار پسزمینه (Relay) ردیفهای ارسالنشده را میخواند و به Broker میفرستد.
- بعد از ارسال موفق، روی ردیف علامت «ارسال شد» میزند.
چرا این کار مشکل را حل میکند؟
- اگر تراکنش موفق شد، پیام حتماً در جدول هست.
- اگر برنامه افتاد، بعد از بالا آمدن، کار پسزمینه همان ردیف را پیدا میکند و میفرستد.
- اگر Broker چند دقیقه قطع بود، پیامها در جدول صبر میکنند. هیچ پیامی گم نمیشود.
ایده Inbox: پیام تکراری را بشناس
پیام تکراری فقط از Outbox نمیآید. Kafka و RabbitMQ هم ممکن است یک پیام را دوباره تحویل دهند. مثلاً مصرفکننده پیام را پردازش کرده، ولی قبل از تأیید (Ack یا Commit کردن Offset) افتاده است.
الگوی Inbox آینه Outbox در طرف گیرنده است:
- هر پیام یک شناسه یکتا (MessageId) دارد.
- سرویس انبار جدولی به اسم InboxMessages دارد. کلید اصلی آن، شناسه پیام است.
- در یک تراکنش، هم شناسه پیام را در Inbox ثبت میکند و هم کالا را رزرو میکند.
- اگر همین پیام دوباره برسد، ثبت شناسه با خطای کلید تکراری شکست میخورد. کل تراکنش برمیگردد و کالا دوباره رزرو نمیشود.
- فقط بعد از ذخیره موفق، پیام را تأیید (Ack) میکنیم.
مثال زنده
گزینهها را روشن یا خاموش کن و هر سناریو را اجرا کن. ببین در هر حالت چند رزرو در انبار ثبت میشود:
کد
ذخیره سفارش و پیام در یک تراکنش
جدول Outbox ساده است: شناسه پیام، نوع، متن و زمان ارسال.
public sealed class OutboxMessage
{
public long Id { get; init; } // keeps the order of messages
public Guid MessageId { get; init; } = Guid.NewGuid(); // consumers use it to find duplicates
public required string Type { get; init; }
public required string Payload { get; init; }
public DateTimeOffset? SentAt { get; set; }
}
در EF Core، یک بار صدا زدن متد SaveChangesAsync خودش یک تراکنش است. پس هر دو ردیف با هم ذخیره میشوند:
app.MapPost("/orders", async (PlaceOrder cmd, ShopDbContext db, CancellationToken ct) =>
{
var order = Order.Place(cmd.CustomerId, cmd.Lines);
db.Orders.Add(order);
db.OutboxMessages.Add(new OutboxMessage
{
Type = nameof(OrderPlaced),
Payload = JsonSerializer.Serialize(new OrderPlaced(order.Id, order.Lines))
});
await db.SaveChangesAsync(ct); // one transaction: order + message
return Results.Created($"/orders/{order.Id}", order.Id);
});
کار پسزمینه که پیامها را میفرستد
public sealed class OutboxRelay(
IServiceScopeFactory scopes, IMessagePublisher publisher, ILogger<OutboxRelay> logger)
: BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken ct)
{
using var timer = new PeriodicTimer(TimeSpan.FromSeconds(1));
while (await timer.WaitForNextTickAsync(ct))
{
try { await SendBatchAsync(ct); }
catch (Exception ex) when (!ct.IsCancellationRequested)
{
logger.LogError(ex, "Outbox batch failed. Will try again.");
}
}
}
private async Task SendBatchAsync(CancellationToken ct)
{
await using var scope = scopes.CreateAsyncScope();
var db = scope.ServiceProvider.GetRequiredService<ShopDbContext>();
var batch = await db.OutboxMessages
.Where(m => m.SentAt == null)
.OrderBy(m => m.Id)
.Take(100)
.ToListAsync(ct);
foreach (var message in batch)
{
// IMessagePublisher is our own small wrapper around Kafka or RabbitMQ.
await publisher.PublishAsync(message.MessageId, message.Type, message.Payload, ct);
message.SentAt = DateTimeOffset.UtcNow;
await db.SaveChangesAsync(ct);
}
}
}
اگر برنامه بین فرستادن و ذخیره علامت بیفتد، پیام دوباره فرستاده میشود. این طبیعی است. Inbox در طرف گیرنده آن را میگیرد.
مصرفکننده با Inbox
public sealed class OrderPlacedHandler(InventoryDbContext db)
{
public async Task HandleAsync(Guid messageId, OrderPlaced message, CancellationToken ct)
{
db.InboxMessages.Add(new InboxMessage(messageId, DateTimeOffset.UtcNow));
db.Reservations.Add(Reservation.For(message.OrderId, message.Lines));
try
{
await db.SaveChangesAsync(ct); // inbox row + reservation, together
}
catch (DbUpdateException ex) when (IsDuplicateKey(ex))
{
// Same MessageId was processed before: do nothing.
}
// Ack the message only after this method returns.
}
private static bool IsDuplicateKey(DbUpdateException ex) =>
ex.InnerException is PostgresException { SqlState: PostgresErrorCodes.UniqueViolation };
}
تشخیص خطای کلید تکراری به نوع دیتابیس بستگی دارد. اینجا PostgreSQL با کتابخانه Npgsql است.
آماده: MassTransit
لازم نیست همه اینها را خودت بنویسی. کتابخانه MassTransit یک Outbox آماده برای EF Core دارد. همان جدولها و کار پسزمینه را میسازد و در طرف گیرنده هم پیام تکراری را تشخیص میدهد:
builder.Services.AddMassTransit(x =>
{
x.AddEntityFrameworkOutbox<ShopDbContext>(o =>
{
o.UsePostgres();
o.UseBusOutbox(); // published messages go to the outbox table first
});
x.UsingRabbitMq((context, cfg) => cfg.ConfigureEndpoints(context));
});
قانونهای مهم
- پیام و داده در یک تراکنش. اگر ردیف Outbox در تراکنش دیگری ذخیره شود، همان مشکل Dual Write برمیگردد.
- انتظار پیام تکراری را داشته باش. Outbox تحویل حداقل یک بار میدهد. پس هر مصرفکننده باید Idempotent باشد.
- هر پیام یک شناسه یکتا داشته باشد. این شناسه در همه تلاشهای ارسال ثابت است. گیرنده با آن تکراریها را میشناسد.
- تأیید پیام بعد از ذخیره. اگر قبل از پردازش Ack کنی و برنامه بیفتد، پیام گم میشود.
- روی چند Pod مراقب باش. اگر چند نسخه از کار پسزمینه یک ردیف را با هم بخوانند، پیام چند بار میرود. یا فقط یک نسخه را اجرا کن، یا ردیفها را قفل کن (مثلاً با دستور FOR UPDATE SKIP LOCKED در PostgreSQL). اگر ترتیب پیامها مهم است، این انتخاب را با دقت انجام بده.
- جدولها را تمیز کن. ردیفهای قدیمی Outbox و Inbox را بعد از مدتی پاک کن. جدول Inbox باید حداقل به اندازه بیشترین زمان ممکن برای رسیدن پیام تکراری نگه داشته شود.
- کار بیرون از دیتابیس را جدا ببین. Inbox فقط تغییرات همان دیتابیس را دقیقاً یک بار میکند. اگر پیامک بفرستی و قبل از ذخیره بیفتی، پیامک دوباره میرود. برای این کارها، یک شناسه یکتا به سرویس بیرونی بده، اگر پشتیبانی میکند.
اشتباههای رایج
| اشتباه | نتیجه | راه درست |
|---|---|---|
| ذخیره در دیتابیس و بعد فرستادن مستقیم پیام | با افتادن برنامه، پیام گم میشود. | الگوی Outbox. |
| فرستادن پیام داخل تراکنش، قبل از ثبت تراکنش | پیام میرود، ولی تراکنش برمیگردد. | فرستادن فقط از جدول Outbox. |
| فکر کردن به «دقیقاً یک بار» | پیام تکراری کار را دو بار انجام میدهد. | مصرفکننده Idempotent یا Inbox. |
| اول چک کردن Inbox، بعد کار، بعد ثبت | دو پیام همزمان هر دو از چک رد میشوند. | ثبت شناسه و کار در یک تراکنش با کلید یکتا. |
| تأیید پیام قبل از پردازش | اگر برنامه بیفتد، پیام گم میشود. | تأیید بعد از ذخیره موفق. |
| چند Relay بدون قفل | هر پیام چند بار فرستاده میشود. | یک نسخه، یا قفل ردیف. |
| پاک نکردن جدولها | جدول بزرگ و کوئری کند میشود. | پاک کردن دورهای ردیفهای قدیمی. |
چه وقت Outbox و Inbox؟
مناسب
- تغییر داده و پیام باید با هم هماهنگ بمانند.
- گم شدن پیام هزینه دارد: سفارش بدون رزرو، پرداخت بدون رسید.
- چند سرویس با پیام و Saga با هم کار میکنند.
لازم نیست
- پیام فقط یک اعلان کماهمیت است و گم شدنش مهم نیست.
- اصلاً دیتابیسی تغییر نمیکند و فقط یک پیام فرستاده میشود.
- همه چیز در یک برنامه و یک دیتابیس است. یک تراکنش معمولی کافی است.
خلاصه در شش خط
- نوشتن در دیتابیس و فرستادن پیام، تراکنش مشترک ندارند (Dual Write).
- با Outbox، پیام همراه با داده در یک تراکنش ذخیره میشود.
- یک کار پسزمینه پیامها را از جدول برمیدارد و میفرستد.
- نتیجه، تحویل حداقل یک بار است. پس پیام تکراری داریم.
- با Inbox، گیرنده شناسه پیام را همراه با کارش در یک تراکنش ثبت میکند تا تکراریها اثر نداشته باشند.
- پیام را بعد از ذخیره تأیید کن، جدولها را تمیز کن و مراقب چند Pod باش.