В очакване на старта на нов поток по курса подготивихме превод на интересен материал.

Преглед
Ще обсъдим доста популярен модел, с който приложенията използват няколко хранилища за данни, където всяко хранилище служи за определена цел, например, за съхранение на каноничната форма на данни (MySQL и т.н.), осигуряване на разширени възможности за търсене (ElasticSearch и т.н.), кеширане (Memcached и т.н.) и други. Обикновено при използване на няколко хранилища за данни едно от тях работи като основно хранилище, а останалите като производни хранилища. Единствената пречка е как да се синхронизират тези хранилища за данни.
Разгледахме редица различни модели, които се опитват да решат проблема с синхронизацията на няколко хранилища, като двойна запись, разпределени транзакции и т.н. Въпреки това, тези подходи имат съществени ограничения по отношение на реалната употреба, надеждността и техническото обслужване. Освен синхронизацията на данните, на някои приложения също е необходимо да се обогатяват данните, извиквайки външни услуги.
За решаване на тези проблеми е създаден Delta. Delta в крайна сметка представлява координирана, управлявана от събития платформа за синхронизация и обогатяване на данни.
Съществуващи решения
Двойна запис
За да синхронизирате две хранилища за данни, може да се използва двойна запис, която извършва записа в едно хранилище, а след това веднага го записва и в друго. Първият запис може да бъде повторен, а вторият – прекратен, ако първият завърши неуспешно след изчерпване на опитите. Обаче, двете хранилища за данни могат да спрат да се синхронизират, ако записът във второто хранилище завърши неуспешно. Този проблем обикновено се решава чрез създаване на процедура за възстановяване, която периодично може да прехвърля данни от първото хранилище във второто или да го прави само ако бъдат открити различия в данните.
Проблеми:
Извършването на процедурата за възстановяване е специфична работа, която не може да бъде повторно използвана. Освен това, данните между хранилищата остават несинхронизирани, докато не се извърши процедурата за възстановяване. Решението става по-сложно, ако се използват повече от две хранилища за данни. И накрая, процедурата за възстановяване може да увеличи натоварването на изходния източник на данни.
Таблица на логовете за промени
Когато в набор от таблици настъпват промени (например, вмъкване, актуализиране и изтриване на записи), записите за промените се добавят в таблицата на логовете като част от същата транзакция. Друг поток или процес постоянно запитва събития от таблицата на логовете и ги записва в едно или повече хранилища за данни, като при нужда премахва събитията от таблицата на логовете след потвърждение на записа от всички хранилища.
Проблеми:
Този шаблон трябва да бъде реализиран като библиотека и в идеалния случай без да се променя кодът на приложението, което я използва. В многоезична среда реализацията на такава библиотека трябва да съществува на всеки необходим език, но осигуряването на последователност в работата на функциите и поведението между езиците е много трудно.
Друг проблем е свързан с получаването на промени в схемата в системи, които не поддържат транзакционни промени в схемата [1][2], като например MySQL. Затова шаблонът за извършване на промяна (например, промяна на схемата) и транзакционното записване на това в таблицата на логовете за промени не винаги ще работи.
Разпределени Транзакции
Разпределените транзакции могат да се използват за разделяне на транзакция между няколко хетерогенни хранилища за данни, така че операцията или да бъде зафиксирована във всички използвани хранилища, или да не бъде зафиксирована в нито едно от тях.
Проблеми:
Разпределените транзакции представляват много голям проблем за хетерогенни хранилища данни. По своята природа те могат да разчитат само на най-ниския общ знаменател на участващите системи. Например, XA транзакциите блокират изпълнението, ако в процеса на приложението се случи неуспех по време на подготовката. Освен това, XA не осигурява откритие на задръствания и не поддържа оптимистични схеми за управление на паралелизъм. В допълнение, някои системи като ElasticSearch не поддържат XA или някаква друга хетерогенна модел на транзакции. Така че осигуряването на атомарност на записа в различни технологии за съхранение на данни остава много сложна задача за приложенията [3].
Delta
Delta беше разработен за премахване на ограниченията на съществуващите решения за синхронизация на данни и позволява също така обогатяване на данните в реално време. Нашата цел беше да абстрахираме всички тези сложни моменти от разработчиците на приложения, за да могат да се фокусират изцяло върху реализирането на бизнес функционалността. В следващите редове ще опишем 'Movie Search', реалния случай на приложение на Delta от Netflix.
В Netflix широко се прилага микросервисна архитектура и всеки микросервис обикновено обслужва един тип данни. Основната информация за филма е в микросервиза, наречен Movie Service, а свързаните с него данни, като информация за продуценти, актьори, доставчици и т.н., се управляват от няколко други микросервиза (включително Deal Service, Talent Service и Vendor Service).
Бизнес потребителите в Netflix Studios често имат нужда от търсене по различни критерии на филми, затова е много важно те да имат възможност да търсят във всички данни, свързани с филмите.
Преди появата на Delta, екипът за търсене на филми трябваше да получава данни от няколко микросервиза, преди да индексира данните за филми. Освен това, екипът трябваше да разработи система, която периодично обновяваше търсещия индекс, запитвайки за промени от други микросервизи, дори когато нямаше реални промени. Тази система бързо стана сложна и трудно поддържана.

Фигура 1. Системата за polling преди Delta
След като започнете да използвате Delta, системата е опростена до система, управлявана от събития, както е показано на следващата рисунка. Събитията CDC (Change-Data-Capture) се изпращат в теми на Keystone Kafka чрез Delta-Connector. Приложението Delta, изградена с помощта на Delta Stream Processing Framework (базирано на Flink), получава CDC-събития от темата, обогатява ги, извиквайки други микросервизи, и накрая предава обогатените данни в индекса за търсене в Elasticsearch. Целият процес става почти в реално време; т.е. веднага щом изменението бъде записано в хранилището за данни, индексите за търсене се обновяват.

Рисунка 2. Пайплайн на данни при използване на Delta
В следващите раздели ще опишем работата на Delta-Connector, който се свързва с хранилището и публикува CDC-събития на транспортно ниво, представляващо инфраструктура за предаване на данни в реално време, насочваща CDC-събития в теми на Kafka. А в самия край ще говорим за структурата на обработка на потоците Delta, която разработчиците на приложения могат да използват за логика на обработка и обогатяване на данни.
CDC (Change-Data-Capture)
Създадохме CDC-сервис, наречен Delta-Connector, който може да фиксира ангажирани промени от хранилището на данни в реално време и да ги записва в поток. Промените в реално време се вземат от дневника на транзакциите и дамповете на хранилището. Дамповете се използват, тъй като дневниците на транзакциите обикновено не съ хранят цялата история на промените. Промените обикновено се сериализират като събития Delta, така че получателят не трябва да се притеснява откъде идва промяната.
Delta-Connector поддържа няколко допълнителни функции, като:
- Възможност да пишете в персонализирани изходни данни, без да използвате Kafka.
- Възможност за активиране на ръчни дампове по всяко време за всички таблици, за определена таблица или за конкретни основни ключове.
- Дамповете могат да се вземат на части, така че не е необходимо да започвате отначало в случай на срив.
- Не е необходимо да поставяте блокировки на таблиците, което е много важно, за да се избегне блокиране на трафика при запис в базата данни от нашия сервис.
- Висока наличност благодарение на резервни инстанции в AWS Availability Zones.
Сега поддържаме MySQL и Postgres, включително при разгръщане в AWS RDS и Aurora. Също така поддържаме Cassandra (multi-master). Повече подробности за Delta-Connector можете да научите в това .
Kafka и транспортният слой
Транспортният слой на събитията Delta е построен на услугата за обмяна на съобщения на платформата .
Исторически, публикуването на съобщения в Netflix е оптимизирано с оглед на повишаване на достъпността, а не на дълготрайността (виж ). Компромисът е потенциално несъответствие на данни в брокера при различни крайни сценарии. Например, unclean leader election отговаря за това, че получателят потенциално дублира или губи събития.
С Delta искахме да постигнем по-сериозни гаранции за дълготрайност, за да осигурим доставка на CDC-събития в производни хранилища. Затова предложихме специално проектиран Kafka клъстер като обект от първостепенна важност. Можете да видите някои настройки на брокера в таблицата по-долу:

В клъстерите Keystone Kafka, unclean leader election обикновено е включен за осигуряване на достъпност на издателя. Това може да доведе до загуба на съобщения в случай, че несинхронизирана реплика бъде избрана за лидер. За новия високо надежден Kafka клъстер, параметърът unclean leader election е изключен, за да се предотврати загубата на съобщения.
Също така увеличихме replication factor от 2 на 3 и minimum insync replicas от 1 на 2. Издателите, пишещи в този клъстер, изискват acks от всички останали, гарантирайки, че 2 от 3 реплики ще имат най-актуалните съобщения, изпратени от издателя.
Когато инстанцията на брокера приключи, нова инстанция замества старата. Въпреки това новият брокер трябва да навакса несинхронизираните реплики, което може да отнеме няколко часа. За да намалим времето за възстановяване при този сценарий, започнахме да използваме блочно хранилище (Amazon Elastic Block Store) вместо локалните дискове на брокерите. Когато новата инстанция заменя приключилата инстанция на брокера, тя прикрепя EBS тома, който е имала приключилата инстанция, и започва да наваксва новите съобщения. Този процес съкращава времето за ликвидиране на задълженията от няколко часа до няколко минути, тъй като новата инстанция не трябва вече да репликира от празно състояние. В цялост, отделните жизнени цикли на хранилището и брокера значително намаляват въздействието от смяната на брокера.
За да увеличим още повече гаранцията за доставка на данни, използвахме за откриване на всяка загуба на съобщения в екстремни условия (например, разсинхронизация на часовниците на лидера на раздела).
Stream Processing Framework
Нивото на обработка в Delta е изградена на базата на платформата Netflix SPaaS, която осигурява интеграция на Apache Flink с екосистемата на Netflix. Платформата предоставя потребителски интерфейс, който управлява разгръщането на задания на Flink и оркестрацията на клъстери Flink върху нашата платформа за управление на контейнери Titus. Интерфейсът също управлява конфигурациите на задачите и позволява на потребителите да правят промени в конфигурацията динамично, без необходимост от повторно компилиране на задачите Flink.
Delta предоставя фреймворк за обработка на потоци (stream processing framework) за данни на базата на Flink и SPaaS, който използва система, основана на анотации DSL (Domain Specific Language), за да абстрахира техническите детайли. Например, за да определи стъпката, с която ще бъдат обогатявани събитията, извиквайки външни услуги, потребителите трябва да напишат следния DSL, а фреймворкът ще създаде на неговата основа модел, който ще се изпълни от Flink.

Фигура 3. Пример за обогатяване на DSL в Delta
Фреймворкът за обработка не само съкращава кривата на обучение, но и осигурява общи функции за обработка на потоци, като дедупликация, схематизиране, както и гъвкавост и отказоустойчивост за решаване на общи проблеми в работата.
Delta Stream Processing Framework се състои от два основни модула: модулът DSL & API и модулът Runtime. Модулът DSL & API предоставя DSL и UDF (User-Defined-Function) API, което позволява на потребителите да напишат собствена логика за обработка (например филтриране или трансформации). Модулът Runtime осигурява реализация на парсера DSL, който изгражда вътрешно представяне на стъпките за обработка в модели DAG. Компонентът Execution интерпретира моделите DAG, за да инициализира фактическите оператори Flink и в крайна сметка да стартира приложението Flink. Архитектурата на фреймворка е илюстрирана на следния рисунък.

Рисунок 4. Архитектура на Delta Stream Processing Framework
Този подход има няколко предимства:
- Потребителите могат да се фокусират върху своята бизнес логика, без да е необходимо да се задълбочават в спецификата на Flink или структурата на SPaaS.
- Оптимизацията може да се извършва по прозрачен за потребителите начин, а грешките могат да бъдат поправяни без необходимост от извършване на промени в кода на потребителя (UDF).
- Работата с приложенията Delta е улеснена за потребителите, тъй като платформата предлага гъвкавост и устойчивост на откази от самото начало и събира множество подробни метрики, които могат да се използват за известия.
Използване в продукция
Delta работи в продукция от над година и играе ключова роля в много приложения на Netflix Studio. Тя помогна на екипите да реализират случаи на употреба като индексиране на търсене, съхранение на данни и работни процеси, управлявани от събития. По-долу е представен преглед на архитектурата на платформата Delta на високо ниво.

Рисунок 5. Архитектура на Delta на високо ниво.
Благодарности
Бихме искали да благодарим на следните личности, които участваха в създаването и развитието на Delta в Netflix: Allen Wang, Charles Zhao, Jaebin Yoon, Josh Snyder, Kasturi Chatterjee, Mark Cho, Olof Johansson, Piyush Goyal, Prashanth Ramdas, Raghuram Onti Srinivasan, Sandeep Gupta, Steven Wu, Tharanga Gamaethige, Yun Wang и Zhenzhong Xu.
Източници
- Martin Kleppmann, Alastair R. Beresford, Boerge Svingen: Онлайн обработка на събития. Commun. ACM 62(5): 43–49 (2019). DOI:
: «Data Build Tool за хранилище Amazon Redshift».
Източник: habr.com
