Distributed Systems Consensus Vocabulary: Raft, Paxos, CAP, and Linearizability Explained (англійською)
Освоєння словникового запасу розподіленого консенсусу для старших інженерів: вибори лідерів Raft, кворум, лінійність проти серіалізованості, теореми CAP і PACELC, split- brain, векторні годинники, 2PC і saga.
Розподілений консенсус є однією з найбільш технічно вкрай складних областей інженерії бекенду і платформи - і однією з найбільш багатих на словниковий запас. Коли ваша команда обговорює «чи потрібна нам тут лінійність, чи достатньо причинної послідовності?» або «чи можемо ми терпіти розділений мозок, якщо ми використовуємо фехтування лідерів?», Точність має значення. Неправильное слово стоит вам часов в обзоре дизайна.
У цій статті наведено словниковий запас, який вам знадобиться для участі у обговореннях щодо проектування розподілених систем, успішного проходження інтерв’ ю з керівниками проектів з підтримки та інфраструктури, а також для розуміння таких документів, як Raft і Paxos.
Прикладом цієї проблеми є проблема мови — мова є мовною проблемою
** Розподілений консенсус ** Проблема отримання згоди декількох вузлів мережі щодо єдиного значення, навіть у разі невдач вузлів і затримок мережі. Розв’язання за допомогою консенсусних протоколів, таких як Raft і Paxos.
“Розподілений консенсус потрібний, коли декілька реплік мають погодитися з порядком записів — інакше різні репліки можуть бачити різні історії.”
**Допускається помилка **
Вміння розподіленої системи продовжувати працювати правильно, коли деякі вузли не працюють. Система, яка терпить f невдач потребує принаймні 2f + 1 вузлів для консенсусу Raft-style.
“Наш кластер Raft має 5 вузлів — він може терпіти 2 одночасні невдачі вузлів без втрати можливості вибору лідерів або виконання записів.”
Кворум Більшість вузлів, які мають погодитися на продовження операції. У кластері з 5 вузлів кворум = 3. Дії, які не досягають кворуму, блокуються до тих пір, поки не буде достатньої кількості доступних вузлів.
- “Ми втратили кворум, коли другий вузол зламався — кластер перестав приймати записи, поки ми не повернули третій вузол в мережу.” *
** Мережевий розділ ** Неможливість зв’ язку вузлів один з одним — кластер розділено на дві або більше ізольованих груп (розділів). Мережеві розділи є найскладнішим сценарієм збою в розподілених системах.
- “Мережний розділ розділив наш кластер з 5 вузлів на групу з 3 вузлів і групу з 2 вузлів. Розділ з 3 вузлами мав кворум і продовжував працювати. 2-вузловий розділ правильно відмовився від записів.”*
Розділ 1: Raft — Словник сучасного алгоритму консенсусу
Raft — це алгоритм консенсусу, який використовується в etcd (Kubernetes), CockroachDB, TiKV, і багатьох інших сучасних розподілених системах. Зрозуміти словник Raft є важливим для інженерів платформи.
Лідер, Послідовник, Кандидат Три стани, у яких може перебувати вузол Raft. Лідер обробляє всі запити клієнтів і реплікує записи до підписників. Коли лідер зазнає невдачі, послідовники стають кандидатами і починають вибори. Кандидат, який отримає голоси від кворуму, стає новим лідером.
- “Вузол лідерів зазнав аварії. Через 150 мс три послідовники перевищили свій тайм-аут, оголосили себе кандидатами і почали просити голосів. Вузол з найбільш актуальним журналом виграв вибори.»*
Відповідний термін Число, що монотонна зростає, що визначає період керування. Кожні вибори починають новий термін. Терміни використовуються для виявлення застарілих повідомлень — вузол, який отримує повідомлення з більшим терміном, знає, що він пропустив вибори і оновлює себе.
“Лідер був у 12-му класі. Коли він впав і нові вибори завершилися, новий лідер був на 13-му терміні. Ми можемо побачити це в журналах etcd.”
** Журнал реплікації ** Процес, за допомогою якого лідер надсилає записи журналу (команди клієнта) всім підлеглим. Запис вважається ** переданим ** тільки після того, як лідер реплікував його до кворуму послідовників.
- “Реплікація журналу була повільною — лідер чекав на повільного підлеглого в іншій зоні доступності, щоб той підтвердив кожен запис перед перенесенням. Ми вирішили це, скинувши повільний наслідувач нижче кворуму.”*
Біт серця Періодичне повідомлення, надіслане лідером всім підлеглим, щоб запобігти їхньому початку виборів. Якщо послідовник не отримує серцевий ритм протягом часу виборів, він вважає, що лідер помер і починає вибори.
Розділ 2: Paxos і ZooKeeper
Пакс Оригінальний алгоритм розподіленого консенсусу, описаний Leslie Lamport. Відомий тим, що його важко зрозуміти і реалізувати правильно. Більшість реальних реалізацій використовують Multi-Paxos (оптимізація для реплікованих станових машин) або переходять на Raft.
- “ZooKeeper використовує Zab, варіант Paxos. Kafka покладається на ZooKeeper для координації брокерів — саме тому кластери Kafka мали ZooKeeper як залежність так довго.”*
** Zab (Атомна трансляція ZooKeeper) ** Протокол консенсусу, який використовується Apache ZooKeeper. Створено для первинної резервної копії у системах, таких як Kafka broker coordination. Концептуально схожий на Paxos.
Розділ 3: Консистенція моделей
** Лінійність (сильна послідовність) ** Найсильніша модель послідовності. Кожне читання повертає останній запис, і всі дії виконуються миттєво у одному місці у часі. Лінійність є дорогою — вона вимагає координації між репліками на кожній операції.
“типово, etcd забезпечує лінійні читання — кожне читання проходить через лідер, щоб переконатися, що воно відображає всі записані дії. Це додає затримку, але гарантує коректність для рішень контролера Kubernetes. ”
** Серіалізація ** Модель послідовності для транзакцій (кількох операцій). Трансакции виконуються одна за одною (якби вони були серіалізовані), навіть якщо вони виконуються одночасно. Відрізняється від лінійності, яка застосовується до окремих операцій.
- “База даних надає можливість виконання операцій, які можна послідовно виконувати — одночасні операції розкладаються так, щоб результати були еквівалентними певному порядку виконання операцій. Це запобігає запису і фантомному читання.»*
** Причинна послідовність ** Операції, які є причинно- наслідковими, з’ являються у тому ж порядку для всіх вузлів. Операції без причинно- наслідкового зв’ язку можуть бути розглянуті в різних порядках різними вузлами. Слабша, ніж лінійність, але значно дешевша.
- “Ми використовуємо причинну послідовність для гілки відповідей на коментарі — відповіді завжди з’ являються після коментаря, на який вони посилаються, але два непов’ язаних коментарі можуть з’ являтися в різних порядках для різних користувачів.” *
Відповідність можливості Якщо буде надано достатньо часу без подальших оновлень, всі репліки збігатимуться до одного значення. Найслабша корисна модель послідовності — не надає гарантій про те, коли відбувається збіг або як виглядають проміжні стани.
- “Корзинка використовує можливу послідовність — якщо ви додасте елемент до вашого телефону, коли не буде доступу до мережі WiFi, його буде синхронізовано, коли з’ єднання буде відновлено. Ми обробляємо конфлікти з політикою останнього запису-переможця. “*
Розділ 4: CAP і PACELC теореми
** Теорема КАП ** Фундаментальна теорема (Brewer, 2000): розподілена система може гарантувати не більше двох з трьох властивостей одночасно:
- ** Послідовність: ** кожне читання отримує найновіший запис
- ** Доступність: ** кожен запит отримує відповідь (не помилку)
- ** Допуск розділів: ** система продовжує працювати, незважаючи на мережеві розділи
На практиці, мережеві розділи є неминучими — тому справжній вибір полягає в тому, чи надавати перевагу послідовності (CP) або доступності (AP) під час розділу.
- “Cassandra є системою з доступом до мережі — під час розділу вона визначає пріоритет доступності і приймає записи з обох сторін розділу. Ви отримуєте остаточну послідовність, але жоден мережевий розділ не може зробити його відмовитися від записів.”*
** Теорема ПАСЕЛКА ** Розширення CAP, створене Daniel Abadi: навіть якщо не існує розділу (P), розподілена система повинна обирати між Затримкою (L) і Послідовністю (C). Це захоплює реальний компроміс, який працює постійно, а не тільки під час невдач.
- “DynamoDB є PA/ EL — доступним для розділів, і компроміс у звичайній роботі — це низька затримка за рахунок сильної послідовності. Саме тому DynamoDB рекомендує послідовне читання за замовчуванням для програм, чутливих до затримки. *
** КП проти Системи АП
- ** CP systems ** (Consistent + Partition- tolerant): надає перевагу зникненню доступу до даних, ніж обслуговуванню застарілих даних. Приклади: HBase, ZooKeeper, тощо, Spanner.
- ** AP systems ** (Available + Partition- tolerant): надає перевагу продовженню обслуговування запитів навіть з потенційно застарілими даними. Приклади: Cassandra, DynamoDB (можливо), CouchDB.
Розділ 5: Розділення мозку і фехтування
Раздробленная голова Сценарій, де мережевий розділ спричиняє дві частини кластера, щоб кожна вважала, що вони є активним лідером / основним - обидві сторони приймають записи незалежно, що призводить до розбіжностей даних, які можуть бути неможливими для автоматичного примирення.
- “Подія split- brain призвела до того, що дві основні репліки бази даних приймали конфліктні записи протягом 4 хвилин. Примирення розбіжного стану вимагало вручну втручання і певної втрати даних.”*
Лідер Фехтування Механізм, щоб запобігти колишньому лідеру приймати рішення після того, як він був скинутий з престолу - забезпечуючи, що в той же час є тільки один авторитетний лідер. Реалізовано за допомогою токенів епохи, замків на основі оренди або STONITH (Shoot The Other Node In The Head).
“Ми використовуємо фехтування на основі епох — кожен новий лідер має вищий номер епохи. Стародавні повідомлення від скинутого лідера відкидаються всіма вузлами.»
СТОНИТ (Встреляй другого узла в голову) Метод обмеження, який примушує вузол, який, як вважається, поводиться неправильно, припинити роботу — відключивши живлення, відкликавши доступ до мережі або надіславши команду віддаленого вимкнення. Брутально, но надежно.
“Наш кластер використовує STONITH через IPMI — якщо вузол не відповідає на агента огорожі, ми відключаємо живлення цього сервера через інтерфейс управління поза діапазоном.”
Розділ 6: Години в розподілених системах
Время в Лампорте Логічний годинник, який забезпечує часткове впорядкування подій у розподіленій системі. Кожен вузол збільшує свій лічильник за кожною подією; під час отримання повідомлення він встановлює свій лічильник на значення max( local, received) + 1. Не вловлює причинно-наслідкову зв’язок повністю.
** Векторний годинник ** Розширення часових штампів Lamport, яке захоплює причинно- наслідкові зв’ язки. Кожен вузол підтримує вектор лічильників для всіх вузлів. Можна виявити одночасні (не пов’ язані між собою) і причинно- наслідкові події (одна з яких сталася раніше за іншу). Використовується у розподілених базах даних для виявлення конфліктів.
- “ДинамоDB спочатку використовував векторні годинники для виявлення конфліктних записів. Пізніше вони перейшли на простіший підхід last-write-wins для масштабованості.”*
** Гібридний логічний годинник (HLC) ** Об’ єднує фізичний час і логічний час — залишається близьким до часу на стіні (для зручності читання і сумісності з SQL) при збереженні властивостей впорядкування Lamport. Використовується CockroachDB і YugabyteDB для глобально розподіленого SQL.
“CockroachDB використовує HLCs, тому що час транзакції близький до реального часу — це дозволяє географічно розподіленим транзакціям читати дані з точним до мілісекунди часовим обмеженням.”
Розділ 7: Розподілені транзакції
** Двофазний затвердження (2PC) ** Протокол розподіленої транзакції, у якому координатор запитує у всіх учасників про голосування (фаза підготовки), і якщо всі проголосують « так », він затверджує транзакцію на всіх вузлах (фаза затвердження). Надає гарантії ACID для декількох вузлів. Головна слабкість: якщо координатор зазнає аварії після підготовки, але перед затвердженням, учасники будуть заблоковані на неопределенный термін (« проблема блокування »).
- “Ми використовуємо 2PC для операцій запису між осколками у нашій системі платежу — це забезпечує, що або всі осколки здійснюють платіж, або всі відкидають його. Ризик невдачі координатора зменшується за рахунок зайвого координатора з виборами лідера на основі Raft. ”*
** Saga (як альтернатива розподіленій транзакції) ** Транзакція з тривалим часом виконання розділена на послідовність локальних транзакцій, кожна з яких має компенсаційну транзакцію, яка скасує її дію, якщо пізніший крок зазнає невдачі. Уникає проблеми блокування 2PC, але не забезпечує повної ізоляції ACID.
“Ми використовуємо Sagas для потоку виконання замовлень через 5 служб — кожен крок виконується локально і видає подію. Якщо служба доставки не спрацює, компенсуючі транзакції відновлюють резервування складу і повертають плату. “
Розмовна мова
| Context | Phrase |
|---|---|
| Recommending consistency level | ”For the leader election use case, we need linearizability — strong consistency is non-negotiable. For the analytics read replica, eventual consistency is fine.” |
| Explaining a failure scenario | ”If we have a 3-node consensus cluster and two nodes go down simultaneously, we lose quorum. The surviving node will refuse writes until quorum is restored.” |
| CAP trade-off | ”Given our SLA requires 99.99% availability, we should build an AP system and handle conflict resolution at the application layer. A CP system will reject writes during any partition.” |
| Explaining split-brain risk | ”Without a fencing mechanism, a slow leader could rejoin the cluster after a partition and attempt to write — resulting in split-brain. Leader fencing via epoch tokens prevents this.” |
Зв’язані ресурси
- Програмне забезпечення для архітектури — словник системного дизайну
- Розподілені системи словникової практики — вправи з словниковим запасом
- Технічний інтерв’ю System Design англійською — використовуючи цей словник в інтерв’ю