Kafka
کافکا یک دفتر ثبت بزرگ برای پیامها است. پیامها بعد از خواندن پاک نمیشوند و هر خواننده جای خودش را نگه میدارد. ترتیب فقط داخل یک Partition حفظ میشود و هر Partition در هر گروه فقط یک خواننده دارد.
نویسنده: bezzad
مشکل: یک سفارش، چند سرویس
مشتری در فروشگاه اینترنتی سفارش ثبت میکند. بعد از آن، چند سرویس باید کاری انجام دهند:
- سرویس ایمیل باید رسید خرید را بفرستد.
- سرویس انبار باید کالا را آماده کند.
- سرویس گزارش باید فروش امروز را بهروز کند.
راه ساده این است که سرویس سفارش هر سه را مستقیم صدا بزند. ولی این راه سه مشکل دارد:
- وابستگی زیاد. اگر سرویس ایمیل خراب باشد، ثبت سفارش هم کند یا خراب میشود.
- تغییر سخت. برای هر سرویس جدید، کد سرویس سفارش باید تغییر کند.
- بار ناگهانی. در روز حراج، همه سرویسها باید همان لحظه همه درخواستها را جواب دهند.
راه بهتر این است: سرویس سفارش فقط یک خبر ثبت میکند، «سفارش ثبت شد». هر سرویسی که لازم دارد، هر وقت آماده بود آن را میخواند. کافکا جای نگهداشتن همین خبرها است.
ایده اصلی: یک دفتر ثبت، نه یک صف
کافکا را مثل یک دفتر حساب تصور کن. هر خبر جدید در یک خط تازه، آخر دفتر نوشته میشود. چیزی وسط دفتر عوض نمیشود.
- پیام بعد از خواندن پاک نمیشود. پیام تا مدت مشخصی (Retention) میماند. پیشفرض این مدت یک هفته است و قابل تغییر است.
- هر خواننده جای خودش را دارد. هر خواننده یک عدد نگه میدارد: «تا خط چندم را خواندهام». به این عدد Offset میگوییم.
- پس خوانندهها کار هم را خراب نمیکنند. سرویس ایمیل و سرویس گزارش هر دو همه پیامها را میخوانند، هر کدام با سرعت خودش.
- میشود دوباره خواند. یک سرویس جدید میتواند از اول دفتر شروع کند و همه خبرهای قدیمی را ببیند.
Topic، Partition و Offset
هر نوع خبر یک Topic دارد. مثلاً همه خبرهای سفارش در topic به نام orders میروند.
یک دفتر تنها برای حجم زیاد کند است. پس هر topic به چند بخش تقسیم میشود. به هر بخش Partition میگوییم. هر Partition یک دفتر جدا است که روی یک سرور نگه داشته میشود (و چند کپی روی سرورهای دیگر دارد).
کافکا از کجا میداند هر پیام به کدام Partition برود؟
- هر پیام یک کلید (Key) دارد. ما در فروشگاه، شناسه سفارش را کلید میگذاریم.
- کافکا از کلید یک عدد میسازد (Hash). باقیمانده این عدد بر تعداد Partition ها، شماره Partition است.
- پس کلید یکسان، همیشه Partition یکسان. همه پیامهای سفارش o-3 در یک Partition و به همان ترتیب ارسال قرار میگیرند.
- اگر کلید نگذاریم، پیامها بین Partition ها پخش میشوند و هیچ ترتیبی بین آنها حفظ نمیشود.
چند خواننده برای یک کار: Consumer Group
سرویس ایمیل روی چند pod اجرا میشود. نمیخواهیم هر pod به هر مشتری یک ایمیل جدا بفرستد. پس همه pod های ایمیل را در یک Consumer Group میگذاریم.
قانونهای گروه سادهاند:
- هر گروه همه پیامهای topic را میخواند. گروه ایمیل و گروه انبار هر دو همه سفارشها را میبینند.
- داخل یک گروه، هر Partition فقط یک خواننده دارد. پس هیچ پیامی در یک گروه همزمان دو بار پردازش نمیشود.
- یک خواننده میتواند چند Partition داشته باشد.
- هر گروه Offset خودش را دارد. اگر گروه انبار عقب باشد، گروه ایمیل منتظر آن نمیماند.
نتیجه مهم قانون دوم: بیشترین تعداد خواننده مفید در یک گروه، برابر تعداد Partition ها است. اگر شش Partition و ده pod داشته باشی، چهار pod بیکار میمانند.
مثال زنده
دکمهها را بزن. تعداد Consumer را تا هشت بالا ببر و ببین چه میشود. بعد سفارشهای مشتری بزرگ را بفرست:
کلید پیام در دکمه اول شناسه سفارش است. در دکمه دوم کلید شناسه مشتری است.
دو چیز را در این مثال دیدی:
- وقتی تعداد Consumer عوض میشود، گروه دوباره تقسیم میشود. به این کار Rebalance میگوییم. در این مدت خواندن پیام کند یا متوقف میشود.
- وقتی یک کلید خیلی پرکار است، یک 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 نمینویسد.
خواندن پیام (Consumer)
خواننده باید بعد از کار، به کافکا بگوید «تا اینجا را انجام دادم». به این کار Commit کردن Offset میگوییم. زمان Commit خیلی مهم است:
- اگر قبل از پردازش Commit کنی، و بعد سرویس بمیرد، پیام گم میشود. به این حالت at-most-once میگوییم.
- اگر بعد از پردازش Commit کنی، و قبل از Commit سرویس بمیرد، پیام دوباره میآید. به این حالت at-least-once میگوییم.
- پس حالت دوم را انتخاب میکنیم. گم شدن پیام بدتر از تکرار است. تکرار را با یک خواننده 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. شناسه هر پیام را در جدولی با کلید یکتا ثبت کن. اگر قبلاً ثبت شده بود، کار را دوباره انجام نده.
انتخاب کلید و تعداد Partition
کلید دو کار با هم انجام میدهد: ترتیب و پخش بار. این دو گاهی با هم دعوا دارند.
- اول بپرس ترتیب کجا لازم است. در فروشگاه ما، ترتیب پیامهای یک سفارش مهم است. ترتیب بین دو سفارش مختلف مهم نیست.
- کوچکترین کلیدی را انتخاب کن که ترتیب لازم را حفظ کند. شناسه سفارش بهتر از شناسه مشتری است. یک مشتری بزرگ با هزار سفارش، با کلید سفارش روی همه Partition ها پخش میشود.
- از اول کمی جا برای رشد بگذار. بعداً میشود Partition اضافه کرد، ولی نمیشود کم کرد.
- اضافه کردن Partition جای بعضی کلیدها را عوض میکند. پس در آن لحظه ترتیب پیامهای آن کلیدها ممکن است به هم بریزد.
- عددی بگذار که خوب تقسیم شود. مثلاً دوازده Partition بین ۲، ۳، ۴ یا ۶ pod مساوی تقسیم میشود.
اشتباههای رایج
| اشتباه | نتیجه | راه درست |
|---|---|---|
| فکر کنی ترتیب در کل topic حفظ میشود | پیام «لغو» قبل از «ثبت» پردازش میشود. | پیامهای مرتبط را با یک کلید بفرست. |
| بیشتر کردن pod از تعداد Partition | pod های اضافه بیکار میمانند و فقط هزینه دارند. | تعداد pod را حداکثر برابر Partition بگذار. |
| کلید خیلی درشت، مثل شناسه مشتری | یک Partition داغ میشود و بقیه بیکارند. | کلید ریزتر، مثل شناسه سفارش. |
| Commit قبل از پردازش | با هر crash یا deploy، پیام گم میشود. | اول پردازش، بعد Commit. |
| فکر کنی پیام تکراری نمیآید | مشتری دو ایمیل یا دو پیامک میگیرد. | خواننده Idempotent بساز. |
| ذخیره در دیتابیس و بعد ارسال به کافکا، بدون Outbox | سفارش هست، ولی پیامش هیچ وقت نمیرسد. | الگوی Outbox. |
| مرتب کردن رویدادها با ساعت سرورها | ساعت سرورها با هم فرق دارد و ترتیب غلط میشود. | ترتیب Partition یا شماره نسخه از یک منبع. |
چه وقت Kafka؟
مناسب
- چند سرویس مختلف همان رویدادها را لازم دارند.
- حجم پیام زیاد است.
- باید بشود پیامهای قدیمی را دوباره خواند.
- ترتیب رویدادهای هر موجودیت (مثلاً هر سفارش) مهم است.
نامناسب
- فقط یک صف کار ساده با چند کارگر لازم داری. یک صف مثل RabbitMQ سادهتر است.
- مسیریابی پیچیده بر اساس نوع پیام لازم داری.
- هر پیام باید جدا تأیید یا رد شود و بقیه پشت آن منتظر نمانند.
- سیستم کوچک است و تیم توان نگهداری یک کلاستر را ندارد.
خلاصه در هفت خط
- کافکا یک دفتر ثبت است. پیام بعد از خواندن پاک نمیشود.
- هر topic چند Partition دارد. ترتیب فقط داخل یک Partition حفظ میشود.
- کلید یکسان یعنی Partition یکسان. کلید را با فکر انتخاب کن.
- داخل یک Consumer Group، هر Partition فقط یک خواننده دارد.
- خواننده بیشتر از تعداد Partition، بیکار میماند.
- اول پردازش کن، بعد Commit. پس پیام تکراری میآید و خواننده باید Idempotent باشد.
- برای ارسال مطمئن از دیتابیس به کافکا، از Outbox استفاده کن.