Kafka и микросервизите: преглед

Kafka и микросервизите: преглед

Здравейте на всички. В тази статия ще обясня защо в Авито преди девет месеца избрахме Kafka и какво представлява тя. Ще споделя един от кейсовете за употреба - брокер на съобщения. Накрая ще поговорим за ползите, които получихме от прилагането на подхода Kafka като услуга.

Проблема

Kafka и микросервизите: преглед

За начало, малко контекст. Преди известно време започнахме да се отклоняваме от монолитната архитектура и в настоящия момент в Авито вече има няколко стотин различни услуги. Те имат свои хранилища, свой технологичен стек и отговарят за своята част от бизнес логиката.

Една от проблемите с голямото количество услуги е комуникацията. Услуга А често иска да научи информация, която притежава услуга Б. В този случай услуга А се обръща към услуга Б чрез синхронен API. Услуга В иска да знае какво става при услуги Г и Д, а те от своя страна се интересуват от услуги А и Б. Когато „любопитните“ услуги стават много, взаимовръзките между тях се превръщат в заплетен клубок.

В същото време, във всеки един момент, услуга А може да стане недостъпна. Какво да прави в този случай услуга Б и всички останали услуги, зависещи от нея? А ако за изпълнението на бизнес операция е необходимо да се извърши верига от последователни синхронни повиквания, вероятността за отказ на цялата операция става още по-висока (и тя нараства с дължината на тази верига).

Избор на технология

Kafka и микросервизите: преглед

Добре, проблемите са ясни. Можем да ги елиминираме, като създадем централизиран система за обмен на съобщения между услугите. Сега на всяка от услугите им е достатъчно да знаят само за тази система за обмен на съобщения. В допълнение, самата система трябва да бъде устойчива на откази и хоризонтално мащабируема, а също така при аварии да съхранява буфер за последваща обработка.

Нека сега изберем технологията, на която ще бъде реализирано доставката на съобщения. За целта първо ще разберем какво очакваме от нея:

  • съобщенията между услугите не трябва да се губят;
  • съобщенията могат да се дублират;
  • съобщенията могат да се съхраняват и четат за период от няколко дни (персистентен буфер);
  • услугите могат да се абонират за интересни за тях данни;
  • няколко услуги могат да четат едни и същи данни;
  • съобщенията могат да съдържат детайлен, голям payload (пренос на състояние, свързано с събития);
  • понякога е необходима гаранция за реда на съобщенията.

Също така беше критично важно за нас да изберем максимално мащабируема и надеждна система с висока пропускна способност (не по-малко от 100k съобщения по няколко килобайта в секунда).

На този етап се разделихме с RabbitMQ (трудно е да бъде поддържан стабилен при високи rps), PGQ от SkyTools (недостатъчно бърз и слабо мащабируем) и NSQ (неперсистентен). Всички тези технологии се използват в компанията ни, но не подхождаха за решаваната задача.

След това започнахме да разглеждаме нови за нас технологии — Apache Kafka, Apache Pulsar и NATS Streaming.

Първи отхвърлихме Pulsar. Решихме, че Kafka и Pulsar са доста сходни решения. И въпреки че Pulsar е проверен от големи компании, по-нов е и предлага по-ниска латентност (в теория), решихме да оставим Kafka от тези две като de facto стандарт за такива задачи. Вероятно ще се върнем към Apache Pulsar в бъдеще.

И така, останаха двама кандидати: NATS Streaming и Apache Kafka. Проучихме подробно и двете решения и и двете подхождаха за задачата. Но в крайна сметка се притеснявахме от относителната младост на NATS Streaming (и от това, че един от основните разработчици, Tyler Treat, реши да напусне проекта и да започне свой собствен — Liftbridge). В същото време режимът на Clustering на NATS Streaming не предоставяше възможност за силно хоризонтално мащабиране (вероятно това вече не е проблем след добавянето на режима на partitioning през 2017 година).

Въпреки това, NATS Streaming е страхотна технология, написана на Go и имаща поддръжка от Cloud Native Computing Foundation. За разлика от Apache Kafka, тя не се нуждае от Zookeeper за работа (възможно е скоро да можем да кажем същото и за Kafka), тъй като вътрешно реализира RAFT. В същото време NATS Streaming е по-лесен за администриране. Не изключваме, че в бъдеще отново ще се върнем към тази технология.

И все пак, днес нашият победител е Apache Kafka. На нашите тестове тя показа, че е достатъчно бърза (над един милион съобщения в секунда на четене и запис при обем на съобщенията от 1 килобайт), достатъчно надеждона, добре мащабируема и проверена на практика от големи компании. Освен това, Kafka поддържа поне няколко големи търговски компании (например, ние ползваме версията на Confluent), а също така Kafka разполага с развита екосистема.

Преглед на Kafka

Преди да започнем, веднага препоръчвам отлична книга — «Kafka: The Definitive Guide» (има и в руския превод, но термините малко объркват мозъка). В него може да се намери информация, необходима за базисно разбиране на Kafka и дори малко повече. Самата документация от Apache и блогът от Confluent също са отлично написани и лесно четими.

И така, нека погледнем на това как е устроена Kafka от височина. Основната топология на Kafka включва producer, consumer, broker и zookeeper.

Брокер

Kafka и микросервизите: преглед

Зa съхранението на вашите данни отговаря брокерът (broker). Всички данни се съхраняват в бинарен вид, а брокерът малко знае за това какво представляват и каква е тяхната структура.

Всеки логически тип събития обикновено се намира в своето отделно топик (topic). Например, събитието за създаване на обява може да попада в топик item.created, а събитието за неговото изменение — в item.changed. Топиците могат да се разглеждат като класификатори на събития. На ниво топик могат да се зададат такива конфигурационни параметри, като:

  • обем на съхраняваните данни и/или тяхната възраст (retention.bytes, retention.ms);
  • фактор на излишък на данни (replication factor);
  • максимален размер на едно съобщение (max.message.bytes);
  • минималното число на съгласувани реплики, при което в топик може да се запишат данни (min.insync.replicas);
  • възможност за извършване на failover на несинхронна изоставаща реплика с потенциална загуба на данни (unclean.leader.election.enable);
  • и още много други (https://kafka.apache.org/documentation/#topicconfigs).

От своя страна, всеки топик се разделя на една или повече партиции (partition). Именно в партицията в крайна сметка попадат събитията. Ако в кластера има повече от един брокер, партициите ще бъдат разпределени равномерно по всички брокери (наколкото е възможно), което ще позволи мащабиране на натоварването за запис и четене в един топик наведнъж на няколко брокера.

На диска данните за всяка партиция се съхраняват във вид на файлове сегменти, по подразбиране равни на един гигабайт (контролира се чрез log.segment.bytes). Важна особеност — изтриването на данни от партиции (при задействане на retention) става точно по сегменти (невъзможно е да се изтрие едно събитие от партиция, може да се изтрие само цял сегмент, и то само неактивен).

Zookeeper

Zookeeper играе ролята на хранилище на метаданни и координатор. Той е този, който може да каже, живи ли са брокерите (може да се види това от гледна точка на zookeeper чрез zookeeper-shell с командата ls /brokers/ids), кой от брокерите е контролерът (get /controller), в синхронно ли са партициите със своите реплики (get /brokers/topics/topic_name/partitions/partition_number/state). Освен това именно към zookeeper първо ще се насочат producer и consumer, за да разберат на кой брокер какви теми и парции се намират. В случаи, когато за тема е зададен replication factor по-голям от 1, zookeeper ще посочи кои партиции са лидери (в тях ще се извършва запис и от тях ще се извършва четене). При падане на брокера именно в zookeeper ще бъде записана информация за новите лидер-партиции (от версия 1.1.0 асинхронно, и това е важно).

В по-старите версии на Kafka zookeeper беше отговорен и за съхранение на офсети, но в сегашно време те се съхраняват в специална тема __consumer_offsets на брокера (въпреки че все още можете да използвате zookeeper за тези цели).

Най-лесният начин да превърнете данните си в тиква всъщност е загуба на информация от zookeeper. В такъв сценарий ще бъде много трудно да разберете какво и откъде трябва да четете.

Producer

Producer — това е най-често сервис, който извършва директно записване на данни в Apache Kafka. Producer избира тема, в която ще бъдат съхранявани неговите тематични съобщения, и започва да записва информация в нея. Например, producer може да бъде сервис за обяви. В такъв случай той ще изпраща в тематични теми събития като „обява създадена“, „обява актуализирана“, „обява изтрита“ и т.н. Всяко събитие в този случай представлява двойка ключ-стойност.

По подразбиране всички събития се разпределят между партициите на темата по метода round-robin, ако ключ не е зададен (загубвайки подреденост), и чрез MurmurHash (ключ), ако ключът присъства (подреденост в рамките на една партиция).

Тук веднага трябва да се отбележи, че Kafka гарантира реда на събитията само в рамките на една партиция. Но на практика често това не е проблем. Например, можете гарантирано да добавяте всички изменения на едно и също обявление в една партиция (по този начин запазвайки реда на тези изменения в рамките на обявлението). Също така може да се предава последователен номер в едно от полетата на събитието.

Consumer

Kafka и микросервизите: преглед

Consumer е отговорен за получаването на данни от Apache Kafka. Ако се върнем към примера по-горе, consumer-ът може да бъде услуга за модериране. Тази услуга ще бъде абонирана за темата на услугата за обяви и при появата на нова обява ще я получи и анализира за съответствие с определени зададени политики.

Apache Kafka запомня последните събития, които consumer-ът е получил (за това се използва служебна тема __consumer__offsets), като по този начин гарантира, че при успешно четене consumer-ът няма да получи същото съобщение два пъти. Въпреки това, ако се използва опцията enable.auto.commit = true и напълно се остави работата по следене на положението на consumer-а в темата на изцяло на Kahka, може да загубите данни. В производствен код, обикновено позицията на консумиращия се контролира ръчно (разработчикът управлява момента, в който задължително трябва да се извърши commit на прочетеното събитие).

В случаите, когато един consumer не е достатъчен (например, потокът от нови събития е много голям), могат да се добавят още няколко consumer-и, свързвайки ги в consumer group. Consumer group логически представлява точно такъв самия consumer, но с разпределение на данни между участниците в групата. Това позволява на всеки от участниците да вземе своя дял от съобщенията, като по този начин се увеличава скоростта на четене.

Резултати от теста

Kafka и микросервизите: преглед

Няма да пиша много пояснителен текст тук, просто ще споделя получените резултати. Тестовете бяха проведени на 3 физически машини (12 CPU, 384GB RAM, 15k SAS DISK, 10GBit/s Net), брокерите и zookeeper бяха разположени в lxc.

Тестиране на производителността

В хода на тестовете бяха получени следните резултати.

  • Скоростта на запис на съобщения с размер 1KB едновременно от 9 producer-а е 1300000 събития в секунда.
  • Скоростта на четене на съобщения с размер 1KB едновременно от 9 consumer-а е 1500000 събития в секунда.

Тест за отказоустойчивост

В хода на тестовете бяха получени следните резултати (3 брокера, 3 zookeeper).

  • Неочаквано прекратяване на един от брокерите не води до спиране или недостъпност на клъстера. Работата продължава в нормален режим, но оставащите брокери поемат по-голямо натоварване.
  • Неочакваното спиране на двама брокера в случай на клъстер от три брокера и min.isr = 2 води до недостъпност на клъстера за запис, но остава достъпен за четене. В случай, че min.isr = 1, клъстерът продължава да бъде достъпен както за четене, така и за запис. Въпреки това, този режим противоречи на изискването за висока надеждност на данните.
  • Неочакваното спиране на един от сървърите Zookeeper не води до спиране или недостъпност на клъстера. Работата продължава нормално.
  • Неочакваното спиране на двама сървъри Zookeeper води до недостъпност на клъстера, докато поне един от сървърите Zookeeper не бъде възстановен. Това твърдение е валидно за клъстер Zookeeper от 3 сървъра. След проучвания беше решено да се увеличи клъстерът Zookeeper до 5 сървъра за повишаване на устойчивостта.

Kafka като услуга

Kafka и микросервизите: преглед

Убедихме се, че Kafka е отлична технология, която решава задачата, поставена пред нас (реализация на брокер за съобщения). Въпреки това, решихме да забраним на услугите да се свързват директно с Kafka и я затворихме зад услугата data-bus. Защо направихме това? Наистина има няколко причини.

  • Data-bus пое всички задачи, свързани с интеграцията с Kafka (реализация и настройка на consumer-и и producer-и, мониторинг, алармиране, логиране, мащабиране и т.н.). По този начин интеграцията с брокера за съобщения е максимално опростена.

  • Data-bus ни позволи да се абстрахираме от конкретен език или библиотека за работа с Kafka.

  • Data-bus позволи на другите услуги да се абстрахират от слоя за съхранение. Може би в някакъв момент ще сменим Kafka с Pulsar и при това никой няма да забележи (всички услуги знаят само за API data-bus).

  • Data-bus пое валидирането на схемите на събитията.

  • С помощта на data-bus е реализирана автентикация.

  • Под прикритие на data-bus можем без downtime и незабелязано да актуализираме версиите на Kafka, централизирано да управляваме конфигурациите на producer-и, consumer-и, брокери и т.н.

  • Data-bus ни позволи да добавяме нужните функции, които липсват в Kafka (като например одит на теми, контрол над аномалиите в клъстера, създаване на DLQ и т.н.).

  • Data-bus позволява реализирането на failover централизирано за всички услуги.

В момента, за да започнете да изпращате събития в съобщителния брокер, е достатъчно да свържете малка библиотека в кода на услугата си. Това е всичко. Имате възможност да пишете, четете и мащабирате с една линия код. Цялата реализация е скрита от вас, навън стърчат само няколко параметъра като размер на партидата. Под капака услугата data-bus стартира необходимия брой инстанции на producer-и и consumer-и в Kubernetes и им предоставя необходимата конфигурация, но всичко това е прозрачно за вашата услуга.

Разбира се, няма универсално решение, и подобен подход има свои ограничения.

  • Data-bus трябва да се поддържа самостоятелно, за разлика от външните библиотеки.
  • Data-bus увеличава броя на взаимодействията между услугите и съобщителния брокер, което води до намаляване на производителността в сравнение с голата Kafka.
  • Не всичко може да се скрие толкова лесно от услугите, не искаме да дублираме функционалността на KSQL или Kafka Streams в data-bus, затова понякога трябва да позволим на услугите да се свързват директно.

В нашия случай предимствата надвиха недостатъците и решението да скрием съобщителния брокер зад отделна услуга се оправда. През годината на експлоатация нямаме сериозни инциденти и проблеми.

P.S. Благодаря на моята приятелка, Екатерина Обаляева, за страхотните илюстрации към тази статия. Ако са ви харесали, тук ще намерите още илюстрации.

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster