Распределённые системы Консенсус: Paxos, Raft, выбор лидера, репликация лога
0%

Консенсус: Paxos, Raft, выбор лидера, репликация лога

Консенсус: 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.

Отсюда же следует, что «выбрать лидера» — не более простая задача, чем консенсус: договориться о лидере значит договориться об одном значении. Все схемы вида «выберем лидера по-простому, а консенсус потом» математически безнадёжны.

Слово «детерминированный» в определении 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_узла) с лексикографическим сравнением.

Два раунда

Правила акцептора умещаются в две строки псевдокода, правило 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, но проектировался с явной целью понятности. Три решения обеспечивают её:

  1. Сильный лидер. Записи текут только от лидера к последователям. Последователь никогда не создаёт записи сам и никогда не спорит о содержимом.
  2. Выборы отделены от репликации. Смена лидера — самостоятельный, полностью описанный подпротокол, а не побочный эффект гонки бюллетеней.
  3. Ограничение на выборы. Лидером может стать только узел, чей лог не старше лога большинства. Это избавляет от отдельной фазы догонки лога новым лидером.

Состояния узла

Термы вместо часов

Время в 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, иначе выборы затягиваются на несколько раундов подряд.

Репликация лога

Анатомия реплицированного лога Raft: индексы, термы, commitIndex, расхождение хвоста

Лидер шлёт 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 — асимптотически то же самое.

Правило коммита — место, где ломаются самописные реализации

Ловушка коммита в Raft: запись на большинстве ещё не значит закоммичена

Очевидное правило «запись реплицирована на большинство → она закоммичена» неверно. Разбор на картинке выше — это 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 применение записи лога к машине состояний и ничего не обещает про клиентский запрос. Клиент находится снаружи реплицированной машины состояний, и канал до него протоколом не покрыт.

Это задача двух генералов в чистом виде: нельзя построить надёжный канал поверх ненадёжного за конечное число сообщений. Никакой алгоритм консенсуса эту границу не сдвигает, и никакая настройка таймаутов её не обходит. Клиент, получивший таймаут, обязан повторить запрос; система, получившая повтор, обязана уметь его распознать.

Что делают вместо. Раздел 6.3 диссертации Онгаро описывает штатное решение — линеаризуемую семантику для клиента поверх Raft:

  1. Клиент регистрирует сессию, и сама регистрация идёт через лог, чтобы clientId был одинаковым на всех репликах.
  2. Каждая команда несёт порядковый номер в рамках сессии.
  3. Машина состояний хранит таблицу clientId → (последний номер, сохранённый ответ).
  4. Пришла команда с уже виденным номером — она не применяется, возвращается сохранённый ответ.
  5. Сессии истекают, и срок истечения обязан вычисляться детерминированно — по индексу лога или по времени, записанному в саму запись лога. Истечение по локальным часам реплики разведёт состояния (см. раздел 1).

Иначе говоря, дедупликация переезжает внутрь машины состояний и становится частью того же атомарного применения, что и сам эффект. Это общий рецепт, который в треке повторяется в каждой второй статье: at-least-once доставка плюс детерминированный ключ идемпотентности плюс дедупликация в той же атомарной операции, что и эффект. Развёрнутый разбор паттернов — в гарантиях доставки и идемпотентности.

Как это выглядит в реальных API:

  • etcd не продаёт «exactly-once», он продаёт сравнение с ревизией: транзакция txn вида «если mod_revision(key) == 42, то put, иначе fail». Повтор после таймаута безопасен, потому что второй раз условие не выполнится. Это compare-and-swap поверх лога, а не магия.
  • ZooKeeper прямо документирует проблему: CONNECTIONLOSSrecoverable ошибка, и клиент не знает, применилась ли операция. Классическая ловушка — 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. Типичные ошибки

  1. Писать свой Raft. Между «прошли тесты» и «корректно» лежат persistence перед ответом, правило коммита прошлых термов, pre-vote, изменение состава, снапшоты, дедупликация клиентских запросов — каждый пункт отдельный класс инцидентов.
  2. Не делать fsync ради latency. Работает годами, ломается молча при одновременной потере питания в стойке.
  3. Считать, что чтение с лидера линеаризуемо. Без ReadIndex, лизы или CheckQuorum — нет.
  4. Чётный размер кластера и добавление узла вместо удаления мёртвого. Кворум растёт быстрее, чем доступность.
  5. Растягивать кворум на регионы, не поменяв SLA записи. Один RTT через океан не сжимается ничем.
  6. Хранить в etcd/ZooKeeper данные, а не метаданные. Каждая запись — fsync на большинстве узлов; это координационный примитив, а не база.
  7. Недетерминизм в машине состояний. now(), rand(), порядок обхода map, внешние вызовы — реплики разъедутся без единой ошибки в логе.
  8. Принимать «есть на большинстве» за «закоммичено» и полагаться на 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 консенсус не даёт: клиент вне машины состояний. Работающая замена — сессии клиента с порядковыми номерами и дедупликация внутри того же атомарного применения, что и эффект.

Источники

Что дальше

Консенсус даёт одну согласованную последовательность решений внутри одной группы. Но бизнес-операция обычно затрагивает несколько таких групп и несколько сервисов, и там появляются свои протоколы и свои способы всё сломать: Распределённые транзакции: 2PC, saga, outbox, компенсации.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов