Распределённые системы Репликация: лидер и последователи, кворумы, конфликты, CRDT
0%

Репликация: лидер и последователи, кворумы, конфликты, CRDT

Репликация: лидер и последователи, кворумы, конфликты, CRDT

Если бы данные никогда не менялись, репликация была бы задачей на scp. Скопировали файл на пять машин, поставили балансировщик — готово. Вся сложность появляется ровно в тот момент, когда данные начинают меняться, и весь предмет этой статьи — про то, как договориться, какие изменения считаются существующими.

Это не красивая формулировка, а операционально точная. Каждая схема репликации — это ответ на один вопрос: в какой момент клиенту можно сказать «ОК»? И каждый ответ покупается конкретным, воспроизводимым классом отказов. Сказали «ОК» после записи в локальный WAL лидера — купили окно потерянных коммитов при failover. Сказали после подтверждения от всех реплик — купили остановку записи при падении одной машины. Сказали после кворума — купили отсутствие линеаризуемости и необходимость разбирать конфликты. Бесплатного варианта в этом списке нет, и его наличие не зависит от вендора.

Репликацию заводят ровно по трём причинам, и их полезно держать раздельно, потому что они друг другу противоречат: отказоустойчивость (узел умрёт — данные останутся; важна долговечность подтверждённой записи), латентность (копия рядом с пользователем отвечает за 5 мс, а не за 150; важна географическая близость, а она прямо противоречит синхронности) и масштабирование чтений (десять реплик читают в десять раз больше; платим отставанием).

Дальше по тексту предполагается, что вы уже знакомы с моделями отказов (особенно с тем, что отказ узла неотличим от медленной сети), с логическими часами и с моделями согласованности. Репликация — это механизм; согласованность — наблюдаемое свойство, которое из этого механизма получается.

Три топологии и один вопрос

Три топологии репликации: один лидер, несколько лидеров, без лидера

Один лидер Несколько лидеров Без лидера
Кто принимает запись ровно один узел любой лидер в своём регионе клиент пишет напрямую на w узлов
Конфликты записи невозможны по конструкции неизбежны неизбежны
Порядок записей глобальный частичный частичный
Точка отказа выборы лидера разрешение конфликтов тихое расхождение реплик
Примеры PostgreSQL, MySQL, MongoDB, Kafka-партиция BDR, Active-Active MySQL, CouchDB Cassandra, Riak, Dynamo, Voldemort

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

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

Что именно реплицируется

Реплицируется не «данные», а поток изменений, и формат этого потока определяет половину практических проблем.

Уровень Что едет по проводу Где применяется Чем ломается
Statement-based текст SQL-запроса MySQL binlog_format=STATEMENT недетерминизм: NOW(), RAND(), UUID(), LIMIT без ORDER BY, побочные эффекты триггеров
Физический (WAL/redo) байты изменённых страниц PostgreSQL streaming replication, Oracle Data Guard реплика жёстко привязана к версии и разрядности сервера; мажорный апгрейд без остановки невозможен
Логический (row-based) строки «до/после» MySQL binlog_format=ROW, PostgreSQL logical replication больше трафика; исторически не переносит DDL и последовательности — проверяйте по версии
Триггерный всё, что напишете сами Bucardo, Londiste, самописное гибко и очень дорого по CPU; легко получить рассинхрон в собственном коде

Сценарий отказа: MySQL со statement-based репликацией и запросом INSERT INTO sessions (id, expires) VALUES (UUID(), NOW() + INTERVAL 1 HOUR). На лидере одни значения, на реплике другие. Ничего не падает, ошибок нет. Через месяц кто-то делает failover, и половина активных сессий внезапно оказывается просроченной, а половина живёт лишний час. Именно из-за таких случаев MySQL с 5.7 по умолчанию использует ROW, а MIXED автоматически переключается на построчный формат при обнаружении недетерминированной конструкции.

Синхронно, асинхронно и коварное «полусинхронно»

  • Асинхронно. Лидер коммитит локально и сразу отвечает клиенту. Латентность записи не зависит от реплик, отказ реплики не влияет на запись. Цена — окно потерянных подтверждённых коммитов (ненулевой RPO).
  • Синхронно. Лидер ждёт подтверждения реплики. RPO ноль, но латентность записи = RTT до реплики, а падение реплики останавливает запись целиком.
  • Полусинхронно. Ждём одну реплику из нескольких. В PostgreSQL это synchronous_standby_names = 'ANY 1 (s1, s2, s3)' — компромисс, который действительно работает.

Тонкость PostgreSQL, которую часто пропускают: synchronous_commit имеет пять уровней, и три из них — не то, что вы думаете. remote_write означает «реплика приняла WAL в память», а не «записала на диск»: одновременный отказ питания на лидере и реплике потеряет данные. on — записала и сбросила на диск, но не применила: сразу после failover чтение на новой реплике может ещё не видеть эту транзакцию. remote_apply — применила и видна в запросах. Каждый шаг стоит латентности.

Сценарий отказа: тихая деградация MySQL semi-sync. По умолчанию, если реплика не отвечает дольше rpl_semi_sync_source_timeout (10 секунд), источник автоматически переключается в асинхронный режим и продолжает подтверждать коммиты. В логе появляется одна строка:

[Warning] [MY-011166] [Repl] Timeout waiting for reply of binlog (file: binlog.000842, pos: 194883021), semi-sync up to file binlog.000842, position 194883021.
[Note] [MY-011158] [Repl] Semi-sync replication switched OFF.

После этого система месяцами работает как асинхронная, дашборд «semi-sync enabled = ON» показывает конфиг, а не факт, и никто не знает, что RPO уже не ноль. Мониторить нужно Rpl_semi_sync_source_status (реальное состояние) и Rpl_semi_sync_source_no_tx (число коммитов, ушедших без подтверждения) — а не строчку в my.cnf.

Сколько именно данных теряется при failover

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

FATAL:  requested timeline 2 is not a child of this server's history
DETAIL:  Latest checkpoint is at 0/8A000028 on timeline 1, but in the history
         of timeline 2 the server forked off from that timeline at 0/89FFE120.

Разница 0/8A000028 − 0/89FFE120 — это ровно объём WAL с подтверждёнными клиенту транзакциями, которых нет в новой истории. Их можно прочитать через pg_waldump и вручную восстановить бизнес-смысл; pg_rewind их молча выбросит. Первое, что нужно делать после аварийного failover, — снять дамп расходящегося WAL, а не запускать pg_rewind. У MongoDB это ещё нагляднее: вернувшийся узел складывает откатанные документы в BSON-файлы в каталоге rollback/ — буквально папка с потерянными данными; в логе это Rolling back {oplog entries} и Wrote rollback file.

Split brain: два лидера — это не теория

Каноничный публичный разбор — инцидент GitHub 21 октября 2018. 43-секундное разделение сети между восточным и западным побережьем США; Orchestrator честно отработал свою логику и промоутил кластер на западе; когда связность восстановилась, обе половины имели записи, которых не было у другой. Итог — более суток деградации, потому что автоматическое слияние двух разошедшихся историй MySQL невозможно в принципе.

Механика проблемы всегда одна и та же и следует из моделей отказов: таймаут не отличает мёртвый узел от медленного. Старый лидер может быть жив и просто застрять — в 20-секундной stop-the-world паузе GC, в свопе, в перегруженном гипервизоре. Он проснётся и продолжит писать так, будто ничего не произошло. Как это выглядит в приложении: на новом лидере начинают сыпаться

ERROR: duplicate key value violates unique constraint "orders_pkey"
DETAIL: Key (id)=(88214) already exists.

потому что старый лидер выдал те же значения последовательности, что и новый. Ошибки уникальности после failover — почти всегда симптом split brain, а не «случайности».

Единственная надёжная защита — фенсинг, то есть активное лишение старого лидера права писать:

  • STONITH — выключить узел через IPMI или API облака. Грубо и надёжно.
  • Лиза (lease) с TTL — лидер обязан продлевать аренду; не продлил, значит через TTL сам себя понижает. Требует ограничения на дрейф часов, см. время и часы и координацию.
  • Fencing token — монотонно растущий номер эпохи; хранилище отвергает запись с устаревшим токеном. Так работают эпохи в Kafka и RBD exclusive-lock в Ceph.

Ключевой принцип: «лидер по таймауту» — это не лидер. Лидерство должно быть решением кворума, и именно этим занимается консенсус, см. Paxos и Raft.

Отставание реплик глазами пользователя

Асинхронная реплика отстаёт всегда. Вопрос лишь в том, какие аномалии это создаёт для конкретного пользователя:

  • Нарушение read-your-writes. Пользователь загрузил аватар, обновил страницу — старый аватар. Или страшнее: оплата вернула 200, а /orders пуст.
  • Нарушение monotonic reads. Два подряд запроса ушли на разные реплики, вторая отстаёт сильнее — комментарий появился и исчез. Пользователь думает, что сходит с ума.
  • Нарушение consistent prefix reads. При шардировании ответ приезжает раньше вопроса: реплики разных партиций отстают по-разному.

Практическое лечение — не «сделаем всё синхронным», а маршрутизация по свежести. Клиент носит с собой позицию своей последней записи (LSN в PostgreSQL, operationTime в MongoDB, токен сессии), и роутер выбирает только те реплики, которые до неё дотянулись:

import random

def route_read(user_lsn: str | None, replicas: list["Replica"], leader: "Leader"):
    """Читаем с реплики, только если она догнала LSN последней записи
    пользователя. user_lsn кладём в подписанную куку при каждой записи.
    Сложность: O(R) по числу реплик на запрос, R обычно 3-10."""
    if user_lsn is None:
        return random.choice(replicas)          # пользователь ничего не писал

    fresh = [r for r in replicas if r.replay_lsn() >= user_lsn]
    if fresh:
        return random.choice(fresh)
    return leader                               # честно деградируем к лидеру

Три рабочих варианта, от дешёвого к дорогому: sticky-сессия на реплику в течение N секунд после записи; чтение с лидера только на критичных путях (после POST в рамках той же сессии); LSN-токен как выше. В PostgreSQL 18 появилась функция pg_wal_replay_wait(), которая позволяет реплике подождать нужный LSN на своей стороне — до этого приходилось опрашивать pg_last_wal_replay_lsn() вручную.

Кворумы: репликация без лидера

Второй большой класс — leaderless-репликация, популяризованная статьёй Amazon Dynamo: Amazon’s Highly Available Key-value Store (SOSP 2007). Клиент сам пишет на несколько узлов и сам читает с нескольких.

Кворумы: правило w + r больше n и три случая

Правило w + r > n (n — число реплик ключа, w — сколько подтвердили запись, r — сколько опрошено при чтении) гарантирует ровно одно: множества записи и чтения пересекаются хотя бы по одному узлу, значит читатель увидит хотя бы одну копию последнего подтверждённого значения. Cassandra даёт это как настройку на каждый запрос: ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL.

И теперь — чего кворум НЕ даёт. Это самая частая ошибка в понимании Dynamo-подобных систем.

  1. Кворум — не линеаризуемость. Пересечение множеств гарантирует «увидим значение», а не «увидим одно и то же значение все и одновременно». Подробный разбор различий — в моделях согласованности.
  2. Запись на w узлов не атомарна. Пока она разъезжается, один читатель уже видит новое значение, а параллельный — ещё старое. Монотонности чтений тоже нет: следующий запрос того же клиента может попасть на отставший набор и вернуться назад во времени.
  3. Частично успешная запись не откатывается. Если при w=3 подтвердили только двое, клиент получит ошибку — но на двух узлах данные останутся и никуда не денутся. Read repair позже разнесёт их на остальных, и «неудачная» запись материализуется.

Сценарий отказа: сервис получил таймаут на записи в Cassandra и решил, что записи не было. Ретраит с другим значением (например, пересчитанным статусом). Оба значения живы, побеждает то, у которого больше timestamp — то есть ретрай. Пока это upsert идемпотентного состояния — нормально. Как только это counter или «добавить в список» — данные искажаются, и в логах нет ничего, кроме WriteTimeoutException. Отсюда практическое правило Cassandra-приложений: любая запись должна быть идемпотентным upsert’ом состояния, а не дельтой.

Read repair и anti-entropy

Расхождения чинят двумя механизмами:

  • Read repair (при чтении) — при чтении координатор видит разные версии на разных узлах и дописывает свежую на отставшие. Чинит только то, что читают. Холодные данные расходятся вечно.
  • Anti-entropy — фоновое сравнение через деревья Меркла (Dynamo, §4.7): узлы обмениваются хешами диапазонов и спускаются по дереву только там, где хеши разошлись. Стоимость сравнения — O(log N) обменов на расхождение вместо передачи всего диапазона. В Cassandra это nodetool repair.

Сценарий отказа — воскрешение удалённых данных, классика Cassandra. Удаление там не удаляет, а пишет tombstone; tombstone живёт gc_grace_seconds (по умолчанию 10 дней) и потом вычищается компакцией. Если узел был offline дольше этого срока и вернулся без полного repair, у него сохранилась старая живая запись, а tombstone у остальных уже исчез — read repair добросовестно распространит «живое» значение обратно, и удалённый пользователь возвращается в базу. В логах не будет ничего. Это самый неприятный вид отказа: молчаливый, отложенный на недели, обнаруживаемый только жалобой. Отсюда правило — nodetool repair должен проходить до конца чаще, чем gc_grace_seconds, и за этим нужен отдельный мониторинг.

Sloppy quorum и hinted handoff

Если часть узлов из preference list ключа недоступна, Dynamo-системы могут записать значение на любые другие живые узлы, пометив это «подсказкой» (hint). Формально кворум набран, доступность записи сохранена. Фактически — как показано в третьей панели SVG выше — пересечения с читающим кворумом больше нет, и правило w + r > n перестаёт что-либо значить. Данные приедут по назначению позже, когда узлы вернутся.

Это осознанный размен: sloppy quorum повышает доступность записи и понижает гарантии чтения. В Cassandra окно хранения подсказок ограничено (max_hint_window_in_ms, по умолчанию 3 часа); если узел не вернулся за это время, подсказки просто выбрасываются, и восстановление возможно только через repair.

Цепочка вместо кворума

Третья структура, которая редко попадает в учебники, но крутит реальные хранилища, — chain replication (van Renesse, Schneider, OSDI 2004). Реплики выстроены в цепь HEAD → … → TAIL: запись входит в голову, последовательно проходит все звенья и подтверждается клиенту только хвостом; чтения обслуживает тоже хвост. Конвейер записи в HDFS устроен ровно так же — клиент шлёт пакеты первому datanode, тот следующему, ack возвращается по цепочке назад.

Ради чего это делают: линеаризуемость без голосования на каждый запрос. Хвост по построению содержит ровно множество подтверждённых записей, поэтому чтение с него не требует опроса кворума — один узел, один RTT, полная согласованность. CRAQ (USENIX ATC 2009) снимает главное ограничение схемы — то, что читать можно только с хвоста: любое звено отвечает само, если у него нет «грязных» (ещё не подтверждённых хвостом) версий, иначе спрашивает у хвоста номер актуальной версии. Чтения масштабируются линейно, линеаризуемость сохраняется.

Сценарий отказа: выпало среднее звено. Запись останавливается целиком, пока внешний координатор (на практике — etcd или ZooKeeper, см. координацию) не пересоберёт цепь и не дольёт соседям недостающие пакеты. Chain replication по своей природе CP: она не умеет деградировать в «запишем куда получится», и в этом её честность. Плата — латентность записи линейна по длине цепи: три звена по 2 мс дают 6 мс, а не 2 мс, как у кворума с параллельной рассылкой.

Конфликты: три честных ответа

Как только записи одного ключа принимаются больше чем в одном месте (несколько лидеров или leaderless), конкурентные записи неизбежны. «Конкурентные» здесь — строгий термин: две записи конкурентны, если ни одна не находится в причинно-следственной связи с другой, см. логические часы.

Ответов ровно три, и других не существует.

Ответ 1: LWW — выбрать одну и потерять другую

Last Write Wins: у каждой записи есть timestamp, побеждает больший. Это поведение Cassandra по умолчанию.

Вопрос на миллион: по чьим часам «последняя»? Ответ — по настенным часам того узла, который принял запись. А значит:

Сценарий отказа: на одном узле NTP скакнул на 3 секунды вперёд (или chrony применил step вместо slew после долгого рассинхрона). Все записи, принятые этим узлом, получают timestamp из будущего. В течение следующих трёх секунд любая корректная запись на других узлах молча проигрывает. Данные не «повреждены» — они просто отсутствуют. Диагностика: SELECT WRITETIME(column) FROM ... показывает timestamp из будущего. В логах нет ничего.

LWW безопасен ровно в двух ситуациях: у ключа один логический владелец (конкурентных записей физически не бывает), либо потеря записи приемлема по бизнесу (кэш, сессии, телеметрия последнего значения). Во всех остальных случаях LWW — это тихая потеря данных, встроенная в архитектуру.

Ответ 2: обнаружить конфликт и отдать его приложению

Честный путь: хранить обе версии и попросить приложение слить их. Для этого нужен способ отличать «новее» от «конкурентно» — версионные векторы.

Терминологическая тонкость, на которой все спотыкаются: векторные часы отслеживают события у процессов, версионные векторы — версии реплик объекта. Структура похожа, семантика разная. Dynamo и Riak используют версионные векторы.

VV = dict[str, int]  # реплика -> счётчик версий объекта

def descends(a: VV, b: VV) -> bool:
    """a доминирует b: всё, что видел b, видел и a."""
    return all(a.get(node, 0) >= counter for node, counter in b.items())

def concurrent(a: VV, b: VV) -> bool:
    """Ни одна версия не доминирует другую — это конфликт."""
    return not descends(a, b) and not descends(b, a)

def merge_vv(a: VV, b: VV) -> VV:
    """Поточечный максимум. Это join полурешётки: коммутативен,
    ассоциативен, идемпотентен."""
    return {n: max(a.get(n, 0), b.get(n, 0)) for n in a.keys() | b.keys()}

Сложность всех трёх операций — O(k) по времени и памяти, где k — число реплик, когда-либо писавших этот ключ. И вот здесь прячется практическая проблема: k растёт. В Riak вектор обрезают по размеру и возрасту (small_vclock, big_vclock, young_vclock, old_vclock). Обрезка безопасна в одну сторону: она может превратить причинно-связанные версии в «ложно конкурентные» (лишние siblings), но не может потерять данные.

Riak с allow_mult=true возвращает клиенту все конкурентные версии — siblings — и требует их слить. Каноничный пример из Dynamo — корзина покупок: слияние = объединение множеств товаров. Работает, пока пользователь только добавляет. Как только он что-то удалил, а параллельная реплика этого не видела, удалённый товар возвращается в корзину. Amazon описывал это как известное свойство системы; именно эта боль и мотивировала следующий ответ.

CRDT: третий ответ — сделать конфликт невозможным

Идея настолько простая, что кажется жульничеством: если операция слияния двух состояний коммутативна, ассоциативна и идемпотентна, то порядок доставки, дубликаты и повторы перестают иметь значение. Реплики сходятся к одному состоянию автоматически, без координации, без лидера, без консенсуса. Формально состояния образуют join-полурешётку, а merge — операцию наименьшей верхней границы; каждая локальная операция монотонно двигает состояние вверх по решётке. Каноническая работа — Shapiro, Preguiça, Baquero, Zawirski, A comprehensive study of Convergent and Commutative Replicated Data Types (INRIA RR-7506, 2011).

Два семейства различаются тем, что едет по сети, и требования к сети у них противоположные:

  • CvRDT (state-based). Передаём всё состояние, применяем merge. Каналу разрешено терять, дублировать и переставлять сообщения — идемпотентность и коммутативность это переваривают. Плата: трафик. Лечится delta-state CRDT — передаём только дельту решётки.
  • CmRDT (op-based). Передаём операции. Требуется доставка ровно один раз в причинном порядке: повторно применённый «инкремент» испортит счётчик. То есть сложность не исчезла, а переехала в слой доставки — и там она называется «exactly-once», о чём ниже.

Реализация: PN-Counter и OR-Set

from dataclasses import dataclass, field
from itertools import count


@dataclass
class PNCounter:
    """Счётчик с инкрементом и декрементом = два G-Counter'а: value =
    sum(P) - sum(N). Оба вектора только растут, merge — поточечный максимум."""
    node: str
    p: dict[str, int] = field(default_factory=dict)
    n: dict[str, int] = field(default_factory=dict)

    def inc(self, delta: int = 1) -> None:
        self.p[self.node] = self.p.get(self.node, 0) + delta

    def dec(self, delta: int = 1) -> None:
        self.n[self.node] = self.n.get(self.node, 0) + delta

    def value(self) -> int:
        return sum(self.p.values()) - sum(self.n.values())

    def merge(self, other: "PNCounter") -> None:
        for k in self.p.keys() | other.p.keys():
            self.p[k] = max(self.p.get(k, 0), other.p.get(k, 0))
        for k in self.n.keys() | other.n.keys():
            self.n[k] = max(self.n.get(k, 0), other.n.get(k, 0))


@dataclass
class ORSet:
    """Observed-Remove Set: add выигрывает у конкурентного remove. Каждое
    добавление получает уникальный тег; remove гасит только наблюдённые теги."""
    node: str
    _seq: count = field(default_factory=lambda: count(1))
    elements: dict[object, set[tuple[str, int]]] = field(default_factory=dict)
    tombstones: set[tuple[str, int]] = field(default_factory=set)

    def add(self, x) -> None:
        self.elements.setdefault(x, set()).add((self.node, next(self._seq)))

    def remove(self, x) -> None:
        self.tombstones |= self.elements.get(x, set())

    def value(self) -> set:
        return {x for x, tags in self.elements.items() if tags - self.tombstones}

    def merge(self, other: "ORSet") -> None:
        for x, tags in other.elements.items():
            self.elements.setdefault(x, set()).update(tags)
        self.tombstones |= other.tombstones

Проверка ключевого свойства OR-Set: реплика A добавляет элемент; реплика B параллельно его удаляет, не видя нового тега; после слияния элемент остаётся. Это и есть add-wins — и это осознанное решение о семантике, а не «правильный ответ».

Сложность: merge — O(|E| + |T|) по времени, где E — теги элементов, T — надгробия; память растёт с числом операций, а не с числом элементов. Отсюда главная эксплуатационная проблема наивного OR-Set: надгробия не собираются мусором. Практические реализации используют ORSWOT (OR-Set Without Tombstones), где вместо явных надгробий хранится версионный вектор реплики — память тогда пропорциональна числу реплик, а не числу удалений.

Границы применимости CRDT

Это не серебряная пуля, и её границы очень чёткие:

  • CRDT не выражают глобальные инварианты. «Баланс не уходит в минус», «мест не больше 100», «логин уникален» — это утверждения обо всех репликах сразу, они требуют координации. Никакая алгебра слияния их не даст; нужен консенсус, см. Paxos и Raft. Компромисс — bounded counter / escrow: каждой реплике заранее выдаётся квота, внутри квоты она работает без координации.
  • Сходимость — не корректность. Реплики гарантированно сойдутся к одному состоянию. Никто не обещал, что это состояние кто-то хотел. Одновременные «удалить товар» и «изменить количество» в add-wins-множестве дадут товар с новым количеством — формально корректно, для пользователя странно.
  • Метаданные растут. Теги, версионные векторы, надгробия. Для текстов (RGA, YATA) это особенно чувствительно; Automerge и Yjs потратили годы на компактное представление.
  • Семантику выбираете вы. add-wins или remove-wins — это продуктовое решение, замаскированное под структуру данных.

Где применяют в проде: Redis Enterprise CRDB (активная георепликация на CRDT), Riak DT, Azure Cosmos DB в multi-master (LWW по умолчанию плюс пользовательская процедура разрешения), Yjs и Automerge для совместного редактирования, Figma и Linear — для локальной офлайн-модели.

Как это устроено в реальных системах

etcd: репликация через консенсус

etcd не «реплицирует», а прогоняет каждую запись через Raft: запись закоммичена, когда её приняло большинство. Отсюда свойства: подтверждённая запись не теряется при отказе меньшинства, зато при потере кворума кластер останавливает запись — это выбор CP в терминах CAP. Типичные строки при проблемах:

etcdserver: request timed out, possibly due to previous leader failure
raft2026/03/14 02:11:04 INFO: 8e9e05c52164694d [term 47] received MsgTimeoutNow, starting campaign
etcdserver: read-only range request "key:\"/registry/pods\"" took too long (1.42s)

Последняя строка — самая частая в продакшене Kubernetes, и почти всегда она означает не проблему Raft, а медленный диск: fsync WAL на каждом коммите. Метрика etcd_disk_wal_fsync_duration_seconds важнее любой сетевой.

Kafka: ISR и три ручки, которые всё решают

Сценарий отказа: acks=all, replication.factor=3, min.insync.replicas=1. Продюсер уверен, что «all» значит «три копии». Реально all означает «все реплики из текущего ISR», а ISR может сжаться до одного лидера:

[2026-03-14 02:11:07,431] INFO [Partition orders-3 broker=2] Shrinking ISR from 2,3,1 to 2 (kafka.cluster.Partition)

С этого момента подтверждение приходит от единственной машины. Лидер падает — сообщения исчезают, продюсер не видел ни одной ошибки. Правило: min.insync.replicas = replication.factor − 1, то есть 2 при RF=3.

Второй сценарий: unclean.leader.election.enable=true. Лидером становится реплика вне ISR, у которой лога меньше. Разница удаляется:

[2026-03-14 02:12:44,902] WARN [ReplicaFetcher replicaId=1, leaderId=2] Truncating partition orders-3 to local high watermark 148812 (kafka.server.ReplicaFetcherThread)

Truncating ... to local high watermark — это буквально строка «удаляем подтверждённые сообщения ради доступности». Иногда это осознанный выбор (метрики), но он должен быть именно осознанным.

Cassandra, PostgreSQL, MongoDB, Spanner

Система Модель Ключевые ручки Где ломается
Cassandra leaderless, кворум на запрос LOCAL_QUORUM, gc_grace_seconds, max_hint_window_in_ms пропущенный repair → воскрешение удалённых; LWW + дрейф часов
PostgreSQL один лидер, WAL synchronous_commit, synchronous_standby_names, слоты репликации забытый слот копит WAL до ENOSPC и роняет лидера — ставьте max_slot_wal_keep_size
MongoDB replica set, oplog writeConcern {w:"majority", j:true}, readConcern w:1 + failover → файлы в rollback/
Kafka лидер партиции + ISR acks, min.insync.replicas, unclean election усечение лога, тихое сжатие ISR
Spanner Paxos-группа на каждый split TrueTime, commit-wait латентность записи упирается в границу неопределённости часов

Spanner стоит отдельного слова: там реплика — не «копия», а участник Paxos-группы, а линеаризуемость на глобальном масштабе достигается тем, что коммит ждёт максимальную неопределённость часов (обычно единицы миллисекунд) перед возвратом клиенту. Это редкий случай, когда физическое время встроено в протокол корректности, а не только в мониторинг. Первоисточник — Spanner: Google’s Globally-Distributed Database (OSDI 2012), продолжение темы — в времени и часах.

Почему exactly-once — обычно миф

Репликация — это частный случай доставки сообщений: «применить запись W на реплике R ровно один раз». Поэтому разговор про exactly-once начинается именно здесь.

Аргумент невозможности простой и не требует математики. Канал может потерять что угодно, в том числе подтверждение. Отправитель, не получив ack, не может отличить «сообщение не дошло» от «дошло и было применено, но ack потерялся». Значит у него ровно два варианта поведения:

  • не ретраить — получаем at-most-once, теряем сообщения;
  • ретраить — получаем at-least-once, получаем дубликаты.

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

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

  • Репликация БД. Реплика знает свою позицию в логе (LSN, GTID, oplog timestamp). Повторная отправка того же куска WAL просто игнорируется — применение идемпотентно по позиции. Это и есть дедупликация, встроенная в протокол.
  • Idempotent producer в Kafka. Каждое сообщение несёт тройку (producer_id, epoch, sequence_number) для партиции; брокер отбрасывает дубликат последовательности. Работает внутри Kafka и в пределах сессии продюсера. Это дедупликация в брокере, а не магия в сети.
  • Kafka transactions / EOS. Атомарность цепочки consume → process → produce внутри Kafka плюс read_committed у консьюмера. Ровно в тот момент, когда обработчик делает побочный эффект наружу — HTTP-вызов, письмо, списание, — гарантия заканчивается. Kafka не участвует в вашей транзакции с платёжным шлюзом.
  • Приложение. Ключ идемпотентности и эффект фиксируются одной транзакцией:
BEGIN;

-- 1. Пытаемся «занять» ключ идемпотентности
INSERT INTO processed_messages (idempotency_key, applied_at)
VALUES ($1, now())
ON CONFLICT (idempotency_key) DO NOTHING
RETURNING idempotency_key;
-- Вернулось 0 строк => сообщение уже применено:
-- ROLLBACK, подтверждаем оффсет, эффект НЕ повторяем.

-- 2. Иначе применяем эффект здесь же, в той же транзакции
UPDATE accounts SET balance = balance - $2 WHERE id = $3;

COMMIT;

Критично именно «в той же транзакции». Если отметка об обработке и сам эффект коммитятся раздельно, вы получили распределённую транзакцию с двумя участниками и все её проблемы — см. 2PC, saga и outbox.

Отдельная ловушка: op-based CRDT из предыдущего раздела требуют exactly-once в причинном порядке. То есть «CRDT избавляют от координации» верно только для state-based; op-based просто перекладывают ту же проблему на транспорт. Подробный разбор гарантий доставки — в идемпотентности и доставке.

Как выбирать

Типичные ошибки

  1. Считать реплику бэкапом. Репликация мгновенно повторяет DROP TABLE и повреждение логическим багом. Бэкап — это точка во времени и проверенное восстановление, а не «у нас есть standby».
  2. Мониторить конфиг вместо факта. semi-sync = ON в my.cnf не означает, что semi-sync работает прямо сейчас. Проверяйте состояние, а не намерение.
  3. Мерить лаг репликации в байтах. Разница LSN не говорит, сколько секунд отстаёт реплика: 1 МБ WAL может быть одной секундой, а может быть часом. Мерьте pg_last_xact_replay_timestamp() / Seconds_Behind_Source, и лучше — через heartbeat-таблицу с записью раз в секунду.
  4. acks=all без min.insync.replicas. Разобрано выше: «all» — это «все из ISR».
  5. Автоматический failover без фенсинга. Вы не ускорили восстановление, вы автоматизировали split brain.
  6. Считать, что кворум даёт линеаризуемость. w + r > n даёт пересечение множеств, и только его.
  7. LWW «по умолчанию, разберёмся потом». Не разберётесь: потерянные записи не оставляют следов ни в логах, ни в метриках.
  8. Repair как разовое действие. Не проходящий до конца nodetool repair — это отложенное воскрешение удалённых данных.
  9. Тестировать failover только «убив процесс». Реальные отказы — это частичное разделение сети, замедление в 50 раз и GC-паузы, а не kill -9. Про это — тестирование распределённых систем и Jepsen.

Мини-итог

  • Репликация — это договор о том, в какой момент запись считается существующей. Всё остальное — следствия.
  • Топология определяет жанр вашего инцидента: выборы лидера, разрешение конфликтов или тихое расхождение реплик.
  • Асинхронная репликация всегда имеет ненулевой RPO, и он измерим: расходящийся хвост WAL, файлы в rollback/, усечённые сегменты Kafka. Автоматический failover без фенсинга опаснее ручного — таймаут не отличает мёртвый узел от медленного.
  • w + r > n гарантирует пересечение множеств и ничего сверх этого; sloppy quorum отменяет и это. Chain replication достигает того же результата другой ценой — линейной латентностью записи и жёсткой CP-семантикой.
  • Конфликты решаются ровно тремя способами: потерять одну версию (LWW), отдать обе приложению (версионные векторы), сделать конфликт невозможным (CRDT). CRDT сходятся без координации, но не выражают глобальных инвариантов — для них нужен консенсус.
  • Exactly-once как свойство канала невозможен. Работающая замена — at-least-once плюс дедупликация в том же атомарном шаге, что и эффект.
  • Самые дорогие отказы репликации — молчаливые: деградировавший semi-sync, сжавшийся ISR, пропущенный repair. В логах их нет, поэтому они должны быть в метриках.

Источники

Что дальше

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

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

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

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

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