تصور کنید برای انتقال دادهها بین دو خوشه کافکا، مجبور باشید یک زیرساخت پیچیده و مجزا را مدیریت کنید که هر لحظه احتمال خطا در آن وجود دارد. این سردرد عملیاتی، تعریف تاریخی فرآیند جابهجایی دادهها بین خوشههای Apache Kafka بوده است. در ۲۲ سپتامبر ۲۰۲۶، شرکت رد هت (Red Hat) جزئیاتی را منتشر کرد که نشان میدهد چگونه استاندارد KIP-1279 با ادغام مستقیم تکثیر بینخوشهای در بروکر کافکا، نیاز به ابزارهای خارجی را بهکلی از بین میبرد.
سالهاست که سازمانها برای توزیع جغرافیایی، رعایت مرزهای قانونی و انطباق (Compliance)، جداسازی تیمها یا تفکیک نسخهها، از چندین خوشه (Cluster) — شبیه به داشتن چندین شعبه از یک بانک در شهرهای مختلف برای دسترسی سریعتر مشتریان — استفاده میکنند. اما طبق گزارشهای فنی، ایجاد یک کپی دقیق و وفادار از دادهها در خوشهای دیگر، همواره نیازمند فرآیندهای خارجی و تحمل ریسکهای عملیاتی بالا بوده است. اکثر تیمها به MirrorMaker 2 (MM2) متکی بودند که از زمان نسخه ۲.۴ کافکا به عنوان استاندارد شناخته میشد. MM2 در واقع مجموعهای از کارکنان (Workers) Kafka Connect بود که دادهها را از یک مبدأ مصرف (Consume) کرده و در مقصد تولید (Produce) میکردند.
همانطور که در تحلیلهای پیشین ما دربارهی بهینهسازی زیرساختهای داده اشاره کردیم، حذف لایههای واسط همواره منجر به کاهش تأخیر میشود. این رویکرد بومی، معماری بنیادی جابهجایی داده را تغییر میدهد. اکنون بروکر مقصد بهجای تماشای دادهها از دور یا استفاده از یک پل خارجی، خود به یک شرکت فعال تبدیل شده است. این بروکر با استفاده از همان پروتکل fetch که دنبالکنندههای (Followers) داخلی از آن استفاده میکنند، رکوردهای تایید شده (Committed) را دریافت کرده و آنها را بهصورت بایت-به-بایت به لاگهای محلی اضافه میکند. این یعنی دیگر نیازی به زیرساختهای خارجی، جداول ترجمه آفست و فشردهسازی مجدد نیست.
معماری Mirroring بومی
به نقل از مستندات فنی KIP-1279، سه مؤلفه اصلی در هر بروکر مقصد برای مدیریت این فرآیند همکاری میکنند. شکل ۱ نشان میدهد که این اجزا چگونه به یکدیگر و به خوشه مبدأ متصل میشوند:
- MirrorMetadataManager (MMM): این بخش ارکستراتوری است که با پیادهسازی رابط
MetadataPublisherبه تغییرات لاگ متادیتای KRaft واکنش نشان میدهد. هنگامی که کنترلر یک رکوردMirrorTopicStateChangeRecordرا مینویسد، MMM روی لیدر پارتیشن مربوطه، انتقال وضعیت (ایجاد، شروع، توقف، توقف موقت،e resume، بازیابی یا حذف) را هدایت میکند. این مدیر یک اتصال Admin client به خوشه مبدأ حفظ میکند و بهطور پیشفرض هر ۶۰ ثانیه یکبار متاداتا را بهروز میکند تا موضوعات (Topics) جدیدی که با الگوهای include/exclude مطابقت دارند را شناسایی کند، تنظیمات را همگامسازی نماید و آفستهای گروههای مصرفکننده را دریافت کند. همچنین، MMM تایید میکند که ID خوشه مبدأ تغییر نکرده باشد تا از فساد خاموش دادهها (Silent Data Corruption) جلوگیری شود. - ClusterMirrorCoordinator (CMC): مدیریت پایداری وضعیت را بر عهده دارد و از همان الگوی Coordinator که در هماهنگکنندههای گروه و تراکنش استفاده میشود، پیروی میکند. این بخش شاردهای (Shards) یک موضوع فشرده داخلی به نام
__mirror_stateرا مدیریت میکند که بهطور پیشفرض دارای ۵۰ پارتیشن و ضریب تکثیر ۳ است. وضعیت هر پارتیشن Mirror بهصورت یک رکورد کلید-مقدار ذخیره میشود و از کنترل همزمانی خوشبینانه (Optimistic Concurrency Control) از طریق leader epoch و state epoch fencing استفاده میکند. - MirrorFetcherThread (MFT): موتور اصلی است که رکوردها را از مبدأ میکشد و در واقع توسعهیافتهی
AbstractFetcherThreadکافکا است. این بخش از یکNetworkClientاختصاصی با اعتبارنامههای احراز هویت مجزا برای هر Mirror استفاده میکند تا زمینههای (Contexts) SASL/SSL ایزوله بمانند. مدیر fetcher، رشتهها را بر اساس یک شناسه سهبعدی (شناسه fetcher، نقطه انتهایی بروکر مبدأ، نام mirror) کلیدبندی میکند تا توازن بار دقیق و پاسخ سریع به تغییرات لیدر مبدأ امکانپذیر شود.
مزایای فنی نسبت به MirrorMaker 2
یکی از بزرگترین دستاوردهای این معماری، حذف بار پردازشی CPU برای باز کردن (Decompress) و بستن مجدد (Recompress) فایلهای فشرده است. فرقی نمیکند یک دسته داده از gzip, snappy, lz4 یا zstd استفاده کند؛ دادهها در مقصد دقیقاً به همان شکل اصلی خود میرسند. این انتقال بایت-به-بایت، انتخابهای اصلی تولیدکننده در فشردهسازی را حفظ کرده و چرخه زمانبر باز/بستن را حذف میکند.
حفظ آفست (Offset) — که مثل شماره صفحه در یک کتاب است و میگوید مصرفکننده دقیقاً کجا متوقف شده — پیروزی دوم این سیستم است. در تنظیمات قبلی، آفستها اغلب دارای فقدان داده بودند یا از طریق یک موضوع (Topic) ترجمه میشدند. اکنون لاگ مقصد دقیقاً همان آفستهای مبدأ را حفظ میکند، حتی شکافهایی (Gaps) که بر اثر فشردهسازی موضوع (Topic Compaction) ایجاد شدهاند. این بدان معناست که گروههای مصرفکننده میتوانند بدون نیاز به جداول پیچیده ترجمه آفست، عملیات Failover را انجام دهند؛ زیرا آفست تایید شده در مبدأ، همان آفست تایید شده در مقصد است.
کنترل پهنای باند و منابع
به دلیل ادغام در بروکر، کنترل پهنای باند اکنون بخشی از مدیریت منابع داخلی بروکر است. بروکر مقصد یک حد مجاز برای نرخ تکثیر (Replication Rate Limit) اعمال میکند تا از اشباع شدن شبکه توسط فرآیند Mirroring جلوگیری کند.
در سمت مبدأ نیز، ترافیک fetch مربوط به mirror مانند درخواستهای استاندارد مصرفکننده تلقی میشود. این یعنی مکانیزمهای موجود برای سهمیهبندی کلاینت (Client Quota) بدون هیچ تغییری اعمال میشوند و مدیران میتوانند ترافیک تکثیر را با همان ابزارهایی که برای مصرفکنندگان تولیدی استفاده میکنند، محدود (Throttle) کنند.
مقایسه: MM2 در برابر Cluster Mirroring
| ویژگی | MirrorMaker 2 | Cluster Mirroring |
|---|---|---|
| استقرار | کارکنان خارجی Connect | ادغام شده در بروکر |
| فشردهسازی | باز و بسته کردن مجدد | انتقال بایت-به-بایت |
| آفستها | ترجمه تقریبی/با خطا | کاملاً یکسان |
| جایگزینی مصرفکننده | پرسوجو از موضوع همگامساز | مستقیم و بدون ترجمه |
| انتخابات ناپاک | بدون مدیریت | همگرایی کامل لاگ |
| سازگاری مبدأ | Kafka 2.0+ | Kafka 2.1+ |
| نظارت | ابزارهای خاص Connect | متریکهای استاندارد JMX بروکر |
چرخه حیات پارتیشن Mirror
یک پارتیشن برای تضمین سازگاری، از یک ماشین وضعیت (State Machine) قطعی عبور میکند. شکل ۲ این جریان را نشان میدهد، جایی که هر وضعیتی در صورت بروز خطا میتواند به وضعیت FAILED منتقل شود:
۱. LOG_ALIGNMENT: بروکر لاگ محلی را با مبدأ تراز میکند. اگر یک رکورد Last Mirror Epoch (LME) یافت شود، بروکر رکوردها را فقط تا آخرین آفست و epoch مربوط به mirror برش میزند تا تضادها برطرف شود. اگر LME وجود نداشته باشد (تکثیر برای اولین بار یا مبدأ پشتیبانی نشده)، لاگ تا صفر برش خورده و تکثیر از ابتدا شروع میشود.
۲. EPOCH_FENCING: بروکر یک درخواست BumpLeaderEpochs به کنترلر ارسال میکند و epoch لیدر محلی را ۱۰ واحد افزایش میدهد (با آستانه re-bump برابر با ۳). این کار از رد شدن لیدر توسط مصرفکنندگان به دلیل قدیمی بودن (Stale) در صورتی که epoch مبدأ از epoch محلی بیشتر باشد، جلوگیری میکند.
۳. MIRRORING: رشته fetcher رکوردها را میکشد و با تکثیر دادهها توسط دنبالکنندههای محلی، High Watermark را جلو میبرد. آفستهای گروه بهطور دورهای از مبدأ همگام شده و به محدوده آفست معتبر در مقصد محدود (Clamp) میشوند.
۴. ULE_RECOVERY: اگر یک انتخابات لیدر ناپاک (Unclean Leader Election) در مبدأ رخ دهد، لیدر مقصد تضاد لاگ را تشخیص میدهد. اگر mirror.unclean.leader.election.enable فعال باشد، پارتیشن وارد این وضعیت شده، fetcher را حذف میکند و منتظر میماند تا تمام نسخههای تخصیص یافته (نه فقط اعضای ISR) با انتهای لاگ برشخورده همگرا شوند و سپس تکثیر را از سر میگیرد.
۵. PAUSING / PAUSED: اپراتورها میتوانند تکثیر را متوقف کنند. در این حالت رشتههای fetcher تخریب شده و پارتیشن فقط خواندنی میماند. بازگشت از این وضعیت، پارتیشن را مستقیماً به حالت MIRRORING برمیگرداند.
۶. STOPPING / STOPPED: مسیر جایگزینی (Failover). بروکر fetcherها را حذف کرده، رکورد LME را ثبت میکند، epoch محلی را افزایش داده، نشانگرهای ABORT را برای تراکنشهای در جریان اضافه میکند و یک رکورد کنترلی MIRROR_PID_RESET مینویسد تا وضعیت تولیدکننده منقضی شود. سپس پارتیشن برای نوشتن باز میشود.
مدیریت خطا و بازیابی
وقتی پارتیشنی به وضعیت FAILED میرود، سیستم بهطور ساده متوقف نمیشود. بلکه یک مکانیزم تلاش مجدد خودکار با استفاده از backoff نمایی همراه با jitter فعال میشود. این فرآیند تا تعداد دفعات حداکثری که توسط کاربر تنظیم شده، ادامه مییابد.
با این حال، برخی خطاها به عنوان «غیرقابل تلاش» (Non-retryable) طبقهبندی میشوند. این موارد نیاز به دخالت دستی اپراتور دارند. نمونههایی از این خطاها شامل تغییر در ID خوشه مبدأ یا حذف موضوع در مبدأ است که هر دو باعث شکست اعتماد بنیادی و نگاشت بین خوشهها میشوند.
بازیابی از فاجعه و جایگزینی (Failover)
بازیابی اکنون به چند دستور ساده در خط فرمان تبدیل شده است. شکل ۳ سه فاز عملیاتی را نشان میدهد: عملیات عادی، جایگزینی و بازگشت.
عملیات عادی: خوشه A ترافیک تولیدی را مدیریت میکند. یک Mirror به نام a-to-b بهطور مداوم موضوعات را به خوشه B تکثیر میکند. اپراتورها با دستورات زیر Mirror را ایجاد و شروع میکنند:kafka-cluster-mirrors.sh --bootstrap-server B:9092 --create --mirror a-to-b --mirror-config mirror.propertieskafka-cluster-mirrors.sh --bootstrap-server B:9092 --start --mirror a-to-b --topics ".*"
نظارت از طریق فلگ --describe انجام میشود که به اپراتورها اجازه میدهد پیشرفت تکثیر را بهصورت لحظهای دنبال کنند:kafka-cluster-mirrors.sh --bootstrap-server B:9092 --describe --mirror a-to-b
جایگزینی (Failover): اگر خوشه A از کار بیفتد، یک دستور واحد در خوشه B تکثیر را متوقف کرده و موضوعات را برای نوشتن باز میکند:kafka-cluster-mirrors.sh --bootstrap-server B:9092 --stop --mirror a-to-b
تولیدکنندگان و مصرفکنندگان به خوشه B منتقل میشوند. چون آفستها یکسان هستند و آفستهای گروههای مصرفکننده همگام شدهاند، برنامهها بدون نیاز به پردازش مجدد دادهها، کار را از همان نقطه ادامه میدهند.
بازگشت (Failback): پس از بازیابی خوشه A، یک Mirror معکوس ایجاد میشود. سیستم تشخیص میدهد که A مبدأ قبلی بوده است، LME را از زمان Failover جستجو میکند و برشهای افزایشی (Incremental Truncation) را انجام میدهد:kafka-cluster-mirrors.sh --bootstrap-server A:9092 --create --mirror b-to-a --mirror-config reverse.propertieskafka-cluster-mirrors.sh --bootstrap-server A:9092 --start --mirror b-to-a --topics ".*"
زمانی که تکثیر به سطح برابری رسید، اپراتور دستور --stop را در خوشه A اجرا میکند تا انتقال نهایی صورت گیرد.
هدف بازیابی (RPO) و تکثیر همزمان
هدف بازیابی یا RPO به تأخیر تکثیر (Replication Lag) بستگی دارد، زیرا این فرآیند بهصورت ناهمزمان (Asynchronous) است. رکوردهایی که در A تولید شده اما هنوز به B منتقل نشدهاند، در یک Failover برنامهریزی نشده از دست خواهند رفت. برای اکثر موارد بازیابی از فاجعه (DR)، این یک سبکسنگین پذیرفتنی برای اجتناب از جریمههای تأخیر در تکثیر همزمان بینخوشهای است.
با این حال، برای سازمانهایی با نیاز به Zero-RPO، استاندارد KIP-1360 در حال برنامهریزی است تا KIP-1279 را با یک حالت تکثیر همزمان (Synchronous Mirroring Mode) گسترش دهد.
تسهیل مهاجرت خوشهها
KIP-1279 یک راه میانبر برای مهاجرت از خوشههای قدیمی مبتنی بر ZooKeeper به خوشههای مدرن KRaft فراهم میکند. بهطور سنتی، این کار نیازمند ارتقای گامبهگام از طریق هر نسخه اصلی (مثلاً ۲.x به ۳.x و سپس ۴.x) و مهاجرت متاداتا در جای خود بود.
با Mirroring بومی، شما میتوانید یک خوشه KRaft جدید راه بیندازید و دادهها را مستقیماً از مبدأهایی به قدیمی نسخهی ۲.۱ تکثیر کنید و از سازگاری رو به جلوی کلاینت/بروکر در کافکا ۴.۰ بهره ببرید. شکل ۴ این فرآیند را نشان میدهد:
۱. راهاندازی: ایجاد یک Mirror از خوشه قدیمی به خوشه جدید.
۲. تکثیر: شروع Mirroring. مقصد بهطور خودکار موضوعات را شناسایی کرده، پارتیشنهای متناظر با IDهای یکسان ایجاد میکند، تنظیمات را همگامسازی کرده و تکثیر دادهها را آغاز میکند.
۳. انتقال (Cutover): توقف تولیدکنندگان در خوشه قدیمی، نظارت بر تأخیر تا رسیدن به صفر و سپس اجرای دستور --stop در خوشه جدید برای باز کردن دسترسی نوشتن.
مثال دستورات در خوشه جدید:kafka-cluster-mirrors.sh --bootstrap-server new:9092 --create --mirror old-to-new --mirror-config old-cluster.propertieskafka-cluster-mirrors.sh --bootstrap-server new:9092 --start --mirror old-to-new --topics ".*"
نظارت بر پیشرفت:kafka-cluster-mirrors.sh --bootstrap-server new:9092 --describe --mirror old-to-new
در نهایت، وقتی تأخیر صفر شد، انتقال نهایی اجرا میشود:kafka-cluster-mirrors.sh --bootstrap-server new:9092 --stop --mirror old-to-new
این روش بهطور کامل اسکریپتهای تبدیل ZooKeeper و پرشهای نسخهای میانی را دور میزند. خوشه قدیمی تا زمان بازنشستگی نهایی بدون تغییر باقی میماند.
تضمین سازگاری دادهها
برای جلوگیری از فساد دادهها، سیستم یک قانون سختگیرانه (Invariant) را حفظ میکند: Epoch لیدر مقصد باید همیشه بزرگتر یا مساوی Epoch لیدر مبدأ باشد (DLE >= SLE). بدون این قانون، یک مصرفکننده در مقصد ممکن است با یک epoch تایید شده از مبدأ شروع به کار کند که از epoch محلی بیشتر است، در نتیجه درخواست fetch را رد کرده و متوقف شود.
سیستم این فاصله را از طریق سه نوع Bump (افزایش) حفظ میکند:
- bumps واکنشی: زمانی فعال میشوند که epoch دسته دریافت شده به epoch محلی نزدیک شود.
- bumps پیشدستانه: زمانی فعال میشوند که فاصله به کمتر از ۳ برسد.
- bumps دورهای: در طول همگامسازی متادیتای مبدأ رخ میدهند.
هر Bump، مقدار epoch را ۱۰ واحد افزایش میدهد.
همگرایی لاگ و برش
همگرایی لاگ از طریق یک پروتکل برش دو مرحلهای مدیریت میشود. فاز اولیه از رکورد LME در طول LOG_ALIGNMENT استفاده میکند. فاز وضعیت پایدار (Steady-state)، آفستها را در حین پردازش fetch تراز میکند.
بخش fetcher بهطور پویا نسخه fetch خوشه مبدأ را تشخیص میدهد:
- Fetch v12+: fetcher از اطلاعات epoch واگرا در پاسخهای fetch برای برش داخلی (Inline Truncation) تا نقطه دقیق واگرایی استفاده میکند.
- نسخههای قدیمیتر: fetcher قابلیت برش در هنگام fetch را غیرفعال کرده و به درخواست صریح
OffsetsForLeaderEpochبرای یافتن نقطه برش بازمیگردد.
ایمنی تراکنشها و بازنشانی PID
ایمنی تراکنشها از طریق ایزولاسیون READ_UNCOMMITTED در طول تکثیر مدیریت میشود تا تأخیر کاهش یابد. این بدان معناست که رکوردهای تراکنشهای تایید نشده پیش از آنکه مبدأ درباره سرنوشت آنها تصمیم بگیرد، به مقصد میرسند.
در طول انتقال به وضعیت توقف (Stopping)، بروکر نشانگرهای ABORT را برای تراکنشهای در جریان اضافه میکند. همچنین یک رکورد MirrorPidResetRecord مینویسد تا تمام ورودیهای وضعیت تولیدکننده در ProducerStateManager منقضی شوند. این کار اجازه میدهد تولیدکنندگان جدید در مقصد، شناسههای تولیدکننده (Producer IDs) تازهای را بدون تداخل دریافت کنند.
این مکانیزم بازنشانی PID بهطور صحیح در تمام توپولوژیهای تکثیر منتشر میشود:
- فعال-غیرفعال: A به B
- زنجیرههای بازگشت: A به B و سپس B به A
- توزیع بادبزنی (Fan-out): A به B و A به C
- زنجیرههای چندگانه: A به B و B به C
این تغییر معماری، کافکا را به سمت یک بافت دادهای (Data Fabric) مستقلتر سوق میدهد. با حذف «واسطه» (Kafka Connect) برای تکثیر، تعداد قطعات متحرکی که میتوانند در زمان بحران خراب شوند، کاهش مییابد. برای مهندسان، این یعنی جایگزینی «دفترچههای راهنمای پیچیده» برای بازیابی از فاجعه با مجموعهای از فراخوانیهای استاندارد API.
اگر به دلیل ریسک پرشهای نسخهای میانی، مهاجرت خوشه را به تعویق انداختهاید، Mirroring بومی یک استراتژی خروج تمیز را فراهم میکند. توصیه میشود KIP-1360 را که قصد معرفی حالت تکثیر همزمان برای نیازهای Zero-RPO را دارد، در کنار ادغامهای آینده با tiered storage و پشتیبانی از موضوعات بدون دیسک (Diskless Topics) دنبال کنید.




گفتگو