Консенсус: Paxos, Raft, выбор лидера, репликация лога
Есть ровно одна задача, к которой сводится большинство «сложных» вещей в распределённых системах: заставить несколько машин согласиться на одном значении и никогда не передумать. Кто лидер шарда. Какая конфигурация кластера актуальна. Какая транзакция закоммичена, а какая нет. Кто держит распределённую блокировку. Какой номер владения секцией считается текущим. Все эти вопросы — один и тот же вопрос, и на него есть один класс ответов, называемый консенсусом.
Соблазн решить его «на коленке» огромен, и почти каждая команда однажды пробует: поставить флаг в Redis, взять таймаут побольше, выбрать лидера по наименьшему IP. Все такие решения работают на демо и ломаются в проде одинаково — тихо, с двумя лидерами и с потерянными подтверждёнными записями. Причина не в криворукости: у задачи есть доказанная нижняя граница сложности, и любое решение, которое выглядит проще Paxos, либо эквивалентно Paxos, либо неверно. Формулировку любят приписывать Лампорту, и она обидно точна.
Дальше предполагается, что вы уже знаете, что отказ узла неотличим от медленной сети, что физическим часам доверять нельзя, что линеаризуемость — это про порядок, а не про свежесть, и что кворум — это способ гарантировать пересечение. Здесь мы собираем перечисленное в работающий механизм и разбираем, где он ломается.
1. Что именно требуется от консенсуса
Пусть каждый из N процессов предлагает значение. Алгоритм консенсуса обязан обеспечить четыре свойства:
| Свойство | Формулировка | Что ломается без него |
|---|---|---|
| Uniform Agreement | никакие два процесса не решают разные значения | два лидера, две версии баланса, split-brain |
| Integrity | процесс решает не более одного раза | «передумали» — закоммиченная транзакция откатилась |
| Validity | решённое значение было кем-то предложено | система придумала данные из воздуха |
| Termination | каждый корректный процесс когда-нибудь решает | кластер завис, запись не проходит |
Первые три — это safety («ничего плохого не произойдёт»), четвёртое — liveness («что-то хорошее в конце концов произойдёт»). Запомните разделение: весь практический консенсус построен на том, что safety держат абсолютно, а liveness — только при удачных условиях. Кластер etcd, потерявший кворум, не портит данные — он перестаёт отвечать на запись. Это не баг, а осознанный размен, и на нём стоит вся отрасль.
Слово «uniform» в Agreement несущее. Не-uniform версия разрешает упавшему процессу решить не то, что решили остальные. Звучит безобидно — процесс же умер. Но в модели crash-recovery он воскреснет, прочитает свой диск и начнёт действовать согласно неправильному решению. Промышленным системам нужен именно uniform-вариант, и именно поэтому диск в консенсусе — часть протокола, а не кеш.
Консенсус, атомарная широковещательная рассылка и реплицированный лог — одно и то же
Это ключевая эквивалентность всей области. Total order broadcast (он же atomic broadcast) — рассылка сообщений с двумя гарантиями: надёжная доставка и одинаковый порядок доставки у всех получателей. Оказывается, total order broadcast и консенсус взаимно сводимы:
- Есть консенсус — есть broadcast: запускайте по экземпляру консенсуса на каждый номер позиции в последовательности.
- Есть broadcast — есть консенсус: разошлите своё значение, решением считайте первое доставленное сообщение.
Практический смысл огромен. Total order broadcast — это ровно реплицированный лог: последовательность команд, одинаковая у всех. А одинаковый лог команд плюс детерминированная машина состояний даёт одинаковое состояние у всех реплик. Это подход replicated state machine, сформулированный Шнайдером в 1990 году, и именно так устроены etcd, ZooKeeper, Consul, контроллеры Kafka в режиме KRaft, диапазоны CockroachDB, регионы TiKV и Paxos-группы Spanner.
запросы в поток"] L --> CONS["Консенсус:
согласовать порядок"] CONS --> LOG1["Лог реплики A
1,2,3,4,5"] CONS --> LOG2["Лог реплики B
1,2,3,4,5"] CONS --> LOG3["Лог реплики C
1,2,3,4,5"] LOG1 --> SM1["Детерминированный
автомат A"] LOG2 --> SM2["Детерминированный
автомат B"] LOG3 --> SM3["Детерминированный
автомат C"] SM1 --> S["Одинаковое состояние
на всех узлах"] SM2 --> S SM3 --> S SM1 -.->|"now, rand, обход map,
внешний HTTP"| X["Состояния разъезжаются
БЕЗ ЕДИНОЙ ОШИБКИ В ЛОГЕ"]
Отсюда же следует, что «выбрать лидера» — не более простая задача, чем консенсус: договориться о лидере значит договориться об одном значении. Все схемы вида «выберем лидера по-простому, а консенсус потом» математически безнадёжны.
Слово «детерминированный» в определении RSM тоже несущее. Сценарий отказа: команда SET session.expires_at = now() + 1h применяется на трёх репликах в разные миллисекунды. Расхождение — сотни миллисекунд, никаких ошибок нигде. Через месяц лидер меняется, TTL сдвигается, и часть сессий живёт лишний час, а часть протухает раньше. Правильно — вычислять время на лидере и класть его в саму запись лога, чтобы применение было чистой функцией от лога.
2. FLP: почему честного алгоритма не существует
В 1985 году Fischer, Lynch и Paterson доказали: в полностью асинхронной системе не существует детерминированного алгоритма консенсуса, гарантирующего завершение, если возможен отказ хотя бы одного процесса. Одного. Не византийского, а честного crash-отказа. И даже при надёжной сети без потерь сообщений.
Интуиция доказательства такая. Конфигурация системы называется бивалентной, если из неё ещё достижимы оба исхода — решить 0 или решить 1. Стартовая конфигурация бивалентна. Дальше показывается, что из любой бивалентной конфигурации можно, задерживая ровно одно сообщение, попасть в другую бивалентную конфигурацию. Значит, существует бесконечное выполнение, в котором система никогда не решает. Планировщику даже не нужно ронять процессы — достаточно «удачно» тасовать задержки.
Корень проблемы тот же, что и в моделях отказов: в асинхронной системе нельзя отличить упавший процесс от медленного. Ждать вечно — нарушить Termination. Не ждать — рискнуть Agreement.
Из FLP есть ровно три выхода, и все промышленные системы используют один или два:
| Обход | Идея | Кто использует |
|---|---|---|
| Частичная синхронность | сеть асинхронна, но начиная с некоторого неизвестного момента GST задержки ограничены | Paxos, Raft, Zab, Viewstamped Replication — то есть весь прод |
| Детекторы отказов | вынести «догадку о смерти» в отдельный модуль; доказано, что слабейший достаточный детектор — это ◇W | академическая формализация того, что таймауты делают руками |
| Рандомизация | подбрасывать монетку; завершение с вероятностью 1, а не гарантированно | Ben-Or, Honey Badger BFT, часть блокчейн-протоколов |
Модель частичной синхронности (Dwork, Lynch, Stockmeyer, 1988) — формализация фразы «сеть обычно нормальная, но иногда лежит». Практический вывод, который стоит выписать на стену: safety не зависит от таймаутов, liveness зависит. Подкрутите election-timeout неправильно — получите бесконечные перевыборы и недоступность, но не расхождение данных. Если получили расхождение — это баг реализации, а не настройка.
3. Самодельный консенсус: четыре способа выстрелить себе в ногу
«Лидер — тот, кто первым захватил ключ в Redis». SET leader node-1 NX PX 30000. Работает до первой паузы GC: узел взял лок, ушёл в stop-the-world на 40 секунд, ключ протух, лок взял другой. Первый вернулся и продолжает считать себя лидером — он же не знает, что был мёртв. Redis тут ни при чём: лизу без fencing-токена безопасной не сделать, разбор в координации. В логах это выглядит так:
02:11:04.812 WARN [gc] pause of 41.3s (G1 Full GC (Allocation Failure))
02:11:46.104 INFO [leader] still leader, processing batch 88213
02:11:46.104 INFO [leader] still leader, processing batch 88213 <-- другой хост, тот же batch
Одинаковый номер батча с двух хостов — классический признак; замечают его обычно не по логам, а по дубликатам в биллинге через сутки.
«Лидер — узел с наименьшим ID из живых». Проблема в слове «живых»: у каждого узла свой список. При частичном разделении (A видит B, B видит C, A не видит C) списки различаются, и алгоритм честно выдаёт разные ответы на разных узлах. Bully algorithm из учебников безопасен только в синхронной модели, которой у вас нет.
«Дождёмся подтверждения от всех узлов, тогда точно согласовано». Требование ответа от всех N узлов даёт safety, но убивает доступность: любой один упавший или залипший узел останавливает систему. Кворум большинства — это ровно тот минимум, который сохраняет пересечение и переживает f отказов.
«Двухфазный коммит для выбора лидера». 2PC блокируется при падении координатора: проголосовавшие «да» не имеют права ни закоммитить, ни откатиться. Протокол без liveness при отказе одного узла — подробности в распределённых транзакциях.
Общий диагноз: во всех вариантах решение принимается на основе локального представления о мире. Консенсус работает потому, что решение требует пересечения кворумов — любые два большинства из N узлов имеют общий узел, и этот узел не даст двум решениям разойтись.
4. Paxos
Paxos — первый доказанный алгоритм консенсуса для crash-recovery модели. Лампорт описал его через аллегорию о парламенте острова Паксос (The Part-Time Parliament, 1998), статью почти никто не понял, и через три года вышло Paxos Made Simple — восемь страниц, из которых собственно протокол занимает две.
Три роли (один процесс обычно совмещает все): proposer предлагает значения, acceptor голосует и хранит состояние, learner узнаёт результат. Кворум — любое большинство акцепторов. Каждое предложение имеет уникальный номер бюллетеня (ballot), монотонно растущий и не пересекающийся между proposer’ами: классически это пара (счётчик, id_узла) с лексикографическим сравнением.
Два раунда
предлагаем "x=7" из промиса A2 P->>A1: Accept(5, "x=7") P->>A2: Accept(5, "x=7") A1-->>P: Accepted(5, "x=7") A2-->>P: Accepted(5, "x=7") Note over P,A3: Большинство приняло — значение выбрано навсегда
Правила акцептора умещаются в две строки псевдокода, правило proposer’а — в одну, и в ней вся суть протокола:
ACCEPTOR:
при Prepare(n):
если n <= promised: NACK
promised := n; FSYNC
вернуть Promise(n, accepted_n, accepted_v)
при Accept(n, v):
если n < promised: NACK
promised := n; accepted_n := n; accepted_v := v; FSYNC
вернуть Accepted(n, v)
PROPOSER:
n := следующий уникальный бюллетень
собрать Promise от большинства акцепторов
если хоть один промис вернул принятое значение:
v := значение с НАИБОЛЬШИМ accepted_n # своё значение отбрасывается
иначе:
v := своё значение
разослать Accept(n, v); большинство Accepted => v выбрано
Именно последнее правило делает Paxos безопасным: как только значение принято большинством, любой будущий кворум промисов обязательно содержит узел, который это значение видел, и новый proposer вынужден повторить его. Пересечение кворумов превращается в невозможность передумать.
"""Однораундовый (single-decree) Paxos. Сеть опущена: важна логика состояний."""
class Acceptor:
def __init__(self, storage):
self.storage = storage # durable: fsync ДО ответа на RPC
self.promised = 0 # наибольший обещанный номер бюллетеня
self.accepted_n, self.accepted_v = 0, None # что уже принято
def prepare(self, n: int):
if n <= self.promised:
return None # NACK: уже обещали кому-то не меньше
self.promised = n
self.storage.persist(self) # обещание обязано пережить перезапуск
return (self.accepted_n, self.accepted_v)
def accept(self, n: int, v):
if n < self.promised:
return False # обещали не мешать старшему бюллетеню
self.promised, self.accepted_n, self.accepted_v = n, n, v
self.storage.persist(self)
return True
def propose(acceptors, n: int, own_value):
"""Возвращает выбранное значение либо None, если раунд не удался."""
promises = [r for a in acceptors if (r := a.prepare(n)) is not None]
if len(promises) * 2 <= len(acceptors):
return None # кворума обещаний нет — раунд провален
# КЛЮЧЕВОЙ момент: чужое уже принятое значение сильнее своего.
best_n, best_v = max(promises, key=lambda p: p[0])
chosen = best_v if best_v is not None else own_value
acks = sum(1 for a in acceptors if a.accept(n, chosen))
return chosen if acks * 2 > len(acceptors) else None
Стоимость одного решения: 2 RTT и 4N сообщений, где N — число акцепторов. Память на акцепторе: O(1) на экземпляр консенсуса (три поля), но экземпляров столько, сколько позиций в логе, — отсюда неизбежность снапшотов. Время на proposer’е: O(N) на раунд, плюс два fsync на критическом пути.
Дуэль proposer’ов и почему это не противоречие с FLP
Сценарий: P1 делает Prepare(5), получает кворум. P2 делает Prepare(6) — акцепторы обещают шестёрке. Accept(5, v) от P1 отвергается, P1 повышает номер до 7, делает Prepare(7), и теперь отвергается Accept(6, v) от P2. И так бесконечно. Safety не нарушена, значение просто никогда не выбирается — это и есть FLP в чистом виде, воспроизводимый руками.
Лечение — не в протоколе, а в дисциплине: выбрать одного различаемого proposer’а (лидера) и позволить остальным предлагать, только если лидер молчит дольше таймаута. Рандомизация таймаутов делает дуэль маловероятной. Обратите внимание на асимметрию: лидер в Paxos нужен для liveness, а не для safety. Два proposer’а — это медленно, но не опасно, в отличие от двух лидеров в наивной схеме из раздела 3.
Multi-Paxos и почему все делают именно его
Гонять две фазы на каждую команду — расточительство. Наблюдение: фаза 1 не зависит от значения, она лишь блокирует младшие бюллетени. Значит, лидер может выполнить Prepare(n) один раз для всех будущих позиций лога и дальше слать только Accept — 1 RTT на команду. Это Multi-Paxos, и он же общий скелет Raft, Zab и Viewstamped Replication.
Что осталось за пределами двух страниц Paxos Made Simple и что пришлось изобретать заново каждой команде: выбор лидера, изменение состава кластера, снапшоты, обработка дыр в логе, восстановление после перезапуска, обнаружение дубликатов клиентских запросов. Инженеры Google описали этот разрыв в Paxos Made Live (PODC 2007) — обязательное чтение, если вам когда-нибудь захочется реализовать консенсус самостоятельно. Спойлер: они писали Chubby год и нашли баги алгоритма компилятором собственного языка спецификаций.
5. Raft: те же гарантии, но так, чтобы человек мог реализовать
Raft (Ongaro, Ousterhout, 2014) решает ту же задачу, что Multi-Paxos, но проектировался с явной целью понятности. Три решения обеспечивают её:
- Сильный лидер. Записи текут только от лидера к последователям. Последователь никогда не создаёт записи сам и никогда не спорит о содержимом.
- Выборы отделены от репликации. Смена лидера — самостоятельный, полностью описанный подпротокол, а не побочный эффект гонки бюллетеней.
- Ограничение на выборы. Лидером может стать только узел, чей лог не старше лога большинства. Это избавляет от отдельной фазы догонки лога новым лидером.
Состояния узла
Термы вместо часов
Время в Raft — это term: монотонно растущий целочисленный счётчик, логические часы кластера. В каждом терме не более одного лидера. Каждое RPC несёт term отправителя, и правило универсально: увидел term больше своего — обнови свой и стань follower; увидел меньше — отвергни сообщение. Одно правило снимает целый класс проблем с устаревшими сообщениями, которые в Paxos приходится разбирать вручную.
Персистентное состояние узла — ровно три поля: currentTerm, votedFor, log[]. Все три обязаны попасть на диск (с fsync) до отправки ответа на RPC.
Сценарий отказа: реализация не сохраняет votedFor, считая это мелочью. Узел голосует за кандидата A в терме 48, падает, поднимается с пустым votedFor, голосует за кандидата B в том же терме 48. Оба набирают большинство. Два лидера в одном терме, оба принимают записи, оба уверены в своей правоте. Признак в логах — два разных ID с одинаковым термом:
raft: 8e9e05c52164694d became leader at term 48
raft: b8e2f0c9a1d4a3f7 became leader at term 48
Тот же эффект даёт --unsafe-no-fsync в etcd и любые «оптимизации» вида «fsync раз в секунду». Формально это не потеря записи, а нарушение модели crash-recovery: узел воскресает с состоянием, которого по протоколу быть не могло.
Выборы
Follower, не получавший heartbeat дольше election timeout, переходит в кандидаты: увеличивает currentTerm, голосует за себя, рассылает RequestVote(term, candidateId, lastLogIndex, lastLogTerm).
Голос отдаётся при трёх условиях: term кандидата не меньше своего, в этом терме ещё не голосовали (или голосовали за него же), и лог кандидата не старше своего. Последнее — election restriction, и сравнение тут не по длине: сначала терм последней записи, и только при равенстве — длина. Длинный лог со старым хвостом проигрывает короткому со свежим, потому что длинный хвост мог быть незакоммиченным мусором прошлого лидера.
Из этого следует Leader Completeness Property: любая закоммиченная запись присутствует в логах всех будущих лидеров. Доказательство — снова на пересечении кворумов: закоммиченная запись есть на большинстве, победитель получил голоса большинства, эти большинства пересекаются, и узел из пересечения не проголосовал бы за кандидата без этой записи.
Split vote — два кандидата стартовали одновременно, никто не набрал большинства. Лечится рандомизацией: election timeout выбирается случайно из диапазона (в etcd — из [T, 2T)), и разброс должен заметно превышать типичный RTT, иначе выборы затягиваются на несколько раундов подряд.
Репликация лога
Лидер шлёт AppendEntries(term, prevLogIndex, prevLogTerm, entries[], leaderCommit). Follower применяет пакет, только если у него на позиции prevLogIndex лежит запись терма prevLogTerm. Отсюда индукцией получается Log Matching Property: совпадение пары (индекс, term) гарантирует совпадение всего префикса. Одна проверка на пакет — и весь лог согласован, без сравнения содержимого.
При отказе лидер уменьшает nextIndex для этого follower’а и пробует снова, пока не найдёт общий префикс; всё, что правее, у follower’а усекается. Наивный откат на единицу за раунд стоит O(k) RTT при расхождении в k записей — поэтому реальные реализации возвращают в отказе conflictTerm и conflictIndex, что сокращает до O(число различающихся термов) раундов. Сценарий, где это заметно: follower был отрезан на час, накопил 200 тысяч чужих записей, и наивная реализация догоняет его 200 тысяч RTT — на практике это выглядит как «узел вечно в состоянии probe и никогда не входит в кворум».
from dataclasses import dataclass, field
@dataclass
class Entry:
term: int
command: object
@dataclass
class RaftState:
current_term: int = 0 # персистентно, fsync до ответа
voted_for: int | None = None # персистентно
log: list[Entry] = field(default_factory=list) # персистентно; log[i] = индекс i+1
commit_index: int = 0 # волатильно
last_applied: int = 0 # волатильно
match_index: dict[int, int] = field(default_factory=dict) # только у лидера
def handle_append_entries(s: RaftState, term, prev_index, prev_term,
entries, leader_commit, persist):
"""Обработка AppendEntries на follower'е. O(len(entries)) по времени."""
if term < s.current_term:
return (s.current_term, False) # отправитель отстал
if term > s.current_term:
s.current_term, s.voted_for = term, None # новый терм — голос сброшен
if prev_index > len(s.log):
return (s.current_term, False) # дыра: у нас лог короче
if prev_index > 0 and s.log[prev_index - 1].term != prev_term:
return (s.current_term, False) # конфликт: лидер откатит nextIndex
for offset, e in enumerate(entries):
idx = prev_index + offset + 1
if idx <= len(s.log):
if s.log[idx - 1].term == e.term:
continue # запись уже совпадает
del s.log[idx - 1:] # усечение расходящегося хвоста
s.log.append(e)
if leader_commit > s.commit_index:
s.commit_index = min(leader_commit, len(s.log))
persist(s) # fsync перед ответом «ок»
return (s.current_term, True)
Сложность: время O(m) на пакет из m записей плюс O(t) на усечение хвоста длины t; амортизированно линейно по объёму лога. Память: O(L) записей до снапшота. Сеть: O(N) сообщений на команду против O(N) у Multi-Paxos — асимптотически то же самое.
Правило коммита — место, где ломаются самописные реализации
Очевидное правило «запись реплицирована на большинство → она закоммичена» неверно. Разбор на картинке выше — это Figure 8 из статьи Raft, и её стоит понять целиком, потому что почти каждая самописная реализация спотыкается ровно здесь. Коротко: запись из прошлого терма может лежать на большинстве и всё равно быть затёртой узлом, который выиграет следующие выборы по election restriction, — потому что его собственный хвост окажется свежее по терму.
Правильное правило: лидер объявляет закоммиченными по счётчику реплик только записи своего текущего терма. Унаследованные записи прошлых термов коммитятся косвенно — вместе с более поздней записью текущего терма. Практическое следствие: новый лидер обязан как можно быстрее записать no-op запись своего терма, иначе унаследованный хвост зависает незакоммиченным, а клиентские запросы, отправленные предыдущему лидеру, висят в таймауте до первой новой записи.
def advance_commit_index(s: RaftState, peers: list[int]) -> int:
"""Лидер продвигает commitIndex. O(len(log) * len(peers)) в худшем случае;
на практике берут медиану отсортированного match_index — O(n log n), n = размер кластера."""
for n in range(len(s.log), s.commit_index, -1):
if s.log[n - 1].term != s.current_term:
break # ЗАПИСЬ ПРОШЛОГО ТЕРМА: коммитить по счётчику НЕЛЬЗЯ (Figure 8)
replicas = 1 + sum(1 for p in peers if s.match_index.get(p, 0) >= n)
if replicas * 2 > len(peers) + 1:
s.commit_index = n
return n
return s.commit_index
6. Выбор лидера в проде: четыре механизма, о которых не пишут в туториалах
Pre-vote. Узел, изолированный на пять минут, всё это время крутит выборы и накручивает currentTerm до, скажем, 4021. Сеть чинится — он рассылает RequestVote(term=4021). Живой лидер с термом 48 видит больший терм, снимает с себя полномочия, кластер уходит на перевыборы. Данные не портятся, но запись встала на несколько секунд, и в логах это выглядит немотивированно:
{"level":"info","msg":"raft.node: 8e9e05c52164694d lost leader b8e2f0c9a1d4a3f7 at term 48"}
{"level":"info","msg":"8e9e05c52164694d [term: 48] received a MsgVote message with higher term from c2a91f... [term: 4021]"}
Разрыв в термах на три порядка — характерная подпись: кто-то долго варился в изоляции. Лечение — фаза pre-vote из диссертации Онгаро (§9.6): кандидат сначала спрашивает «проголосовали бы вы за меня?», не увеличивая свой терм, и повышает терм только получив большинство «да». В etcd это флаг --pre-vote (в свежих версиях включён по умолчанию — сверяйтесь с вашей версией), в hashicorp/raft он есть с 1.4.
CheckQuorum. Лидер, потерявший связь с большинством, обязан сам сложить полномочия, а не ждать, пока его сместят: если за период election timeout не пришли ответы от большинства — переход в follower. Без этого он продолжает отвечать на локальные чтения устаревшими данными, будучи уверенным в своём лидерстве.
Leadership transfer. Плановое снятие узла не должно вызывать выборы по таймауту: лидер догоняет целевой узел до конца лога, шлёт TimeoutNow, тот немедленно начинает выборы и почти гарантированно выигрывает. Простой — миллисекунды вместо секунд. В etcd это etcdctl move-leader, первая команда любого корректного rolling-restart; без неё каждый перезапуск лидера стоит кластеру полного election timeout недоступности записи.
Частичное разделение — самый неприятный случай. Односторонняя связность (лидер слышит всех, его не слышит никто) или разрыв «по диагонали». Симптом — «шторм выборов»: терм растёт на десятки в минуту, лидер меняется, никакая запись не проходит, при этом все узлы формально живы и отвечают на ping. Диагностика — скорость роста etcd_server_leader_changes_seen_total и величина etcd_server_proposals_pending. Лечение архитектурное: пять узлов вместо трёх плюс включённый pre-vote.
7. Чтения: почему GET тоже проходит через консенсус
Самое частое заблуждение: «запись через Raft, а чтение можно с лидера напрямую, он же лидер». Нет. Лидер может быть смещён пять секунд назад и не знать об этом — его просто изолировали. Читая локально, он отдаёт данные до чужих коммитов, и вы получаете stale read в формально линеаризуемой системе.
Сценарий отказа: Kubernetes-контроллер читает Pod с изолированного узла etcd, видит состояние Pending, планирует его на другой ноде. На самом деле Pod уже Running. Итог — два экземпляра пода, оба пишут в один том. Строчка в логе одна и выглядит невинно:
{"level":"warn","msg":"apply request took too long","took":"1.42s","expected-duration":"100ms","request":"key:\"/registry/pods/default/worker-7\" "}
Четыре способа читать честно:
ReadIndex — стандарт де-факто: etcd, TiKV, CockroachDB; стоимость — один раунд heartbeat, без обращения к диску. Follower read на нём же разгружает лидера ценой того же RTT, зато параллелит выдачу.
Lease read быстрее (ноль сетевых раундов), но опирается на физическое время: лидер считает, что его лидерство действительно ещё election timeout минус запас на дрейф часов. Если часы прыгнут вперёд (NTP step, live-migration с заморозкой VM), лиза «истечёт» на других узлах раньше, чем лидер это заметит, и он отдаст устаревшие данные, будучи уверенным в своей правоте. Это ровно та проблема, которую Spanner решает через TrueTime с явными границами неопределённости и commit-wait. Не Google — считайте lease read оптимизацией с известным риском и включайте осознанно.
Отдельно: etcdctl get --consistency=s (serializable) читает локально и не даёт линеаризуемости. Это законный режим, но только если вы понимаете, что покупаете; шкала вариантов — в моделях согласованности.
8. Изменение состава кластера
Замена узла — операция, в которой безопасность теряется чаще всего: во время перехода в кластере сосуществуют два представления о том, что такое «большинство». Пример катастрофы: кластер {A,B,C}, добавляем D и E разом. Часть узлов уже считает конфигурацию {A,B,C,D,E} (кворум 3), часть — старой {A,B,C} (кворум 2). A и B по старой конфигурации избирают лидера, C, D, E по новой — другого. Два лидера, оба с честным кворумом, оба принимают записи.
Raft предлагает два решения:
- Joint consensus. Промежуточная конфигурация
C_old,new, в которой решение требует большинства и в старом, и в новом составе одновременно. Пересечения гарантированы, пока идёт переход. Сложнее в реализации, зато позволяет менять сразу несколько узлов. - Single-server change. Менять по одному узлу за раз: большинства
{A,B,C}и{A,B,C,D}всегда пересекаются. Проще, и именно так делает etcd. Историческая деталь: в первой версии этот вариант содержал ошибку — при наложении нескольких последовательных изменений безопасность терялась; Онгаро исправил формулировку в диссертации, добавив требование дождаться коммита предыдущего изменения перед началом следующего.
Learner / non-voting member — узел, который получает лог, но не голосует и не входит в кворум. Обязателен при добавлении узла в живой кластер: новичок с пустым логом сразу становится частью кворума и тормозит все коммиты, пока догоняет. В etcd это member add --learner, затем member promote после того как подтянется RaftAppliedIndex.
Сценарий отказа, которым убивают продакшн: в кластере из трёх узлов один падает, оператор «на всякий случай» добавляет новый — узлов четыре, кворум 3, живых 3, работает. Падает второй — живых 2 из 4, кворума нет, запись встала. Если бы мёртвый узел сначала удалили, кластер был бы {A,B} с кворумом 2 и пережил бы это. Порядок операций: сначала remove, потом add. Диагностический признак — etcdserver: no leader при формально живых узлах: это почти всегда арифметика кворума, а не сеть.
9. Снапшоты, компакция и медленный follower
Лог растёт вечно, поэтому машина состояний периодически сериализуется в снапшот, а префикс лога удаляется. Три следствия, о которых узнают в инциденте:
- follower, отставший дальше границы компакции, не догоняется по AppendEntries — лидер шлёт
InstallSnapshotцеликом (в большом кластере Kubernetes это гигабайты, конкурирующие с heartbeat за ту же сеть, что само по себе провоцирует перевыборы); - частота снапшотов (
--snapshot-countв etcd) — компромисс между постоянным I/O и долгим восстановлением после перезапуска; - снапшот обязан соответствовать конкретному индексу лога и включать конфигурацию кластера, иначе после восстановления узел не знает, кто ещё входит в кворум.
Признак цикла «отстал → снапшот → снова отстал»:
{"level":"info","msg":"sending database snapshot","bytes":2874993152,"size":"2.9 GB","remote-peer-id":"c2a91f..."}
{"level":"warn","msg":"sending database snapshot failed","error":"write tcp ...: i/o timeout"}
{"level":"info","msg":"sending database snapshot","bytes":2874993152,"size":"2.9 GB","remote-peer-id":"c2a91f..."}
Одинаковый размер и повторы каждые несколько минут — узел никогда не догонит сам. Правильное действие: вывести его из кластера, скопировать снапшот вне протокола (обычным scp), вернуть как learner.
10. Производительность: где именно уходят миллисекунды
Стоимость одного коммита в Raft/Multi-Paxos раскладывается в короткую формулу:
время коммита ≈ fsync(лидер) + RTT до медианного узла кворума + fsync(follower)
Из неё сразу следуют все практические правила.
| Ручка | Эффект | Типичная ошибка |
|---|---|---|
| Диск лидера | fsync WAL на каждый коммит; p99 etcd_disk_wal_fsync_duration_seconds должен быть меньше 10 мс |
etcd на сетевом диске или на общем хосте с шумным соседом |
| Размер кластера | кворум растёт как floor(N/2) + 1, латентность определяется медианным узлом |
«поставим 7 для надёжности» — получаем медленнее и не надёжнее |
| Чётное число узлов | 4 узла терпят тот же 1 отказ, что и 3, но кворум больше | «добавим четвёртый на всякий случай» |
| Батчинг | несколько команд в одну запись лога и один fsync | выключен, потому что «увеличивает latency» — на деле снижает p99 под нагрузкой |
| Конвейеризация | не ждать ответа перед отправкой следующего пакета | последовательная отправка режет пропускную способность в разы |
| Геораспределение | кворум через океан = 150+ мс на запись | «растянем кластер на три региона» без пересмотра SLA записи |
Устойчивость: кластер из 2f + 1 узлов переживает f отказов. N=3 → 1, N=5 → 2, N=7 → 3. Чётные значения бессмысленны: N=4 переживает тот же один отказ, что и N=3, платя большим кворумом на каждой записи. Практический предел — 5; документация etcd прямо не рекомендует идти дальше семи.
Что делают, когда один RTT до кворума слишком дорог:
- Multi-Raft / Paxos-группы. Разбить данные на диапазоны и вести отдельную группу консенсуса на каждый: CockroachDB (диапазоны), TiKV (регионы), Spanner (группы на директорию). Пропускная способность масштабируется числом групп, а heartbeat’ы между общими узлами склеиваются в один пакет (coalesced heartbeats); связь с разбиением данных — в партиционировании.
- Flexible Paxos (2016): достаточно, чтобы пересекались кворумы фазы 1 и фазы 2, а не любые два — кворум фазы 2 можно сделать меньше большинства за счёт более дорогих выборов. EPaxos (2013): без выделенного лидера, непротиворечивые команды коммитятся за 1 RTT, цена — тяжёлый путь восстановления и сложность реализации.
- Witness-реплики — узел, участвующий в голосовании, но не хранящий данные: дешёвый способ получить нечётный кворум в двух ДЦ.
- И, наконец, не использовать консенсус — самая частая правильная оптимизация, см. раздел 14.
11. Консенсус и exactly-once: что он на самом деле гарантирует
Формулировка, ради которой стоит читать этот раздел: консенсус даёт exactly-once применение записи лога к машине состояний и ничего не обещает про клиентский запрос. Клиент находится снаружи реплицированной машины состояний, и канал до него протоколом не покрыт.
Третьего исхода клиент не различает C->>L2: POST /transfer 100 USD (повтор) L2->>F: AppendEntries(перевод 100) F-->>L2: успех L2-->>C: 200 OK Note over C,L2: Деньги списаны ДВАЖДЫ.
Консенсус отработал безупречно.
Это задача двух генералов в чистом виде: нельзя построить надёжный канал поверх ненадёжного за конечное число сообщений. Никакой алгоритм консенсуса эту границу не сдвигает, и никакая настройка таймаутов её не обходит. Клиент, получивший таймаут, обязан повторить запрос; система, получившая повтор, обязана уметь его распознать.
Что делают вместо. Раздел 6.3 диссертации Онгаро описывает штатное решение — линеаризуемую семантику для клиента поверх Raft:
- Клиент регистрирует сессию, и сама регистрация идёт через лог, чтобы
clientIdбыл одинаковым на всех репликах. - Каждая команда несёт порядковый номер в рамках сессии.
- Машина состояний хранит таблицу
clientId → (последний номер, сохранённый ответ). - Пришла команда с уже виденным номером — она не применяется, возвращается сохранённый ответ.
- Сессии истекают, и срок истечения обязан вычисляться детерминированно — по индексу лога или по времени, записанному в саму запись лога. Истечение по локальным часам реплики разведёт состояния (см. раздел 1).
Иначе говоря, дедупликация переезжает внутрь машины состояний и становится частью того же атомарного применения, что и сам эффект. Это общий рецепт, который в треке повторяется в каждой второй статье: at-least-once доставка плюс детерминированный ключ идемпотентности плюс дедупликация в той же атомарной операции, что и эффект. Развёрнутый разбор паттернов — в гарантиях доставки и идемпотентности.
Как это выглядит в реальных API:
- etcd не продаёт «exactly-once», он продаёт сравнение с ревизией: транзакция
txnвида «еслиmod_revision(key) == 42, то put, иначе fail». Повтор после таймаута безопасен, потому что второй раз условие не выполнится. Это compare-and-swap поверх лога, а не магия. - ZooKeeper прямо документирует проблему:
CONNECTIONLOSS— recoverable ошибка, и клиент не знает, применилась ли операция. Классическая ловушка —createс флагомSEQUENTIAL: повтор создаёт второй узел с другим номером, и оба выглядят легитимно. Штатный обход из документации — вкладывать уникальный идентификатор клиента в имя узла и при повторе сначала искать свой. - Kafka с
enable.idempotence=trueдедуплицирует по(producerId, epoch, sequence)— только внутри Kafka и только в рамках сессии продюсера. Транзакции Kafka дают атомарность «прочитал-обработал-записал» внутри Kafka; шаг во внешнюю БД в неё не входит — отсюда outbox-паттерн в распределённых транзакциях.
Симптом того, что этим пренебрегли, всегда один: дубликаты появляются пачками сразу после инцидента, а не равномерно во времени. Первое, что стоит сделать при жалобе на двойные списания, — наложить график дублей на график перевыборов лидера; совпадение всплесков закрывает вопрос за минуту.
12. Семейство протоколов: карта местности
Правая ветвь важна не меньше левой. «Кворум» и «консенсус» — разные вещи: кворумная запись в Cassandra не упорядочивает конкурирующие записи и допускает потерю обновления, а консенсус даёт тотальный порядок. Разбор кворумов — в репликации.
13. Как это выглядит в реальных системах
| Система | Алгоритм | Особенности | Что смотреть при инциденте |
|---|---|---|---|
| etcd | Raft (etcd-io/raft) | ReadIndex по умолчанию, learners, move-leader |
etcd_disk_wal_fsync_duration_seconds, etcd_server_leader_changes_seen_total, etcd_server_proposals_pending |
| ZooKeeper | Zab | FIFO-порядок на клиента, префиксная гарантия, одноразовые watch | fsync-ing the write ahead log ... took Nms, outstanding_requests, состояния LOOKING/LEADING |
| Consul | hashicorp/raft | autopilot, non-voting servers в Enterprise | raft.commitTime, raft.leader.lastContact, строки heartbeat timeout reached |
| Kafka (KRaft) | Raft-подобный (KIP-500) | кворум контроллеров отдельно от данных; данные — ISR, не консенсус | controller.quorum.*, переходы Completed transition to Leader(epoch=N) |
| CockroachDB | Multi-Raft | группа на диапазон, coalesced heartbeats, leaseholder для чтений | ranges.underreplicated, replicas.leaders_not_leaseholders |
| TiKV | Multi-Raft | группа на регион, Placement Driver балансирует | tikv_raftstore_apply_log_duration, миграции регионов |
| Spanner | Multi-Paxos + TrueTime | Paxos-группа на директорию; commit-wait даёт внешнюю согласованность | ширина границ неопределённости TrueTime |
| Chubby | Multi-Paxos | грубые блокировки и sequencer’ы (fencing-токены) | описано в Paxos Made Live |
| Cassandra | Paxos только для LWT | обычные записи — кворумы без консенсуса | casWriteTimeout, стоимость около 4 RTT на LWT |
Две строки заслуживают отдельного внимания. Kafka ведёт консенсус только для метаданных контроллера: сами данные реплицируются через ISR, что не то же самое и разобрано в репликации и очередях. Cassandra по умолчанию консенсус не использует вовсе — LWT (IF NOT EXISTS) включает Paxos точечно и стоит около четырёх раундов, поэтому решение «просто добавим IF NOT EXISTS везде» надёжно кладёт кластер; предметный разбор — в Cassandra и wide-column.
Полезные строки логов, по которым отличают проблему консенсуса от проблемы сети:
# ZooKeeper: диск не тянет — главный источник «случайных» перевыборов
WARN [SyncThread:2:FileTxnLog@338] - fsync-ing the write ahead log in SyncThread:2 took 4021ms
which will adversely affect operation latency. File size is 67108880 bytes.
# Consul: лидер честно уходит, потеряв кворум (CheckQuorum в действии)
[WARN] agent.server.raft: failed to contact quorum of nodes, stepping down
# etcd: heartbeat не успевает — почти всегда I/O, а не сеть
{"level":"warn","msg":"failed to send out heartbeat on time (exceeded the 100ms heartbeat interval)",
"heartbeat-interval":"100ms","expected-duration":"200ms","exceeded-duration":"312.6ms"}
Правило большого пальца при разборе: если терм растёт — смотрите сначала на диск, потом на сеть, и только потом на код. Медленный fsync выглядит в логах ровно как сетевая проблема, и на этом теряют часы.
14. Когда консенсус не нужен
Консенсус — самый дорогой примитив в арсенале, и большая часть систем применяет его не там. Честный чек-лист перед тем, как поднимать кластер:
- Нужен ли глобальный инвариант? «Баланс не уходит в минус», «ровно один активный лидер», «уникальность email» — да, нужен консенсус. «Счётчик просмотров», «набор тегов», «корзина» — нет, достаточно CRDT.
- Можно ли свести к CAS в одном шарде? Если инвариант локален для одного ключа, обычная БД с транзакцией делает это дешевле любого распределённого протокола. Один PostgreSQL с репликой закрывает больше сценариев, чем принято думать; см. транзакции и изоляцию.
- Можно ли заменить согласование компенсацией? Отель, который принял две брони на один номер и извинился апгрейдом, зарабатывает больше, чем отель с блокировкой на каждую операцию. Это ход саги — 08.
- Можно ли вынести консенсус в готовую систему? Собственный Raft — это годы работы и Jepsen-тесты. etcd или ZooKeeper уже написаны и проверены; используйте их как источник истины для конфигурации и лидерства, а данные держите вне их.
Отдельно: консенсус в описанном виде не защищает от византийских отказов. Raft и Paxos предполагают, что узлы либо работают корректно, либо молчат. Узел с испорченной памятью, багом сериализации или скомпрометированный узел ломает их безнаказанно: он может проголосовать дважды в одном терме или подтвердить запись, которой у него нет. Для такой модели нужны BFT-протоколы (PBFT и потомки), требующие 3f + 1 узлов и заметно большего трафика; их применяют там, где участники не доверяют друг другу по построению, — в блокчейн-сетях и межорганизационных реестрах. Подробнее о модели — в моделях отказов.
15. Типичные ошибки
- Писать свой Raft. Между «прошли тесты» и «корректно» лежат persistence перед ответом, правило коммита прошлых термов, pre-vote, изменение состава, снапшоты, дедупликация клиентских запросов — каждый пункт отдельный класс инцидентов.
- Не делать fsync ради latency. Работает годами, ломается молча при одновременной потере питания в стойке.
- Считать, что чтение с лидера линеаризуемо. Без ReadIndex, лизы или CheckQuorum — нет.
- Чётный размер кластера и добавление узла вместо удаления мёртвого. Кворум растёт быстрее, чем доступность.
- Растягивать кворум на регионы, не поменяв SLA записи. Один RTT через океан не сжимается ничем.
- Хранить в etcd/ZooKeeper данные, а не метаданные. Каждая запись — fsync на большинстве узлов; это координационный примитив, а не база.
- Недетерминизм в машине состояний.
now(),rand(), порядок обхода map, внешние вызовы — реплики разъедутся без единой ошибки в логе. - Принимать «есть на большинстве» за «закоммичено» и полагаться на exactly-once. Первое — Figure 8, второе консенсус дать не может в принципе.
Как это проверяют
Консенсус — область, где обычные тесты почти бесполезны: баги проявляются на редких переплетениях отказов, которые не воспроизводятся вручную. Три работающих подхода: формальная спецификация (у Raft есть машинно-проверенная TLA+-модель; Amazon описала тот же опыт для S3 и DynamoDB), Jepsen (реальный кластер, инъекция разделений и скачков часов, проверка истории чекерами Knossos/Elle — отчёты на jepsen.io/analyses) и детерминированная симуляция всего кластера в одном процессе с виртуальным временем (FoundationDB, TigerBeetle, Antithesis), позволяющая воспроизвести баг по seed’у. Подробный разбор — в тестировании распределённых систем.
Мини-итог
- Консенсус — договориться об одном значении и не передумать: Agreement, Integrity, Validity (safety) плюс Termination (liveness).
- Консенсус, total order broadcast и реплицированный лог взаимно сводимы. Лог плюс детерминированный автомат = одинаковое состояние на всех узлах; любой недетерминизм ломает это без единой ошибки в логах.
- FLP запрещает детерминированный алгоритм с гарантией завершения в асинхронной модели. Реальные системы жертвуют liveness, но никогда safety.
- Paxos безопасен благодаря одному правилу: новый proposer обязан повторить уже принятое значение с наибольшим номером. Multi-Paxos амортизирует фазу 1 на весь лог.
- Raft — тот же скелет с явным лидером, термами как логическими часами и ограничением на выборы, которое обеспечивает Leader Completeness.
- Записи прошлых термов нельзя коммитить по счётчику реплик — отсюда no-op сразу после победы на выборах. Чтение с лидера без ReadIndex/лизы/CheckQuorum не линеаризуемо.
- Состав кластера менять по одному узлу, новичка вводить как learner, мёртвого удалять до добавления живого. Медленный диск ломает консенсус так же надёжно, как разорванная сеть.
- Exactly-once консенсус не даёт: клиент вне машины состояний. Работающая замена — сессии клиента с порядковыми номерами и дедупликация внутри того же атомарного применения, что и эффект.
Источники
- Fischer, Lynch, Paterson, Impossibility of Distributed Consensus with One Faulty Process (1985) — та самая невозможность.
- Dwork, Lynch, Stockmeyer, Consensus in the Presence of Partial Synchrony (1988) — модель, в которой работают все реальные алгоритмы; Chandra, Toueg, Unreliable Failure Detectors (1996) — слабейший достаточный детектор отказов.
- Leslie Lamport, The Part-Time Parliament (1998) и Paxos Made Simple (2001).
- Chandra, Griesemer, Redstone, Paxos Made Live — An Engineering Perspective (PODC 2007) — всё, чего нет в статьях про Paxos.
- Ongaro, Ousterhout, In Search of an Understandable Consensus Algorithm (Raft, USENIX ATC 2014) и диссертация Онгаро — главы 4 (состав кластера), 6 (линеаризуемая семантика для клиента), 9 (pre-vote).
- Junqueira, Reed, Serafini, Zab (DSN 2011) — протокол ZooKeeper; Oki, Liskov, Viewstamped Replication (1988) — лидер и view change за десять лет до Paxos.
- Moraru, Andersen, Kaminsky, EPaxos (SOSP 2013) и Howard, Malkhi, Spiegelman, Flexible Paxos (2016).
- Burrows, The Chubby Lock Service (OSDI 2006); Schneider, Implementing Fault-Tolerant Services Using the State Machine Approach (1990) — исходная формулировка RSM.
- Corbett et al., Spanner (OSDI 2012) — Paxos-группы плюс TrueTime; DeCandia et al., Dynamo (SOSP 2007) — как выглядит система, которая консенсус сознательно не использует.
- Raft site с интерактивными визуализациями, документация etcd по настройке, Jepsen analyses, Martin Kleppmann, «Designing Data-Intensive Applications», глава 9.
- Родственное на портале: NewSQL и распределённые БД, паттерны устойчивости.
Что дальше
Консенсус даёт одну согласованную последовательность решений внутри одной группы. Но бизнес-операция обычно затрагивает несколько таких групп и несколько сервисов, и там появляются свои протоколы и свои способы всё сломать: Распределённые транзакции: 2PC, saga, outbox, компенсации.