
В Разгледахме клъстеризацията на RabbitMQ за осигуряване на отказоустойчивост и висока достъпност. Сега ще се задълбочим в Apache Kafka.
Тук единицата за репликация е раздел (partition). Всеки топик има един или няколко раздела. Във всеки раздел има лидер с последователи или без тях. При създаването на топик се посочва броят на разделите и коефициентът на репликация. Обичайната стойност е 3, което означава три реплики: един лидер и двама последователи.

Рис. 1. Четири раздела са разпределени между трима брокери
Всички заявки за четене и запис постъпват при лидера. Последователите периодично изпращат на лидера заявки за получаване на последните съобщения. Потребителите никога не се обръщат към последователите, които съществуват само за излишност и отказоустойчивост.

Сбой на раздела
Когато брокерът се откаже, често спират лидерите на няколко раздела. Във всеки от тях лидер става последовател от друг възел. В действителност не винаги е така, тъй като влияе и факторът на синхронизация: дали има синхронизирани последователи, а ако няма, дали е разрешен преходът към несинхронизирана реплика. Но да не усложняваме.
От мрежата излиза брокер 3 — и за раздел 2 се избира нов лидер на брокер 2.

Рис. 2. Брокер 3 умира и неговият последовател на брокер 2 се избира за нов лидер на раздел 2
След това брокер 1 също излиза и раздел 1 губи своя лидер, чиято роля преминава към брокер 2.

Рис. 3. Остава само един брокер. Всички лидери са на един брокер с нулева излишност
Когато брокер 1 се връща в мрежата, той добавя четирима последователи, осигурявайки известна излишност на всеки раздел. Но всички лидери остават на брокер 2.

Рис. 4. Лидерите остават на брокер 2
Когато брокер 3 се вдига, ние се връщаме към три реплики на раздела. Но всички лидери все още са на брокер 2.

Рис. 5. Несбалансирано разпределение на лидерите след възстановяване на брокерите 1 и 3
Kafka разполага с инструмент за по-добро ребалансиране на лидерите, отколкото RabbitMQ. Там се е налагало да се използва външен плъгин или скрипт, който променя политиките за миграция на главния възел в замяна на намалена излишност по време на миграцията. Освен това, за големи опашки е трябвало да се примирим с недостъпност по време на синхронизация.
В Kafka има концепция за „предпочитани реплики“ за ролята на лидера. Когато се създават раздели на тема, Kafka се опитва да разпредели лидерите равномерно върху възлите и маркира тези първи лидери като предпочитани. С времето, поради рестартиране на сървъри, сривове и нарушаване на свързаността, лидерите могат да се окажат на други възли, както в описания по-горе крайен случай.
За да се поправи това, Kafka предлага два варианта:
- Опция auto.leader.rebalance.enable=true позволява на контролера автоматично да преопредели лидерите обратно на предпочитаните реплики, възстановявайки по този начин равномерното разпределение.
- Администраторът може да стартира скрипта kafka-preferred-replica-election.sh за да преопредели ръчно.

Рис. 6. Реплики след презареждане
Това беше опростена версия на срив, но реалността е по-сложна, въпреки че няма нищо прекалено сложно тук. Всичко се свежда до синхронизирани реплики (In-Sync Replicas, ISR).
Синхронизирани реплики (ISR)
ISR е набор от реплики на раздел, който се счита за „синхронизиран“ (in-sync). Има лидер, но фолловери може да липсват. Фолловерът се счита за синхронизиран, ако е направил точни копия на всички съобщения на лидера до изтичането на интервала replica.lag.time.max.ms.
Фолловерът бива отстранен от набора ISR, ако:
- не е направил заявка за извличане в интервала replica.lag.time.max.ms (счита се за мъртъв)
- не е успял да се актуализира в интервала replica.lag.time.max.ms (счита се за бавен)
Фолловерите правят заявки за извличане в интервала replica.fetch.wait.max.ms, който по подразбиране е 500 мс.
За да обясним ясно целта на ISR, трябва да разгледаме потвържденията от производителя (producer) и някои сценарии на отказ. Производителите могат да избират, кога брокерът изпраща потвърждение:
- acks=0, потвърждение не се изпраща
- acks=1, потвърждение се изпраща след като лидерът запише съобщението в своя локален лог
- acks=all, потвърждение се изпраща след като всички реплики в ISR запишат съобщението в локалните лога
В терминологията на Kafka, ако ISR е запазил съобщението, настъпва неговото „комитване“. Acks=all е най-сигурният вариант, но и с допълнителна забавяне. Нека разгледаме два примера на отказ и как различните опции ‘acks’ взаимодействат с концепцията ISR.
Acks=1 и ISR
В този пример ще видим, че ако лидерът не очаква да получи всяко съобщение от всички последователи, може да настъпи загуба на данни при повреда на лидера. Преходът към несинхронизиран последовател може да бъде разрешен или забранен с настройката unclean.leader.election.enable.
В този пример производителят е конфигуриран с acks=1. Разделът е разпределен между трите брокера. Брокер 3 изостава, той се е синхронизирал с лидера преди осем секунди и сега изостава с 7456 съобщения. Брокер 1 е изостанал само с една секунда. Нашият производител изпраща съобщение и бързо получава отговор ack, без допълнителни разходи за бавни или мъртви последователи, които лидерът не очаква.

Рис. 7. ISR с три реплики
Брокер 2 се проваля и производителят получава грешка при свързването. След прехода на лидерството към брокер 1, губим 123 съобщения. Последователят на брокер 1 е присъствал в ISR, но не се е синхронизирал напълно с лидера, когато той се е провалил.

Рис. 8. При повреда се губят съобщения
В конфигурацията bootstrap.servers производителят е изброил няколко брокера и може да попита друг брокер кой е новият лидер на раздела. След това той установява връзка с брокер 1 и продължава да изпраща съобщения.

Рис. 9. Изпращането на съобщенията се възобновява след кратка пауза
Брокер 3 отстъпва още повече. Той прави запитвания за извличане, но не може да се синхронизира. Това може да се дължи на бавна мрежова връзка между брокерите, проблем с хранилището и т.н. Той е отстранен от ISR. Сега ISR се състои от една реплика — лидера! Производителят продължава да изпраща съобщения и получава потвърждения.

Рис. 10. Последователят на брокер 3 е отстранен от ISR
Брокер 1 пада, а ролята на лидера преминава към брокер 3 с загуба на 15286 съобщения! Производителят получава грешка при свързването. Преходът към лидер извън ISR е бил възможен само заради настройката unclean.leader.election.enable=true. Ако е настроена на false, тогава преходът нямаше да се случи, а всички заявки за четене и писане щяха да бъдат отказани. В този случай чакаме връщането на брокер 1 с непокътнатите му данни в репликата, която отново ще възвърне лидерството.

Рис. 11. Брокер 1 пада. При повреда се губи голямо количество съобщения
Производителят установява връзка с последния брокер и вижда, че той вече е лидер в секцията. Той започва да изпраща съобщения до брокер 3.

Рис. 12. След кратка пауза съобщенията отново се изпращат в секция 0
Видяхме, че освен кратките прекъсвания за създаване на нови връзки и намиране на нов лидер, производителят постоянно е изпращал съобщения. Тази конфигурация осигурява наличност на сметка на последователност (сигурност на данните). Kafka загуби хиляди съобщения, но продължи да приема нови записи.
Acks=all и ISR
Нека повторим този сценарий още веднъж, но с acks=all. Забавянето на брокер 3 е средно четири секунди. Производителят изпраща съобщение с acks=all, и сега не получава бърз отговор. Лидерът изчаква, докато съобщението бъде запазено от всички реплики в ISR.

Рис. 13. ISR с три реплики. Една от тях работи бавно, което води до закъснение в записа
След четири секунди допълнително забавяне брокер 2 изпраща ack. Всички реплики вече са напълно актуализирани.

Рис. 14. Всички реплики запазват съобщения и изпращат ack
Брокер 3 сега изостава още повече и е премахнат от ISR. Забавянето значително намалява, тъй като в ISR не останаха бавни реплики. Брокер 2 вече изчаква само брокер 1, а неговото средно закъснение е 500 мс.

Рис. 15. Реплика на брокер 3 е премахната от ISR
След това брокер 2 пада и лидерството преминава към брокер 1 без загуба на съобщения.

Рис. 16. Брокер 2 пада
Производителят намира нов лидер и започва да изпраща съобщения. Забавянето отново намалява, тъй като сега ISR се състои от една реплика! Следователно опцията acks=all не добавя излишък.

Рис. 17. Реплика на брокер 1 поема лидерството без загуба на съобщения
След това брокер 1 пада и лидерството преминава към брокер 3 с загуба на 14238 съобщения!

Рис. 18. Брокер 1 умира, а прехвърлянето на лидерството с настройката unclean води до значителна загуба на данни
Можехме да не задаваме опцията unclean.leader.election.enable на стойност true. По подразбиране тя е равна на false. Настройката acks=all с unclean.leader.election.enable=true осигурява наличност с известна допълнителна сигурност на данните. Но, както виждате, все още можем да загубим съобщения.
Но какво, ако искаме да увеличим сигурността на данните? Можем да зададем unclean.leader.election.enable = false, но това не е задължително да ни защити от загуба на данни. Ако лидера падне рязко и отнесе данни с себе си, то съобщенията все пак остават изгубени, плюс загубва се наличността, докато администраторът не възстанови ситуацията.
По-добре е да се гарантира излишък на всички съобщения, а в противен случай да се откаже записването. Тогава поне от гледна точка на брокера загубата на данни е възможна само при две или повече едновременни повреди.
Acks=all, min.insync.replicas и ISR
С конфигурацията на темата min.insync.replicas повишаваме нивото на сигурност на данните. Нека отново прегледаме последната част от предишния сценарий, но този път с min.insync.replicas=2.
И така, при брокер 2 има лидер реплика, а последователят на брокер 3 е изтрит от ISR.

Фиг. 19. ISR от две реплики
Брокер 2 пада, а лидерството преминава към брокер 1 без загуба на съобщения. Но сега ISR се състои само от една реплика. Това не отговаря на минималния брой за получаване на записи, и затова брокерът отговаря на опита за запис с грешка. NotEnoughReplicas.

Фиг. 20. Броят на ISR е с един по-малко, отколкото е посочено в min.insync.replicas
Тази конфигурация жертва наличността в името на съгласуваността. Преди да потвърдим съобщението, ние гарантираме, че то се записва поне на две реплики. Това дава на производителя много по-голяма увереност. Загубата на съобщения тук е възможна само при едновременна повреда на две реплики в кратък интервал, докато съобщението не е реплицирано на допълнителен последовател, което е малко вероятно. Но ако сте супер параноик, можете да настроите коефициента на репликация на 5, а min.insync.replicas на 3. Нужни са веднага три брокера да паднат едновременно, за да загубите запис! Разбира се, за такава надеждност ще платите с допълнителна латентност.
Когато наличността е необходима за сигурността на данните
Както и в , понякога наличността е необходима за сигурността на данните. Трябва да помислите за следното:
- Може ли публикуващият просто да върне грешка, а надзорната услуга или потребителят да опитат отново по-късно?
- Може ли публикуващият да запази съобщение локално или в база данни, за да опита отново по-късно?
Ако отговорът е отрицателен, оптимизацията на достъпността повишава сигурността на данните. Ще загубите по-малко данни, ако изберете достъпността вместо отказа от запис. По този начин всичко опира до намирането на баланс, а решението зависи от конкретната ситуация.
Смисъл на ISR
Наборът ISR позволява да изберете оптимален баланс между сигурността на данните и закъснението. Например, да осигурите достъпност при повреда на повечето реплики, минимизирайки влиянието на ненормално работещи или бавни реплики относно закъснението.
Ние сами избираме стойността replica.lag.time.max.ms в съответствие с нашите нужди. По същество този параметър означава каква закъснение сме готови да приемем при acks=all. Стойността по подразбиране е десет секунди. Ако това е твърде дълго за вас, можете да я намалите. Тогава честотата на промените в ISR ще се увеличи, тъй като последователите ще бъдат по-често премахвани и добавяни.
В RabbitMQ просто е набор от огледала, които трябва да се репликират. Бавни огледала въвеждат допълнително закъснение, а отговори от ненормални огледала могат да се чакат до изтичането на времето за живот на пакетите, които проверяват достъпността на всеки възел (net tick). ISR е интересен начин да се избегнат тези проблеми с увеличаване на закъснението. Но рискуваме да загубим излишък, тъй като ISR може да намалее само до лидера. За да избегнем този риск, използвайте настройката min.insync.replicas.
Гаранция за свързаност на клиентите
В настройките bootstrap.servers производител и потребител можете да зададете няколко брокера за свързване на клиентите. Идеята е, че при отключване на един възел остават няколко резервни, с които клиентът може да установи връзка. Те не е задължително да са лидери на раздели, а просто плацдарм за начално зареждане. Клиентът може да ги попита на кой възел се намира лидерът на раздела за четене/запис.
В RabbitMQ клиентите могат да се свързват с всеки възел, а вътрешното маршрутизиране изпраща заявките където трябва. Това означава, че можете да поставите балансировач на натоварването пред RabbitMQ. Kafka изисква клиентите да се свързват с възел, на който е разположен лидерът на съответния раздел. В такава ситуация не може да се постави балансировач на натоварването. Списъкът bootstrap.servers е критично важен, за да могат клиентите да се свързват с необходимите възли и да ги намират след повреда.
Архитектура на консенсуса на Kafka
Досега не разгледахме как клъстерът открива падането на брокера и как се избира нов лидер. За да разберем как Kafka работи с мрежови разрези, първо трябва да разберем архитектурата на консенсуса.
Всеки клъстер Kafka се разгръща заедно с клъстера Zookeeper — това е услуга за разпределен консенсус, която позволява на системата да постигне консенсус относно определено зададено състояние с приоритет на последователността пред наличността. За одобрение на операции за четене и писане е необходимо съгласие на повечето възли Zookeeper.
Zookeeper съхранява състоянието на клъстера:
- Списък на теми, разрези, конфигурация, текущи лидери на реплики, предпочитани реплики.
- Членове на клъстера. Всеки брокер изпраща пинг до клъстера Zookeeper. Ако той не получи пинг в зададения период от време, Zookeeper записва брокера като недостъпен.
- Избор на основни и резервни възли за контролер.
Възел контролер — един от брокерите на Kafka, който отговаря за избора на лидери на реплики. Zookeeper изпраща уведомления на контролера за членството в клъстера и промените в темата, а контролерът трябва да действа в съответствие с тези промени.
Например, вземаме нова тема с десет разреза и коефициент на репликация 3. Контролерът трябва да избере лидер за всеки разрез, опитвайки се оптимално да разпредели лидерите между брокерите.
За всеки разрез контролерът:
- актуализира информацията в Zookeeper относно ISR и лидера;
- изпраща командата LeaderAndISRCommand на всеки брокер, който разполага реплика на този разрез, информирайки брокерите за ISR и лидера.
Когато падне брокер с лидер, Zookeeper изпраща уведомление до контролера, който после избира нов лидер. Отново, контролерът първо актуализира Zookeeper, а след това изпраща команда на всеки брокер, уведомявайки ги за промяната в лидерството.
Всеки лидер е отговорен за набор от ISR. Настройката replica.lag.time.max.ms определя кой ще влезе там. При промяна на ISR лидерът предава нова информация на Zookeeper.
Zookeeper винаги е уведомяван за всякакви промени, за да може в случай на провал ръководството плавно да премине към нов лидер.

Рис. 21. Консенсус Kafka
Протокол на репликация
Разбирането на детайлите на репликацията помага за по-добро разбиране на потенциални сценарии за загуба на данни.
Запитвания за извлечение, Log End Offset (LEO) и Highwater Mark (HW)
Разгледахме, че последователите периодично изпращат на лидера заявки за извличане (fetch). Изходният интервал е 500 мс. Това се различава от RabbitMQ, тъй като при RabbitMQ репликацията се инициира не от огледалото на опашката, а от мастера. Мастерът изпраща измененията на огледалата.
Лидерът и всички последователи запазват смещението на края на лог файла (Log End Offset, LEO) и маркера Highwater (HW). Маркерът LEO съхранява смещението на последното съобщение в локалната реплика, а HW — смещението на последната транзакция. Не забравяйте, че за статус „комит“ съобщението трябва да бъде запазено във всички реплики ISR. Това означава, че LEO обикновено е малко пред HW.
Когато лидерът получи съобщение, той го запазва локално. Последователят прави заявка за извличане, предавайки своя LEO. След това лидерът изпраща пакет от съобщения, започвайки от този LEO, и също предава текущия HW. Когато лидерът получи информация, че всички реплики са запазили съобщението с зададеното смещение, той премества маркера HW. Само лидерът може да премества HW, и така всички последователи научават текущата стойност в отговорите на своята заявка. Това означава, че последователите могат да изостават от лидера както по съобщения, така и относно знанието за HW. Потребителите получават съобщения само до текущия HW.
Обърнете внимание, че „запазено“ (persisted) означава записано в паметта, а не на диск. За производителност, Kafka извършва синхронизация на диска с определен интервал. При RabbitMQ също има такъв интервал, но той ще изпрати потвърждение на публикатора само след като мастерът и всички огледала запишат съобщението на диск. Разработчиците на Kafka, по съображения за производителност, решили да изпращат ack, веднага щом съобщението бъде записано в паметта. Kafka разчита на това, че излишъкът компенсира риска от краткосрочно съхранение на потвърдени съобщения само в паметта.
Провал на лидера
Когато лидерът падне, Zookeeper уведомява контролера, който избира нова реплика на лидера. Новият лидер установява нов маркер HW в съответствие със своя LEO. След това информация за новия лидер получават последователите. В зависимост от версията на Kafka, последователят ще избере един от два сценария:
- Ще отреже локалния лог до известен HW и ще изпрати на новия лидер заявка за съобщения след този маркер.
- Изпраща запитване на лидера, за да узнава HW в момента на избора му за лидер, а след това отсича логовете до това изместване. След това започва да прави периодични запитвания за извадка, започвайки от това изместване.
Последователят може да трябва да отреже логовете по следните причини:
- Когато възникне сблъсък на лидера, първият последовател от набора ISR, регистриран в Zookeeper, печели изборите и става лидер. Всички последователи в ISR, макар и считани за "синхронизирани", може и да не са получили копия на всички съобщения от бившия лидер. Напълно е възможно избраният последовател да няма най-актуалната версия. Kafka гарантира, че няма несъответствия между репликите. Следователно, за да се избегне несъответствие, всеки последовател трябва да отсече своя лог до стойността на HW на новия лидер в момента на избора му. Това е още една причина, поради която настройката acks=all е толкова важна за съгласуваност.
- Съобщенията периодично се записват на диск. Ако всички възли на клъстера се сринат едновременно, на дисковете ще останат реплики с различно изместване. Напълно е възможно, че когато брокерите отново се върнат в мрежата, новият лидер, който ще бъде избран, ще се окаже зад своите последователи, тъй като е записал информацията на диск по-рано от другите.
Присъединяване към клъстера
При присъединяване към клъстера репликите обработват данните по същия начин, както при сблъсък на лидера: проверяват репликата на лидера и отсичат своите логове до неговия HW (в момента на избора). За сравнение, RabbitMQ третира присъединените възли като напълно нови. В двата случая брокерът отхвърля всякакво съществуващо състояние. Ако се използва автоматична синхронизация, мастърът трябва да репликира абсолютно всичкото текущо съдържание в ново огледало по метода "и нека целият свят почака". По време на тази операция мастърът не приема никакви операции за четене или запис. Този подход създава проблеми в големи опашки.
Kafka е разпределен лог, който като цяло съхранява повече съобщения отколкото опашката RabbitMQ, където данните се изтриват от опашката след прочитането им. Активните опашки трябва да останат сравнително малки. Но Kafka е лог с собствена политика на съхранение, която може да установи срок от дни или седмици. Подходът с блокиране на опашката и пълна синхронизация е абсолютно неприемлив за разпределен лог. Вместо това, последователите на Kafka просто прекратяват своя лог до HW лидера (в момента на неговото избиране), ако тяхната копие изостава от лидера. В по-вероятния случай, когато последователят е назад, той просто започва да прави заявки за извличане, започвайки от текущия си LEO.
Нови или присъединили се последователи започват извън ISR и не участват в комитите. Те просто работят в близост до групата, получавайки съобщения толкова бързо, колкото могат, докато не настигнат лидера и не влязат в ISR. Няма блокиране и не е нужно да изхвърляте всичките си данни.
Нарушение на свързаността
Kafka разполага с повече компоненти в сравнение с RabbitMQ, така че тук има по-сложен набор от поведения, когато свързаността в клъстера е нарушена. Но Kafka е проектирана първоначално за клъстери, така че решенията са много добре обмислени.
По-долу са представени няколко сценария на нарушаване на свързаността:
- Сценарий 1. Последователят не вижда лидера, но все още вижда Zookeeper.
- Сценарий 2. Лидерът не вижда нито един последовател, но все още вижда Zookeeper.
- Сценарий 3. Последователят вижда лидера, но не вижда Zookeeper.
- Сценарий 4. Лидерът вижда последователите, но не вижда Zookeeper.
- Сценарий 5. Последователят е напълно изолиран и от другите възли на Kafka, и от Zookeeper.
- Сценарий 6. Лидерът е напълно изолиран и от другите възли на Kafka, и от Zookeeper.
- Сценарий 7. Възелът на контролера на Kafka не вижда друг възел на Kafka.
- Сценарий 8. Контролерът на Kafka не вижда Zookeeper.
За всеки сценарий е предвидено специфично поведение.
Сценарий 1. Последователят не вижда лидера, но все още вижда Zookeeper

Рис. 22. Сценарий 1. ISR от три реплики
Нарушението на свързаността отделя брокер 3 от брокерите 1 и 2, но не и от Zookeeper. Брокер 3 вече не може да изпраща заявки за извличане. След изтичане на време replica.lag.time.max.ms Той се премахва от ISR и не участва в комитите на съобщенията. Веднъж щом свързаността бъде възстановена, той ще възобнови заявките за извличане и ще се присъедини отново към ISR, когато настигне лидера. Zookeeper ще продължи да получава пинги и да счита, че брокерът е жив и здрав.

Рис. 23. Сценарий 1. Брокерът се премахва от ISR, ако не е получен запитване за извличане в рамките на интервала replica.lag.time.max.ms
Няма никакво логично разделение (split-brain) или спиране на възела, както е в RabbitMQ. Вместо това се намалява излишъка.
Сценарий 2. Лидерът не вижда нито един последовател, но все още вижда Zookeeper

Рис. 24. Сценарий 2. Лидер и двама последователи
Нарушаването на мрежовата свързаност отделя лидера от последователите, но брокерът все още вижда Zookeeper. Както в първия сценарий, ISR се свива, но този път само до лидера, тъй като всичките последователи престават да изпращат заявки за извличане. Отново, няма логично разделение. Вместо това се случва загуба на излишъка за новите съобщения, докато свързаността не бъде възстановена. Zookeeper продължава да получава пинги и счита, че брокерът е жив и здрав.

Рис. 25. Сценарий 2. ISR се е свил само до лидера
Сценарий 3. Последовател вижда лидера, но не вижда Zookeeper
Последователят се отделя от Zookeeper, но не и от брокера с лидера. В резултат на това, последователят продължава да прави заявки за извличане и да бъде член на ISR. Zookeeper вече не получава пинги и записва падането на брокера, но тъй като това е само последовател, няма никакви последици след възстановяването.

Рис. 26. Сценарий 3. Последователят продължава да изпраща заявки за извличане на лидера
Сценарий 4. Лидерът вижда последователите, но не вижда Zookeeper

Рис. 27. Сценарий 4. Лидер и двама последователи
Лидерът е отделен от Zookeeper, но не и от брокерите с последователите.

Рис. 28. Сценарий 4. Лидерът е изолиран от Zookeeper
След известно време Zookeeper ще регистрира падането на брокера и ще уведоми контролера. Той ще избере нов лидер сред последователите. Въпреки това, оригиналният лидер ще продължи да мисли, че е лидер и ще продължи да приема записи с acks=1. Последователите вече не му изпращат заявки за извличане, затова той ще ги сметне за мъртви и ще се опита да свие ISR до себе си. Но тъй като няма свързаност с Zookeeper, той няма да може да го направи и в този момент ще се откаже от по-нататъшния прием на записи.
Съобщения acks=all няма да получат потвърждение, защото първо ISR включва всички реплики, а до тях съобщенията не стигат. Когато първоначалният лидер опита да ги премахне от ISR, той няма да може да го направи и изобщо ще спре да получава каквито и да било съобщения.
Клиентите скоро забелязват смяната на лидера и започват да изпращат записи на новия сървър. С веднъж работата на мрежата е възстановена, оригиналният лидер вижда, че вече не е лидер и намалява логовете си до стойността на HW, която е била при новия лидер в момента на срив, за да избегне разминаване на логовете. След това той ще започне да изпраща заявки за избор на новия лидер. Всички записи на оригиналния лидер, ненаправени на новия лидер, се губят. Тоест, ще бъдат загубени съобщения, не потвърдени от първоначалния лидер в тези няколко секунди, когато работеха двама лидери.

Рис. 29. Сценарий 4. Лидерът на брокер 1 става последовател след възстановяване на мрежата
Сценарий 5. Последователят е напълно отделен както от другите узли на Kafka, така и от Zookeeper
Последователят е напълно изолиран както от другите узли на Kafka, така и от Zookeeper. Той просто се премахва от ISR, докато мрежата не бъде възстановена, а след това настигва останалите.

Рис. 30. Сценарий 5. Изолираният последовател се премахва от ISR
Сценарий 6. Лидерът е напълно отделен както от другите узли на Kafka, така и от Zookeeper

Рис. 31. Сценарий 6. Лидер и двама последователи
Лидерът е напълно изолиран от своите последователи, контролера и Zookeeper. За кратък период той ще продължи да получава записи с acks=1.

Рис. 32. Сценарий 6. Изолация на лидера от другите узли на Kafka и Zookeeper
Не получавайки заявки след replica.lag.time.max.ms, той ще опита да компресира ISR до самия себе си, но няма да може да го направи, тъй като няма връзка с Zookeeper, тогава той ще спре да получава записи.
Междувременно Zookeeper ще маркира изолирания брокер като мъртъв, а контролера ще избере нов лидер.

Рис. 33. Сценарий 6. Двама лидери
Първоначалният лидер може да приема записи за няколко секунди, но после спира да получава всякакви съобщения. Клиентите се обновяват на всеки 60 секунди с последните метаданни. Те ще бъдат информирани за смяната на лидера и ще започнат да изпращат записи на новия лидер.

Рис. 34. Сценарий 6. Производителите преминават към новия лидер
Ще бъдат загубени всички потвърдени записи, направени от оригиналния лидер от момента на загуба на свързаност. Както само мрежата бъде възстановена, оригиналният лидер чрез Zookeeper ще открие, че вече не е лидер. След това ще изреже своя лог до HW на новия лидер в момента на избор и ще започне да изпраща заявки като последовател.

Рис. 35. Сценарий 6. Оригиналният лидер става последовател след възстановяване на свързаността на мрежата
В тази ситуация за кратък период може да се наблюдава логическо разделение, но само ако acks=1 и min.insync.replicas също 1. Логическото разделение автоматично приключва или след възстановяване на мрежата, когато оригиналният лидер осъзнае, че вече не е лидер, или когато всички клиенти разберат, че лидерът е сменен и започнат да пишат на новия лидер — в зависимост от това, кое ще се случи по-рано. Във всеки случай ще настъпи загуба на някои съобщения, но само с acks=1.
Има друг вариант на този сценарий, когато непосредствено преди разделението на мрежата последователите изостават, а лидерът е свил ISR до само себе си. След това той се изолира поради загуба на свързаност. Избира се нов лидер, но първоначалният лидер продължава да приема записи, дори acks=all, защото в ISR освен него няма никой друг. Тези записи ще бъдат загубени след възстановяване на мрежата. Единственият начин да се избегне такъв вариант — min.insync.replicas = 2.
Сценарий 7. Kafka контролер не вижда друг Kafka възел
В общи линии, след загуба на свързаност с Kafka възел контролерът няма да може да предаде на него никаква информация за промяната на лидера. В най-лошия случай това ще доведе до краткосрочно логическо разделение, както в сценарий 6. Най-често брокерът просто няма да стане кандидат за лидерство в случай на отказ на последния.
Сценарий 8. Kafka контролер не вижда Zookeeper
От отвалил се контролер Zookeeper няма да получи пинг и ще избере нов Kafka възел за контролер. Оригиналният контролер може да продължи да се представя за такъв, но не получава известия от Zookeeper, затова няма да има никакви задачи за изпълнение. Както само мрежата бъде възстановена, той ще разбере, че вече не е контролер, а е станал обикновен Kafka възел.
Изводи от сценариите
Виждаме, че загубата на свързаност с последователите не води до загуба на съобщения, а просто временно намалява излишъка, докато мрежата не се възстанови. Това, разбира се, може да доведе до загуба на данни, ако се загубят един или повече възли.
Ако поради загуба на свързаност лидерът се е отделил от Zookeeper, това може да доведе до загуба на съобщения с acks=1. Липсата на връзка с Zookeeper предизвиква краткотрайно логическо разделение с двама лидери. Този проблем се решава с параметъра acks=all.
Параметър min.insync.replicas в две или повече реплики предоставя допълнителни гаранции, че такива краткосрочни сценарии няма да доведат до загуба на съобщения, както в сценарий 6.
Обзор на загубата на съобщения
Нека изброим всички начини за загуба на данни в Kafka:
- Всички сривове на лидера, ако съобщенията са били потвърдени с помощта на acks=1
- Всеки нечист (unclean) преход на лидерството, тоест на последовател извън ISR, дори с acks=all
- Изолация на лидера от Zookeeper, ако съобщенията са били потвърдени с помощта на acks=1
- Пълна изолация на лидера, който вече е свил групата ISR само до себе си. Всички съобщения ще бъдат загубени, дори acks=all. Това е вярно само в случай, че min.insync.replicas=1.
- Съвпадение на сривовете на всички възли в дяла. Тъй като съобщенията се потвърдват от паметта, някои може още да не са записани на диск. След рестартиране на сървърите може да липсват някои съобщения.
Нечистите преходи на лидерството могат да бъдат избегнати, или като се забранят, или като се осигури излишък от поне два. Най-устойчива конфигурация е комбинацията от acks=all и min.insync.replicas повече от 1.
Пряко сравнение на надеждността на RabbitMQ и Kafka
За осигуряване на надеждност и висока достъпност и двете платформи реализират система за основна и вторична репликация. Въпреки това, RabbitMQ има ахилесовата си пета. При присъединяване след срив, възлите отхвърлят своите данни, а синхронизацията се блокира. Този двойният удар поставя под въпрос дълготрайността на големите опашки в RabbitMQ. Ще трябва да се примирите или с намаляване на излишъка, или с дълги блокировки. Намаляването на излишъка увеличава риска от масивна загуба на данни. Но ако опашките са малки, то за осигуряване на излишъка с кратки периоди на недостъпност (няколко секунди) може да се справите с помощта на повторни опити за свързване.
В Kafka няма такъв проблем. Тя отхвърля данни само от момента на разминаване между лидера и следователя. Всички общи данни се запазват. Освен това, репликацията не блокира системата. Лидерът продължава да приема записи, докато новият следовател не го настигне, така че за девопсите присъединяването или възстановяването на кластера става тривиална задача. Разбира се, все още остават проблеми, като пропускната способност на мрежата при репликация. Ако се добавят няколко следователи едновременно, може да се сблъскате с ограничението на пропускната способност.
RabbitMQ превъзхожда Kafka в надеждността при едновременен отказ на няколко сървъра в клъстера. Както вече споменахме, RabbitMQ изпраща потвърждение на паблишера само след като съобщението бъде записано на диск от мастера и всичките му огледала. Но това добавя допълнително забавяне по две причини:
- fsync на всеки няколко стотици милисекунди
- Отказът на огледало може да бъде забелязан само след изтичане на времето на живот на пакетите, които проверяват наличността на всеки възел (net tick). Ако огледалото забавя или пада, това добавя забавяне.
Kafka залага на това, че ако съобщението се съхранява на няколко възела, може да се потвърдят съобщенията, след като те попаднат в паметта. Поради това съществува риск от загуба на съобщения от всякакъв вид (дори acks=all, min.insync.реплики=2) в случай на едновременен отказ.
В общи линии, Kafka демонстрира по-висока производителност и е първоначално проектирана за клъстери. Броят на следователите може да бъде увеличен до 11, ако това е необходимо за надеждност. Коэффициентът на репликация 5 и минималният брой реплики в синхронизирано състояние min.insync.replicas=3 ще направят загубата на съобщения много рядко събитие. Ако вашата инфраструктура може да осигури такъв коефициент на репликация и ниво на излишък, можете да изберете тази опция.
Клъстрирането на RabbitMQ е подходящо за малки опашки. Но дори и малките опашки могат бързо да нараснат при голям трафик. Когато опашките станат големи, трябва да се направи строг избор между наличността и надеждността. Клъстрирането на RabbitMQ е най-подходящо за не типични ситуации, където предимствата на гъвкавостта на RabbitMQ надвишават всички недостатъци на неговото клъстриране.
Едно от противодействията на уязвимостта на RabbitMQ относно големите опашки е да ги разделите на множество по-малки. Ако не изисквате пълно подреждане на цялата опашка, а само на съответстващите съобщения (например, съобщения на конкретен клиент), или изобщо да не подреждате, то този вариант е приемлив: погледнете проекта ми за разделяне на опашката (проектът все още е в начален етап).
Накрая, не забравяйте за редица бъгове в механизмите за клъстеризация и репликация както при RabbitMQ, така и при Kafka. С времето системите станаха по-зрели и стабилни, но нито едно съобщение никога не ще бъде на 100% защитено от загуба! Освен това, в дата центровете се случват мащабни аварии!
Ако съм пропуснал нещо, направил съм грешка или не сте съгласни с някой от тезите, не се колебайте да напишете коментар или да се свържете с мен.
Често ме питат: "Какво да избера, Kafka или RabbitMQ?", "Коя платформа е по-добра?". Истината е, че това наистина зависи от вашата ситуация, текущ опит и т.н. Не се наемам да изразя мнение, тъй като би било твърде опростено да се препоръча някаква платформа за всички случаи на употреба и възможни ограничения. Написах този цикъл статии, за да можете да оформите собственото си мнение.
Искам да кажа, че и двете системи са лидери в тази област. Възможно е да съм малко предубеден, защото на базата на опита си с проектите, повече ценя неща като гарантирано подреждане на съобщения и надеждност.
Виждам други технологии, на които им липсва тази надеждност и гарантирано подреждане, след което гледам на RabbitMQ и Kafka — и осъзнавам невероятната стойност на двете системи.
Източник: habr.com
