Репликация: лидер и последователи, кворумы, конфликты, 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
реплика её не видела L--xF: узел падает до отправки P->>P: health-check провален 3 раза подряд P->>F: promote F->>F: новая timeline 2, LSN на 8 КБ позади C->>F: SELECT order 4711 F-->>C: 0 rows Note over C,F: клиент получил 200 OK на запись,
которой больше не существует
Это не гипотетика — у потерянных байтов есть точный размер, и его можно посмотреть. Когда старый лидер вернётся и попробует подключиться к новому, 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 (n — число реплик ключа, w — сколько подтвердили запись, r — сколько опрошено при чтении) гарантирует ровно одно: множества записи и чтения пересекаются хотя бы по одному узлу, значит читатель увидит хотя бы одну копию последнего подтверждённого значения. Cassandra даёт это как настройку на каждый запрос: ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL.
И теперь — чего кворум НЕ даёт. Это самая частая ошибка в понимании Dynamo-подобных систем.
- Кворум — не линеаризуемость. Пересечение множеств гарантирует «увидим значение», а не «увидим одно и то же значение все и одновременно». Подробный разбор различий — в моделях согласованности.
- Запись на w узлов не атомарна. Пока она разъезжается, один читатель уже видит новое значение, а параллельный — ещё старое. Монотонности чтений тоже нет: следующий запрос того же клиента может попасть на отставший набор и вернуться назад во времени.
- Частично успешная запись не откатывается. Если при
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 просто перекладывают ту же проблему на транспорт. Подробный разбор гарантий доставки — в идемпотентности и доставке.
Как выбирать
подтверждённых записей?"} B -->|"нет"| D{"Записи одного ключа
приходят в разные регионы?"} C -->|"нет, RPO = 0"| E["Кворумная синхронная репликация:
Raft или Paxos — etcd, Spanner;
PG synchronous_standby_names ANY 1 из 2"] C -->|"да, RPO секунды"| F["Один лидер + async-реплики
+ обязательный фенсинг
+ мониторинг реального лага"] D -->|"нет"| G["Один лидер на регион по ключу:
это задача партиционирования"] D -->|"да"| H{"Слияние выражается
алгеброй?"} H -->|"да"| I["CRDT: счётчики, множества, тексты
Redis CRDB, Riak DT, Yjs, Automerge"] H -->|"нет"| J{"Потеря одной из
конкурентных записей допустима?"} J -->|"да"| K["LWW + строгий NTP
+ явный владелец ключа"] J -->|"нет"| L["Версионные векторы и siblings:
слияние в домене приложения"]
Типичные ошибки
- Считать реплику бэкапом. Репликация мгновенно повторяет
DROP TABLEи повреждение логическим багом. Бэкап — это точка во времени и проверенное восстановление, а не «у нас есть standby». - Мониторить конфиг вместо факта.
semi-sync = ONвmy.cnfне означает, что semi-sync работает прямо сейчас. Проверяйте состояние, а не намерение. - Мерить лаг репликации в байтах. Разница LSN не говорит, сколько секунд отстаёт реплика: 1 МБ WAL может быть одной секундой, а может быть часом. Мерьте
pg_last_xact_replay_timestamp()/Seconds_Behind_Source, и лучше — через heartbeat-таблицу с записью раз в секунду. acks=allбезmin.insync.replicas. Разобрано выше: «all» — это «все из ISR».- Автоматический failover без фенсинга. Вы не ускорили восстановление, вы автоматизировали split brain.
- Считать, что кворум даёт линеаризуемость.
w + r > nдаёт пересечение множеств, и только его. - LWW «по умолчанию, разберёмся потом». Не разберётесь: потерянные записи не оставляют следов ни в логах, ни в метриках.
- Repair как разовое действие. Не проходящий до конца
nodetool repair— это отложенное воскрешение удалённых данных. - Тестировать 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. В логах их нет, поэтому они должны быть в метриках.
Источники
- Leslie Lamport, Time, Clocks, and the Ordering of Events in a Distributed System (1978) — причинность и основание для версионных векторов.
- DeCandia et al., Dynamo: Amazon’s Highly Available Key-value Store (SOSP 2007) — кворумы, sloppy quorum, hinted handoff, деревья Меркла.
- Shapiro et al., A comprehensive study of Convergent and Commutative Replicated Data Types (2011) — формальная база CRDT.
- Ongaro, Ousterhout, In Search of an Understandable Consensus Algorithm — Raft (2014) — репликация лога через консенсус; Corbett et al., Spanner: Google’s Globally-Distributed Database (OSDI 2012) — TrueTime и commit-wait.
- van Renesse, Schneider, Chain Replication for Supporting High Throughput and Availability (OSDI 2004) и Object Storage on CRAQ (USENIX ATC 2009) — альтернатива кворумам.
- Martin Kleppmann, «Designing Data-Intensive Applications», главы 5 и 9; GitHub October 21 post-incident analysis — образцовый разбор реального split brain.
- Kafka replication design и PostgreSQL streaming replication — документация, которую стоит читать целиком.
Что дальше
Репликация отвечает на вопрос «сколько копий у одних и тех же данных». Следующий вопрос — «как разложить данные, которые не помещаются на один узел, и что делать, когда узлы приходят и уходят»: Партиционирование и шардирование: стратегии, ребалансировка, горячие ключи.