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

Kafka

کافکا یک دفتر ثبت بزرگ برای پیام‌ها است. پیام‌ها بعد از خواندن پاک نمی‌شوند و هر خواننده جای خودش را نگه می‌دارد. ترتیب فقط داخل یک Partition حفظ می‌شود و هر Partition در هر گروه فقط یک خواننده دارد.

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

نویسنده: bezzad

مشکل: یک سفارش، چند سرویس

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

  • سرویس ایمیل باید رسید خرید را بفرستد.
  • سرویس انبار باید کالا را آماده کند.
  • سرویس گزارش باید فروش امروز را به‌روز کند.

راه ساده این است که سرویس سفارش هر سه را مستقیم صدا بزند. ولی این راه سه مشکل دارد:

  1. وابستگی زیاد. اگر سرویس ایمیل خراب باشد، ثبت سفارش هم کند یا خراب می‌شود.
  2. تغییر سخت. برای هر سرویس جدید، کد سرویس سفارش باید تغییر کند.
  3. بار ناگهانی. در روز حراج، همه سرویس‌ها باید همان لحظه همه درخواست‌ها را جواب دهند.

راه بهتر این است: سرویس سفارش فقط یک خبر ثبت می‌کند، «سفارش ثبت شد». هر سرویسی که لازم دارد، هر وقت آماده بود آن را می‌خواند. کافکا جای نگه‌داشتن همین خبرها است.

ایده اصلی: یک دفتر ثبت، نه یک صف

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

  1. پیام بعد از خواندن پاک نمی‌شود. پیام تا مدت مشخصی (Retention) می‌ماند. پیش‌فرض این مدت یک هفته است و قابل تغییر است.
  2. هر خواننده جای خودش را دارد. هر خواننده یک عدد نگه می‌دارد: «تا خط چندم را خوانده‌ام». به این عدد Offset می‌گوییم.
  3. پس خواننده‌ها کار هم را خراب نمی‌کنند. سرویس ایمیل و سرویس گزارش هر دو همه پیام‌ها را می‌خوانند، هر کدام با سرعت خودش.
  4. می‌شود دوباره خواند. یک سرویس جدید می‌تواند از اول دفتر شروع کند و همه خبرهای قدیمی را ببیند.
فرق با صف معمولی: در یک صف معمولی، مثل صف RabbitMQ، پیام بعد از تأیید خواننده حذف می‌شود. در کافکا پیام می‌ماند و فقط Offset خواننده جلو می‌رود.

Topic، Partition و Offset

هر نوع خبر یک Topic دارد. مثلاً همه خبرهای سفارش در topic به نام orders می‌روند.

یک دفتر تنها برای حجم زیاد کند است. پس هر topic به چند بخش تقسیم می‌شود. به هر بخش Partition می‌گوییم. هر Partition یک دفتر جدا است که روی یک سرور نگه داشته می‌شود (و چند کپی روی سرورهای دیگر دارد).

سرویس سفارشProducertopic: ordersPartition 0o-3o-6o-3o-90123Partition 1o-1o-4o-1012Partition 2o-2o-5o-2o-8o-501234پیام جدید همیشه به آخر ردیف اضافه می‌شودعدد زیر هر خانه، جای آن پیام است
سفارش o-3 دو پیام دارد و هر دو در Partition 0 هستند. کلید یکسان، Partition یکسان.

کافکا از کجا می‌داند هر پیام به کدام Partition برود؟

  1. هر پیام یک کلید (Key) دارد. ما در فروشگاه، شناسه سفارش را کلید می‌گذاریم.
  2. کافکا از کلید یک عدد می‌سازد (Hash). باقی‌مانده این عدد بر تعداد Partition ها، شماره Partition است.
  3. پس کلید یکسان، همیشه Partition یکسان. همه پیام‌های سفارش o-3 در یک Partition و به همان ترتیب ارسال قرار می‌گیرند.
  4. اگر کلید نگذاریم، پیام‌ها بین Partition ها پخش می‌شوند و هیچ ترتیبی بین آن‌ها حفظ نمی‌شود.
قانون مهم ترتیب: کافکا ترتیب را فقط داخل یک Partition تضمین می‌کند. بین دو Partition هیچ ترتیبی نیست. اگر «ثبت سفارش» و «لغو سفارش» با دو کلید متفاوت بروند، ممکن است «لغو» زودتر خوانده شود.

چند خواننده برای یک کار: Consumer Group

سرویس ایمیل روی چند pod اجرا می‌شود. نمی‌خواهیم هر pod به هر مشتری یک ایمیل جدا بفرستد. پس همه pod های ایمیل را در یک Consumer Group می‌گذاریم.

قانون‌های گروه ساده‌اند:

  1. هر گروه همه پیام‌های topic را می‌خواند. گروه ایمیل و گروه انبار هر دو همه سفارش‌ها را می‌بینند.
  2. داخل یک گروه، هر Partition فقط یک خواننده دارد. پس هیچ پیامی در یک گروه همزمان دو بار پردازش نمی‌شود.
  3. یک خواننده می‌تواند چند Partition داشته باشد.
  4. هر گروه Offset خودش را دارد. اگر گروه انبار عقب باشد، گروه ایمیل منتظر آن نمی‌ماند.
topic: ordersPartition 0Partition 1Partition 2Partition 3گروه ایمیلgroup: emailConsumer AConsumer Bگروه انبارgroup: warehouseConsumer C
گروه ایمیل دو خواننده دارد و Partition ها بین آن‌ها تقسیم شده‌اند. گروه انبار یک خواننده دارد که همه را می‌خواند.

نتیجه مهم قانون دوم: بیشترین تعداد خواننده مفید در یک گروه، برابر تعداد Partition ها است. اگر شش Partition و ده pod داشته باشی، چهار pod بیکار می‌مانند.

مثال زنده

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

مثال زنده: شش Partition و یک Consumer Group

کلید پیام در دکمه اول شناسه سفارش است. در دکمه دوم کلید شناسه مشتری است.

بخش‌های topic
مصرف‌کننده‌های گروه

    دو چیز را در این مثال دیدی:

    1. وقتی تعداد Consumer عوض می‌شود، گروه دوباره تقسیم می‌شود. به این کار Rebalance می‌گوییم. در این مدت خواندن پیام کند یا متوقف می‌شود.
    2. وقتی یک کلید خیلی پرکار است، یک Partition داغ می‌شود. همه پیام‌های مشتری بزرگ روی یک Partition هستند. پس فقط یک Consumer روی آن‌ها کار می‌کند. Consumer بیشتر هیچ کمکی نمی‌کند. به این مشکل Hot Partition می‌گوییم.

    ارسال پیام (Producer)

    در .NET معمولاً از کتابخانه Confluent.Kafka استفاده می‌کنیم. کلید پیام، شناسه سفارش است:

    using Confluent.Kafka;
    using System.Text.Json;
    
    var config = new ProducerConfig
    {
        BootstrapServers = "localhost:9092",
        Acks = Acks.All,           // wait until all in-sync replicas have the message
        EnableIdempotence = true   // retries do not write the same message twice
    };
    
    using var producer = new ProducerBuilder<string, string>(config).Build();
    
    var order = new OrderPlaced(OrderId: "o-3", CustomerId: "c-42", Total: 250_000);
    
    await producer.ProduceAsync("orders", new Message<string, string>
    {
        Key = order.OrderId,                 // same key => same partition
        Value = JsonSerializer.Serialize(order)
    });
    
    public sealed record OrderPlaced(string OrderId, string CustomerId, decimal Total);

    دو تنظیم مهم در این کد:

    • تنظیم Acks روی All. کافکا فقط وقتی جواب «ذخیره شد» می‌دهد که همه کپی‌های هماهنگ پیام را داشته باشند. پس اگر یک سرور بمیرد، پیام گم نمی‌شود.
    • تنظیم EnableIdempotence. اگر شبکه قطع شود و Producer دوباره بفرستد، کافکا پیام تکراری را در Partition نمی‌نویسد.
    مراقب Dual Write باش: اگر اول سفارش را در دیتابیس ذخیره کنی و بعد به کافکا بفرستی، ممکن است وسط این دو کار سرویس بمیرد. سفارش ذخیره شده، ولی پیامی نرفته است. راه درست الگوی Outbox است: پیام را در همان تراکنش دیتابیس در یک جدول ذخیره کن و یک کار پس‌زمینه آن را به کافکا بفرستد.

    خواندن پیام (Consumer)

    خواننده باید بعد از کار، به کافکا بگوید «تا اینجا را انجام دادم». به این کار Commit کردن Offset می‌گوییم. زمان Commit خیلی مهم است:

    حالت ۱: اول تأیید، بعد پردازشCommitتأیید خواندنپردازشسرویس وسط کار می‌میردپیام گم شددیگر کسی آن را نمی‌خواندحالت ۲: اول پردازش، بعد تأییدپردازشپیامک رفتCommitقبل از آن می‌میردپیام دوباره می‌آیدپس مصرف‌کننده باید تکرار را بفهمد
    هیچ حالتی «دقیقاً یک بار» نیست. حالت ۲ را انتخاب می‌کنیم و تکرار را در کد خودمان مدیریت می‌کنیم.
    1. اگر قبل از پردازش Commit کنی، و بعد سرویس بمیرد، پیام گم می‌شود. به این حالت at-most-once می‌گوییم.
    2. اگر بعد از پردازش Commit کنی، و قبل از Commit سرویس بمیرد، پیام دوباره می‌آید. به این حالت at-least-once می‌گوییم.
    3. پس حالت دوم را انتخاب می‌کنیم. گم شدن پیام بدتر از تکرار است. تکرار را با یک خواننده Idempotent حل می‌کنیم.
    using Confluent.Kafka;
    
    var config = new ConsumerConfig
    {
        BootstrapServers = "localhost:9092",
        GroupId = "email-service",
        EnableAutoCommit = false,                  // we commit ourselves, after the work
        AutoOffsetReset = AutoOffsetReset.Earliest, // a new group starts from the beginning
        PartitionAssignmentStrategy = PartitionAssignmentStrategy.CooperativeSticky
    };
    
    using var consumer = new ConsumerBuilder<string, string>(config).Build();
    consumer.Subscribe("orders");
    
    while (!ct.IsCancellationRequested)
    {
        var result = consumer.Consume(ct);
    
        // Idempotent: the handler checks if this order already got an email.
        await emailHandler.HandleAsync(result.Message.Key, result.Message.Value, ct);
    
        consumer.Commit(result); // only after the work is done
    }

    چند نکته درباره این کد:

    • تنظیم CooperativeSticky. در Rebalance فقط Partition هایی جابجا می‌شوند که لازم است، نه همه. پس توقف کوتاه‌تر است.
    • یک Consumer در هر thread. شیء Consumer در Confluent.Kafka برای استفاده همزمان از چند thread ساخته نشده است.
    • خواننده Idempotent. شناسه هر پیام را در جدولی با کلید یکتا ثبت کن. اگر قبلاً ثبت شده بود، کار را دوباره انجام نده.
    درباره exactly-once: کافکا با Producer تراکنشی می‌تواند «خواندن از یک topic و نوشتن در topic دیگر» را دقیقاً یک بار انجام دهد. ولی این فقط داخل خود کافکا است. فرستادن ایمیل یا نوشتن در دیتابیس تو داخل این تراکنش نیست. پس برای این کارها باز هم به Idempotency نیاز داری.

    انتخاب کلید و تعداد Partition

    کلید دو کار با هم انجام می‌دهد: ترتیب و پخش بار. این دو گاهی با هم دعوا دارند.

    1. اول بپرس ترتیب کجا لازم است. در فروشگاه ما، ترتیب پیام‌های یک سفارش مهم است. ترتیب بین دو سفارش مختلف مهم نیست.
    2. کوچک‌ترین کلیدی را انتخاب کن که ترتیب لازم را حفظ کند. شناسه سفارش بهتر از شناسه مشتری است. یک مشتری بزرگ با هزار سفارش، با کلید سفارش روی همه Partition ها پخش می‌شود.
    3. از اول کمی جا برای رشد بگذار. بعداً می‌شود Partition اضافه کرد، ولی نمی‌شود کم کرد.
    4. اضافه کردن Partition جای بعضی کلیدها را عوض می‌کند. پس در آن لحظه ترتیب پیام‌های آن کلیدها ممکن است به هم بریزد.
    5. عددی بگذار که خوب تقسیم شود. مثلاً دوازده Partition بین ۲، ۳، ۴ یا ۶ pod مساوی تقسیم می‌شود.
    مانیتورینگ: معیار اصلی سلامت خواننده‌ها Consumer Lag است، یعنی تعداد پیام‌های خوانده‌نشده. آن را برای هر Partition جدا ببین. مجموع کل گروه، یک Partition داغ را پنهان می‌کند.

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

    اشتباه نتیجه راه درست
    فکر کنی ترتیب در کل topic حفظ می‌شود پیام «لغو» قبل از «ثبت» پردازش می‌شود. پیام‌های مرتبط را با یک کلید بفرست.
    بیشتر کردن pod از تعداد Partition pod های اضافه بیکار می‌مانند و فقط هزینه دارند. تعداد pod را حداکثر برابر Partition بگذار.
    کلید خیلی درشت، مثل شناسه مشتری یک Partition داغ می‌شود و بقیه بیکارند. کلید ریزتر، مثل شناسه سفارش.
    Commit قبل از پردازش با هر crash یا deploy، پیام گم می‌شود. اول پردازش، بعد Commit.
    فکر کنی پیام تکراری نمی‌آید مشتری دو ایمیل یا دو پیامک می‌گیرد. خواننده Idempotent بساز.
    ذخیره در دیتابیس و بعد ارسال به کافکا، بدون Outbox سفارش هست، ولی پیامش هیچ وقت نمی‌رسد. الگوی Outbox.
    مرتب کردن رویدادها با ساعت سرورها ساعت سرورها با هم فرق دارد و ترتیب غلط می‌شود. ترتیب Partition یا شماره نسخه از یک منبع.

    چه وقت Kafka؟

    مناسب

    • چند سرویس مختلف همان رویدادها را لازم دارند.
    • حجم پیام زیاد است.
    • باید بشود پیام‌های قدیمی را دوباره خواند.
    • ترتیب رویدادهای هر موجودیت (مثلاً هر سفارش) مهم است.

    نامناسب

    • فقط یک صف کار ساده با چند کارگر لازم داری. یک صف مثل RabbitMQ ساده‌تر است.
    • مسیریابی پیچیده بر اساس نوع پیام لازم داری.
    • هر پیام باید جدا تأیید یا رد شود و بقیه پشت آن منتظر نمانند.
    • سیستم کوچک است و تیم توان نگهداری یک کلاستر را ندارد.

    خلاصه در هفت خط

    1. کافکا یک دفتر ثبت است. پیام بعد از خواندن پاک نمی‌شود.
    2. هر topic چند Partition دارد. ترتیب فقط داخل یک Partition حفظ می‌شود.
    3. کلید یکسان یعنی Partition یکسان. کلید را با فکر انتخاب کن.
    4. داخل یک Consumer Group، هر Partition فقط یک خواننده دارد.
    5. خواننده بیشتر از تعداد Partition، بیکار می‌ماند.
    6. اول پردازش کن، بعد Commit. پس پیام تکراری می‌آید و خواننده باید Idempotent باشد.
    7. برای ارسال مطمئن از دیتابیس به کافکا، از Outbox استفاده کن.