Levelwise
فارسی
سیستم‌های توزیع‌شده

الگوی Outbox و Inbox

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

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

نویسنده: bezzad

مشکل: دو نوشتن، بدون تراکنش مشترک

مشتری سفارش ثبت می‌کند. سرویس سفارش دو کار انجام می‌دهد:

  1. سفارش را در دیتابیس خودش ذخیره می‌کند.
  2. پیام «سفارش ثبت شد» (OrderPlaced) را به Kafka یا RabbitMQ می‌فرستد. سرویس انبار با این پیام کالا را رزرو می‌کند.

کد ساده این‌طور است:

db.Orders.Add(order);
await db.SaveChangesAsync(ct);
await bus.PublishAsync(new OrderPlaced(order.Id, order.Lines), ct);

ظاهراً درست است. ولی دیتابیس و Broker دو سیستم جدا هستند. هیچ تراکنشی هر دو را با هم پوشش نمی‌دهد. به این مشکل Dual Write می‌گویند.

سرویس سفارشدیتابیس۱. ذخیره سفارشسفارش ثبت شدKafka۲. فرستادن پیامپیامی نرسیدبرنامه این‌جا افتادبین دو کار، هیچ تراکنش مشترکی نیست
اگر برنامه درست بین دو خط بیفتد، سفارش ذخیره شده ولی پیامش هیچ وقت نمی‌رود.

چه چیزی ممکن است خراب شود؟

  1. برنامه بین دو خط می‌افتد. مثلاً هنگام انتشار نسخه جدید، Pod بسته می‌شود. سفارش هست، پیام نیست. انبار هیچ وقت کالا را رزرو نمی‌کند.
  2. سرویس Broker چند ثانیه در دسترس نیست. فرستادن پیام خطا می‌دهد. سفارش قبلاً ذخیره شده و برنمی‌گردد.
  3. ترتیب را عوض کنیم، مشکل برعکس می‌شود. اگر اول پیام را بفرستیم و بعد ذخیره شکست بخورد، انبار کالای سفارشی را رزرو می‌کند که اصلاً وجود ندارد.
خطای پنهان: این مشکل در محیط توسعه تقریباً هیچ وقت دیده نمی‌شود. فقط در Production و هنگام انتشار یا قطعی‌های کوتاه پیش می‌آید. برای همین سخت پیدا می‌شود.

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

دیتابیس فقط یک چیز را تضمین می‌کند: همه نوشتن‌های داخل یک تراکنش با هم انجام می‌شوند. پس پیام را هم در همان دیتابیس می‌نویسیم:

  1. در یک تراکنش، هم سفارش را ذخیره می‌کنیم و هم یک ردیف در جدول OutboxMessages. این ردیف متن پیام است.
  2. یک کار پس‌زمینه (Relay) ردیف‌های ارسال‌نشده را می‌خواند و به Broker می‌فرستد.
  3. بعد از ارسال موفق، روی ردیف علامت «ارسال شد» می‌زند.
سرویس سفارشیک تراکنش دیتابیسOrdersOutboxMessagesمتن پیام و زمان ارسالیا هر دو ذخیره می‌شوند، یا هیچ‌کدامکار پس‌زمینهRelayردیف‌های ارسال‌نشده را می‌خواندBrokerمی‌فرستد، بعد علامت می‌زند
دیتابیس تضمین می‌کند سفارش و پیام با هم ذخیره شوند. فرستادن بعداً و جدا انجام می‌شود.

چرا این کار مشکل را حل می‌کند؟

  1. اگر تراکنش موفق شد، پیام حتماً در جدول هست.
  2. اگر برنامه افتاد، بعد از بالا آمدن، کار پس‌زمینه همان ردیف را پیدا می‌کند و می‌فرستد.
  3. اگر Broker چند دقیقه قطع بود، پیام‌ها در جدول صبر می‌کنند. هیچ پیامی گم نمی‌شود.
ولی پیام تکراری داریم: کار پس‌زمینه ممکن است پیام را بفرستد و درست قبل از علامت زدن بیفتد. دفعه بعد همان پیام را دوباره می‌فرستد. پس Outbox تحویل «حداقل یک بار» (at-least-once) می‌دهد، نه «دقیقاً یک بار». طرف گیرنده باید پیام تکراری را تحمل کند.

ایده Inbox: پیام تکراری را بشناس

پیام تکراری فقط از Outbox نمی‌آید. Kafka و RabbitMQ هم ممکن است یک پیام را دوباره تحویل دهند. مثلاً مصرف‌کننده پیام را پردازش کرده، ولی قبل از تأیید (Ack یا Commit کردن Offset) افتاده است.

الگوی Inbox آینه Outbox در طرف گیرنده است:

  1. هر پیام یک شناسه یکتا (MessageId) دارد.
  2. سرویس انبار جدولی به اسم InboxMessages دارد. کلید اصلی آن، شناسه پیام است.
  3. در یک تراکنش، هم شناسه پیام را در Inbox ثبت می‌کند و هم کالا را رزرو می‌کند.
  4. اگر همین پیام دوباره برسد، ثبت شناسه با خطای کلید تکراری شکست می‌خورد. کل تراکنش برمی‌گردد و کالا دوباره رزرو نمی‌شود.
  5. فقط بعد از ذخیره موفق، پیام را تأیید (Ack) می‌کنیم.
Brokermsg 42msg 42دوباره رسیدسرویس انبار، یک تراکنشInboxMessagesReservationsکلید اصلی: شناسه پیامرزرو انجام شدتکراری، رد شدثبت شناسه برای بار دوم شکست می‌خورد، پس کار دوباره انجام نمی‌شود
دیتابیس با کلید یکتا تضمین می‌کند هر پیام فقط یک بار اثر بگذارد، حتی روی چند Pod.
نکته: الگوی Inbox یک روش برای مصرف‌کننده Idempotent است. اگر کار ذاتاً تکرارپذیر باشد، مثلاً «وضعیت سفارش را پرداخت‌شده کن»، شاید جدول Inbox لازم نباشد. درس Idempotency این را کامل‌تر توضیح می‌دهد.

مثال زنده

گزینه‌ها را روشن یا خاموش کن و هر سناریو را اجرا کن. ببین در هر حالت چند رزرو در انبار ثبت می‌شود:

مثال زنده: یک سفارش، از دیتابیس تا انبار
۰سفارش در دیتابیس
-ردیف Outbox
۰پیام در Broker
۰رزرو در انبار
گزینه‌ها را انتخاب کن و یک سناریو را اجرا کن.

    کد

    ذخیره سفارش و پیام در یک تراکنش

    جدول 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));
    });
    یک راه دیگر: به جای کار پس‌زمینه، می‌شود از CDC (Change Data Capture) استفاده کرد. ابزاری مثل Debezium لاگ تغییرات دیتابیس را می‌خواند و ردیف‌های جدید Outbox را به Kafka می‌فرستد. کد کمتری می‌خواهد، ولی یک زیرساخت دیگر به سیستم اضافه می‌کند.

    قانون‌های مهم

    1. پیام و داده در یک تراکنش. اگر ردیف Outbox در تراکنش دیگری ذخیره شود، همان مشکل Dual Write برمی‌گردد.
    2. انتظار پیام تکراری را داشته باش. Outbox تحویل حداقل یک بار می‌دهد. پس هر مصرف‌کننده باید Idempotent باشد.
    3. هر پیام یک شناسه یکتا داشته باشد. این شناسه در همه تلاش‌های ارسال ثابت است. گیرنده با آن تکراری‌ها را می‌شناسد.
    4. تأیید پیام بعد از ذخیره. اگر قبل از پردازش Ack کنی و برنامه بیفتد، پیام گم می‌شود.
    5. روی چند Pod مراقب باش. اگر چند نسخه از کار پس‌زمینه یک ردیف را با هم بخوانند، پیام چند بار می‌رود. یا فقط یک نسخه را اجرا کن، یا ردیف‌ها را قفل کن (مثلاً با دستور FOR UPDATE SKIP LOCKED در PostgreSQL). اگر ترتیب پیام‌ها مهم است، این انتخاب را با دقت انجام بده.
    6. جدول‌ها را تمیز کن. ردیف‌های قدیمی Outbox و Inbox را بعد از مدتی پاک کن. جدول Inbox باید حداقل به اندازه بیشترین زمان ممکن برای رسیدن پیام تکراری نگه داشته شود.
    7. کار بیرون از دیتابیس را جدا ببین. Inbox فقط تغییرات همان دیتابیس را دقیقاً یک بار می‌کند. اگر پیامک بفرستی و قبل از ذخیره بیفتی، پیامک دوباره می‌رود. برای این کارها، یک شناسه یکتا به سرویس بیرونی بده، اگر پشتیبانی می‌کند.

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

    اشتباه نتیجه راه درست
    ذخیره در دیتابیس و بعد فرستادن مستقیم پیام با افتادن برنامه، پیام گم می‌شود. الگوی Outbox.
    فرستادن پیام داخل تراکنش، قبل از ثبت تراکنش پیام می‌رود، ولی تراکنش برمی‌گردد. فرستادن فقط از جدول Outbox.
    فکر کردن به «دقیقاً یک بار» پیام تکراری کار را دو بار انجام می‌دهد. مصرف‌کننده Idempotent یا Inbox.
    اول چک کردن Inbox، بعد کار، بعد ثبت دو پیام همزمان هر دو از چک رد می‌شوند. ثبت شناسه و کار در یک تراکنش با کلید یکتا.
    تأیید پیام قبل از پردازش اگر برنامه بیفتد، پیام گم می‌شود. تأیید بعد از ذخیره موفق.
    چند Relay بدون قفل هر پیام چند بار فرستاده می‌شود. یک نسخه، یا قفل ردیف.
    پاک نکردن جدول‌ها جدول بزرگ و کوئری کند می‌شود. پاک کردن دوره‌ای ردیف‌های قدیمی.

    چه وقت Outbox و Inbox؟

    مناسب

    • تغییر داده و پیام باید با هم هماهنگ بمانند.
    • گم شدن پیام هزینه دارد: سفارش بدون رزرو، پرداخت بدون رسید.
    • چند سرویس با پیام و Saga با هم کار می‌کنند.

    لازم نیست

    • پیام فقط یک اعلان کم‌اهمیت است و گم شدنش مهم نیست.
    • اصلاً دیتابیسی تغییر نمی‌کند و فقط یک پیام فرستاده می‌شود.
    • همه چیز در یک برنامه و یک دیتابیس است. یک تراکنش معمولی کافی است.

    خلاصه در شش خط

    1. نوشتن در دیتابیس و فرستادن پیام، تراکنش مشترک ندارند (Dual Write).
    2. با Outbox، پیام همراه با داده در یک تراکنش ذخیره می‌شود.
    3. یک کار پس‌زمینه پیام‌ها را از جدول برمی‌دارد و می‌فرستد.
    4. نتیجه، تحویل حداقل یک بار است. پس پیام تکراری داریم.
    5. با Inbox، گیرنده شناسه پیام را همراه با کارش در یک تراکنش ثبت می‌کند تا تکراری‌ها اثر نداشته باشند.
    6. پیام را بعد از ذخیره تأیید کن، جدول‌ها را تمیز کن و مراقب چند Pod باش.