Распределённые системы Модели отказов: сбои узлов, сеть, разделение, византийские отказы
0%

Модели отказов: сбои узлов, сеть, разделение, византийские отказы

Модели отказов: сбои узлов, сеть, разделение, византийские отказы

Любой протокол распределённой системы — Raft, 2PC, gossip, ваш самописный failover на трёх bash-скриптах — доказуемо корректен только относительно предположений о том, как всё это может сломаться. Уберите предположение — и доказательство исчезнет вместе с ним. Поэтому первый вопрос при разборе распределённой системы звучит не «как она работает», а «при каких отказах она обещает работать и что происходит ровно на шаг за границей обещания».

Модель отказов — это страховой полис: в нём написано, что покрывается («узел может внезапно умереть») и что нет («узел не будет присылать заведомо ложные данные»). Проблема та же, что со страховкой: первую страницу читают все, раздел «исключения» — никто. А аварии живут в исключениях. Эта статья — про исключения: разберём иерархию моделей от самой удобной до самой параноидальной и для каждой покажем конкретный сценарий, который её ломает, и как он выглядит в логах etcd, Kafka, Cassandra, ZooKeeper, PostgreSQL. Закончим тем, почему «exactly-once» — маркетинговый слоган, а не свойство системы. Карта трека — в статье «Распределённые системы: карта трека». Прикладные средства защиты (размыкатели, переборки, ретраи с бюджетом) разобраны в «Устойчивость: circuit breaker, retry, bulkhead, backpressure»; здесь — фундамент под ними: без модели отказов вы не знаете, от чего именно защищаетесь.

1. Модель системы — это три независимые оси

Когда в статье пишут «алгоритм корректен», за этим стоят три отдельных набора допущений, и путают их чаще всего.

Ось 1 — модель отказов процессов. Что может сделать узел: остановиться? вернуться с частичной амнезией? замолчать выборочно? ответить не вовремя? соврать?

Ось 2 — модель синхронности. Есть ли известная верхняя граница на задержку сообщения и на скорость шага процесса?

  • Синхронная: границы известны и всегда соблюдаются, таймаут — идеальный детектор отказа. Так живут разве что системы жёсткого реального времени на выделенной шине.
  • Асинхронная: границ нет вообще. Отсюда результат FLP — «Impossibility of Distributed Consensus with One Faulty Process» (Fischer, Lynch, Paterson, JACM 1985): не существует детерминированного алгоритма консенсуса, гарантированно завершающегося, если хотя бы один узел может упасть. Не «сложно» — невозможно. Причина обидно проста: получатель не отличает «отправитель умер» от «сообщение ещё в пути».
  • Частично синхронная: границы существуют, но неизвестны, либо начинают соблюдаться после неизвестного момента GST. Формализация — «Consensus in the Presence of Partial Synchrony» (Dwork, Lynch, Stockmeyer, JACM 1988), и это та модель, в которой живут все ваши продакшн-системы: раз в час случается GC-пауза на восемь секунд, но не навсегда.

Практический смысл частичной синхронности стоит выучить наизусть: Raft и Paxos гарантируют безопасность (safety) всегда, а живучесть (liveness) — только в периоды, когда сеть и узлы ведут себя прилично. Это и есть легальный обход FLP: система имеет право «зависнуть» на время шторма, но не имеет права выдать два разных результата. Детали — в «Консенсус: Paxos, Raft, выбор лидера».

Ось 3 — модель коммуникации. Надёжны ли каналы, может ли сеть дублировать, переупорядочивать и портить сообщения, аутентифицированы ли отправители.

2. Иерархия моделей отказов процессов

Модели вложены друг в друга: каждая внешняя допускает всё, что допускает внутренняя, плюс что-то ещё. Чем шире модель, тем дороже протокол — больше узлов, раундов и криптографии.

Иерархия моделей отказов: crash-stop, crash-recovery, пропуски, тайминг, византийские

2.1 Crash-stop (fail-stop)

Узел останавливается и никогда не возвращается. Модель, в которой удобно доказывать теоремы, и которая почти никогда не верна на практике: упавший процесс обычно перезапускается супервизором через пять секунд. Честной она становится только там, где инстанс при сбое заменяется новым — с новым идентификатором и пустым диском; это инженерное решение подогнать реальность под удобную модель, и оно окупается.

Сценарий отказа. Под с in-memory кэшем сессий убит OOM-киллером и вернулся с тем же именем, но пустой памятью. Клиенты с живыми токенами получают 401, ретраят и укладывают сервис авторизации.

kernel:      Out of memory: Killed process 2471 (session-svc) total-vm:8123456kB
session-svc: session store initialized, 0 entries loaded
api-gw:      upstream returned 401 for 1873 requests in 12s (spike x900)
auth-svc:    p99 latency 40ms -> 2100ms, queue depth 8400

Диагностический признак — сервис зелёный по всем проверкам, а клиенты массово получают 401 сразу после рестарта; лечение не «добавить памяти», а признать, что состояние в памяти теряется при каждом рестарте, и сделать восстановление дешёвым.

2.2 Crash-recovery

Узел падает и возвращается, помня только то, что успел записать на устойчивый носитель. Это базовая модель для баз данных и брокеров: Kafka, PostgreSQL, etcd, ZooKeeper спроектированы вокруг цикла «упал → поднялся → доиграл журнал». Сразу возникает неочевидный вопрос: что вообще значит «записано на диск»? Ответ неприятный, см. раздел 8.

Сценарий отказа. Брокер Kafka подтвердил запись при acks=1 (лидер ответил до репликации), затем умер и вернулся, потеряв последние 200 мс. Клиент считает сообщение принятым — в топике его нет.

producer:  acked offset=91422 partition=orders-7
broker-2:  [ReplicaFetcher] Truncating partition orders-7 to offset 91380 (leader epoch 15)
broker-2:  Log truncated: 42 messages removed

Строка Truncating partition ... to offset — это буквально «мы стираем то, что кому-то уже подтвердили». Она выглядит рутинно, потому что усечение лога при смене лидера штатно, — и потому её пропускают. Лечится acks=all плюс min.insync.replicas >= 2; цена — рост p99 записи на внутрикластерный RTT. Второй сценарий, который забывают тестировать: узел вернулся не через пять секунд, а через сорок минут — со старым состоянием и старой эпохой; если протокол не отвергает такие сообщения по эпохе, «воскресший» узел начинает распространять устаревшие данные как свежие.

2.3 Отказы-пропуски (omission)

Узел жив и исполняет код, но часть сообщений не уходит (send-omission) или не доходит (receive-omission). Это не поломка сети: это переполненный буфер сокета, забитый accept-backlog, тред-пул, не успевающий вычитывать соединение, дроп на сетевой карте (nstat -az | grep -E 'ListenDrops|BacklogDrop' — первое, что смотрят).

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

node-3: recv queue full for peer node-1, dropping heartbeat (drops=3187)
node-3: no heartbeat from node-1 in 1.6s, transitioning FOLLOWER -> CANDIDATE
node-1: received RequestVote from node-3 term=42 (current term=41), stepping down
node-3: election timeout, term=43

Кластер входит в цикл переголосований и перестаёт коммитить что-либо, хотя все узлы живы с точки зрения ps и все health-check зелёные; признак в метриках — растущий term/epoch при нулевой пропускной способности записи.

2.4 Тайминговые отказы

Ответ верный, но слишком поздний — или слишком ранний, если разъехались часы. Источники: паузы сборщика мусора, CPU throttling из-за cgroup-квоты, шумные соседи по гипервизору, медленный fsync, распухшие буферы на сетевом пути.

Сценарий отказа, ставший классикой: лидер уходит в stop-the-world паузу на восемь секунд, кластер объявляет его мёртвым и выбирает нового, старый лидер просыпается, не зная, что прошло восемь секунд, и дописывает операцию с правами лидера. Два лидера пишут одновременно.

[GC pause] Total time for which application threads were stopped: 8.42 seconds
etcd:  failed to send out heartbeat on time (exceeded the 100ms timeout for 8.1s)
etcd:  lost the TCP streaming connection with peer 8211f1d0f64f3269
etcd:  the clock difference against peer is too high [1.4s > 1s]
app:   lease 32695fd6a1b2 expired while request in flight

failed to send out heartbeat on time — самая частая строка в аварийных логах etcd, и почти всегда она означает не сеть, а диск или CPU: fsync WAL не уложился в бюджет (проверяется по etcd_disk_wal_fsync_duration_seconds). Отсюда требование держать etcd на выделенных SSD с p99 fsync меньше 10 мс. Тайминговый отказ — место, где первая ось модели упирается во вторую: длинная пауза неотличима от смерти, поэтому любое решение, принятое по таймауту, обязано быть безопасным при ошибке (раздел 6). Природа расхождения часов — в «Время в распределённой системе».

2.5 Произвольные (византийские) отказы

Узел может слать разным получателям противоречивые ответы, подписываться чужим идентификатором, возвращать формально корректные, но неверные данные. Раздел 7 целиком об этом — и о том, что для этого не нужен злоумышленник.

Модель Что предполагаем Минимум узлов для f отказов Кто так живёт
Crash-stop упал — исчез навсегда n ≥ f + 1 для доступности чтения эфемерные воркеры, stateless-поды
Crash-recovery упал, вернулся, помнит только диск n ≥ 2f + 1 для консенсуса etcd, ZooKeeper, Kafka, PostgreSQL
Пропуски теряются отдельные сообщения n ≥ 2f + 1 плюс повторы и дедуп все реальные сети
Тайминговые задержки непредсказуемы плюс лизы и ограждение все, кто использует таймауты
Византийские узел может врать произвольно n ≥ 3f + 1 и криптоподписи блокчейны, авионика, межорганизационные системы

Граница n ≥ 3f + 1 — не эвристика, а доказанный минимум из «The Byzantine Generals Problem» (Lamport, Shostak, Pease, TOPLAS 1982) для случая без цифровых подписей.

3. Что на самом деле делает сеть

«Восемь заблуждений распределённых вычислений» (Peter Deutsch, дополнены James Gosling) — сеть надёжна; задержка нулевая; полоса бесконечна; сеть безопасна; топология неизменна; администратор один; транспорт бесплатен; сеть однородна — до сих пор лучший чек-лист для ревью (разбор Arnon Rotem-Gal-Oz). Что сеть делает вместо этого: теряет пакеты (TCP это скрывает ценой задержки — RTO превращает запрос с медианой 2 мс в «ответ через 3 секунды», и для вызывающего это тайминговый отказ, а не потеря); дублирует (ретрансмиссия TCP плюс ретрай приложения плюс повтор балансировщика — фундамент раздела 10); переупорядочивает (внутри соединения порядок есть, между соединениями пула — нет); задерживает произвольно долго — оценок сверху не существует, механика в «Сетевом стеке ОС»; портит данные — контрольная сумма TCP всего 16 бит, и в «When the CRC and TCP Checksum Disagree» (Stone, Partridge, SIGCOMM 2000) измерено: один из 1100–32000 повреждённых пакетов проходит проверку. На 10 Гбит/с это несколько испорченных сообщений в сутки, так что прикладной CRC32C у Kafka и чексуммы HDFS/ZFS — не паранойя, а арифметика. Эмпирика вместо ощущений: «The Network is Reliable» (Bailis, Kingsbury, ACM Queue 2014) и «An Analysis of Network-Partitioning Failures in Cloud Systems» (Alquraan et al., OSDI 2018) — разбор 136 реальных инцидентов в MongoDB, Kafka, Elasticsearch, HBase, Cassandra. Два вывода стоит помнить наизусть: около 80 % последствий разделений «тихие» — потеря данных, застрявшие операции, нарушенные инварианты, а не честный отказ обслуживания (система не падает, она врёт); и большинство катастрофических последствий вызвано разделениями вокруг одного узла, а не сценарием «отвалилась половина ДЦ». Отсюда контринтуитивное правило: страшно не большое разделение, а маленькое и незамеченное.

4. Детектор отказов: почему нельзя отличить мёртвого от медленного

Ключевая мысль всей теории, и принять её надо буквально: в асинхронной сети нет способа отличить упавший узел от медленного узла или медленной сети. Единственный доступный инструмент — таймаут, а таймаут это ставка, а не измерение. Chandra и Toueg в «Unreliable Failure Detectors for Reliable Distributed Systems» (JACM 1996) формализовали детектор через полноту (каждый упавший узел рано или поздно заподозрен) и точность (живые не подозреваются или перестают подозреваться со временем). Реально достижимый класс — ◇P: детектор может ошибаться сколько угодно, но после GST перестаёт. Этого достаточно для консенсуса — вот почему ложное подозрение лидера в Raft не ломает безопасность, а лишь тратит время на лишние выборы.

Состояние Joining не декоративное: возврат обязан происходить с новой эпохой, иначе задержавшиеся в сети сообщения прошлой жизни узла примут за свежие.

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

Система Параметр По умолчанию
etcd heartbeat-interval / election-timeout 100 мс / 1000 мс
Kafka (потребитель) session.timeout.ms 45 с
Kafka (брокер ↔ контроллер) broker.session.timeout.ms 9 с
ZooKeeper tickTime × syncLimit 2 с × 5
Cassandra phi_convict_threshold 8
Kubernetes node-monitor-grace-period около 40 с

Разброс в два порядка — не небрежность, а разная цена ошибки: лишние выборы в etcd стоят миллисекунды, лишнее выселение подов с узла — минуты работы и всплеск нагрузки на планировщик («Kubernetes»).

φ-accrual: вместо «да/нет» — степень подозрения

Идея из «The φ Accrual Failure Detector» (Hayashibara et al., SRDS 2004), реализованная в Cassandra и Akka: детектор выдаёт не булево значение, а число φ, растущее со временем тишины, а приложение само выбирает порог под цену действия — для дешёвого (перестать слать новые запросы) φ = 5, для дорогого (смена лидера, перебалансировка) φ = 12.

heartbeat(t):  добавить (t - время_предыдущего) в скользящее окно W; время_предыдущего = t
phi(now):      mean = среднее(W); std = max(отклонение(W), min_std)
               P = вероятность того, что интервал из N(mean, std) окажется > now - время_предыдущего
               вернуть -log10(P)
import math
import statistics
from collections import deque


class PhiAccrualDetector:
    """φ = 8 означает вероятность ошибки ~10^-8 при наблюдавшемся распределении задержек."""

    def __init__(self, window: int = 1000, min_std_ms: float = 50.0) -> None:
        self.intervals: deque[float] = deque(maxlen=window)
        self.last_arrival: float | None = None
        # Нижняя граница σ обязательна: в "слишком стабильной" сети дисперсия
        # стремится к нулю, и любая микрозадержка мгновенно даёт φ = 100.
        self.min_std_ms = min_std_ms

    def heartbeat(self, now_ms: float) -> None:  # O(1)
        if self.last_arrival is not None:
            self.intervals.append(now_ms - self.last_arrival)
        self.last_arrival = now_ms

    def phi(self, now_ms: float) -> float:  # O(W), в проде — O(1) через суммы Уэлфорда
        if self.last_arrival is None or len(self.intervals) < 2:
            return 0.0
        mean = statistics.fmean(self.intervals)
        std = max(statistics.pstdev(self.intervals, mu=mean), self.min_std_ms)
        # P(следующий интервал окажется больше elapsed) для нормального N(mean, std)
        p_later = 0.5 * math.erfc((now_ms - self.last_arrival - mean) / (std * math.sqrt(2.0)))
        return -math.log10(max(p_later, 1e-300))

Сложность. heartbeatO(1); phi в этой наивной версии O(W), а с инкрементальными суммами (алгоритм Уэлфорда, устойчивый к катастрофическому сокращению) — O(1). Память O(W) на пира: окно на 1000 интервалов это около 8 КБ, на кластере в 500 узлов — примерно 4 МБ. Главная ценность подхода — адаптивность: в дата-центре с медианой 1 мс порог φ = 8 сработает через ~200 мс тишины, на межрегиональном линке с медианой 80 мс — через несколько секунд. Одна константа вместо ручной настройки таймаутов под каждую площадку.

Сценарий отказа самого детектора. Cassandra в облаке с шумными соседями: гипервизор периодически крадёт CPU, интервалы heartbeat скачут, узлы помечают друг друга мёртвыми по кругу, растёт очередь hinted handoff.

Gossiper: InetAddress /10.0.3.17 is now DOWN
Gossiper: InetAddress /10.0.3.17 is now UP
StorageService: Not marking nodes down due to local pause of 12043771000ns > 5000000000ns

Строка Not marking nodes down due to local pause — встроенная защита: узел заметил, что тормозил сам, и отказался обвинять соседей. Образцовое решение, которое стоит копировать: прежде чем объявить кого-то мёртвым, проверь, не был ли мёртв ты.

5. Разделение сети: четыре разных зверя

«Partition» обычно понимают как «сеть разорвалась пополам». В реальности разделений четыре вида — и опасны не те, о которых думают.

Четыре вида сетевых разделений: полное, частичное, одностороннее, серый отказ

1. Полное разделение. Кластер распадается на непересекающиеся группы. Это, как ни странно, хороший случай: кворум большинства решает проблему математически — две группы не могут одновременно содержать большинство, меньшинство обязано отказаться от записи. Именно этот случай разбирает теорема CAP («CAP и PACELC»).

2. Частичное разделение. A видит B и C, а B и C не видят друг друга. Кворум не спасает автоматически: B и C по очереди объявляют себя кандидатами, каждый собирает голос A, term растёт, коммитов нет. OSDI-2018 называет этот сценарий самым разрушительным.

node-b: became candidate at term 4471
node-c: became candidate at term 4472
node-a: [term 4472] received MsgVote from node-c, stepping down
node-b: became candidate at term 4473

Растущий на тысячи term при живых узлах — самый узнаваемый след частичного разделения. Лечение — pre-vote (перед повышением term узел проверяет, что реально дотягивается до большинства; в etcd включено по умолчанию) и CheckQuorum (лидер сам уходит в follower, если не получил ответов от большинства за период выборов).

3. Одностороннее (simplex) разделение. Пакеты A→B доходят, B→A нет: асимметричные правила firewall, сломанный ECMP-путь, разный MTU в одну сторону, зависший conntrack. A уверен, что всё в порядке — он же отправляет и не видит ошибок, — а B считает A мёртвым. Самый неприятный вид для отладки, потому что «пинг с A проходит»; в логах он опознаётся по тому, что одна сторона печатает heartbeat sent ... (no error), а вторая в это же время no heartbeat ... marking DOWN.

4. Серый отказ (gray failure). Связь формально жива — TCP-проверка отвечает, ping проходит, — но 2 % пакетов теряется и p99 равен 14 секундам. Автоматика молчит, потому что бинарная проверка зелёная, а система стоит. Это суть «Gray Failure: The Achilles’ Heel of Cloud-Scale Systems» (Huang et al., HotOS 2017): расхождение между тем, что видит мониторинг, и тем, что видит клиент.

lb:      upstream node-5 health=OK (GET /healthz 200, 3ms)  <-- статический обработчик
client:  p99 for node-5 = 14200ms, error rate 2.1%, others p99 = 45ms
node-5:  connection pool to postgres exhausted (in_use=200/200, waiters=1840)

Отсюда правило в стандарт сервиса: health-check должен измерять то же, что делает клиент — сквозной запрос с проверкой зависимостей и латентности, а не статический 200; и вывод узла из ротации должен решаться по клиентским перцентилям, а не по бинарной пробе. Разбор, который стоит прочитать целиком, — постмортем GitHub от 21 октября 2018 года: 43 секунды разделения между площадками и 24 часа деградации, потому что автоматический failover MySQL успел переключить мастер, а свести записи обратно автоматически было нельзя. Главный урок: длительность аварии определяется не длительностью разделения, а стоимостью восстановления консистентности после него.

6. Split-brain и ограждение (fencing)

Разделение опасно тем, что порождает два узла, одновременно считающих себя владельцами ресурса. Контринтуитивный момент: распределённая блокировка с TTL сама по себе не защищает. Клиент взял лок, ушёл в GC-паузу дольше TTL, лок достался другому, первый проснулся и продолжил писать, будучи уверенным, что владеет им.

Fencing token — монотонно растущее число, выдаваемое вместе с блокировкой; ресурс обязан отвергать любую операцию с токеном меньше уже виденного. Без этой проверки блокировка — вежливая просьба, а не гарантия. Канонический разбор — у Мартина Клеппманна, «How to do distributed locking»; про лизы и сессии — «Координация: распределённые блокировки, лизы». В готовых системах роль токена уже играют term и revision в Raft/etcd, leader epoch и producer epoch в Kafka (старый транзакционный продюсер получает ProducerFencedException), явный fencing NameNode в HDFS вплоть до sshfence/STONITH, Lease API и resourceVersion в Kubernetes.

-- Ресурс хранит наибольший увиденный токен и отвергает всё, что старее.
-- Ноль изменённых строк = запись отвергнута; это НЕ ошибка сети и не повод для ретрая.
UPDATE resource_state
SET    payload    = :new_payload,
       last_token = :token
WHERE  resource_id = :id
  AND  last_token < :token;
// Главное не «взял лок и работаю», а «контекст умирает вместе с лизой,
// а ресурс сам проверяет монотонность токена».
func (w *Worker) processWithFencing(ctx context.Context, job Job) error {
    sess, err := concurrency.NewSession(w.cli, concurrency.WithTTL(10))
    if err != nil {
        return err
    }
    defer sess.Close()

    mu := concurrency.NewMutex(sess, "/locks/"+job.ResourceID)
    if err := mu.Lock(ctx); err != nil {
        return err
    }
    token := mu.Header().Revision // монотонен в пределах кластера etcd

    ctx, cancel := context.WithCancel(ctx)
    defer cancel()
    go func() {
        <-sess.Done() // лиза потеряна: сеть, пауза, перезапуск etcd
        cancel()      // работа обязана прерваться, а не идти «по инерции»
    }()
    return w.store.WriteFenced(ctx, job.ResourceID, job.Payload, token)
}

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

worker-a: acquired lock order/8891 lease=32695fd6a1b2 ttl=10s, then [pause 14.2s]
worker-b: acquired lock order/8891 lease=32695fd6a1c7 ttl=10s
worker-b: charge 500 USD -> ok, payment_id=p-77120
worker-a: charge 500 USD -> ok, payment_id=p-77121   <-- второе списание, ошибок нет

7. Византийские отказы: не только про злоумышленников

Византийский отказ — произвольное поведение узла: противоречивые ответы разным получателям, чужая подпись, синтаксически корректные, но неверные данные. Термин из работы Лэмпорта 1982 года. Без подписей задача разрешима только при n ≥ 3f + 1: при n = 3f предатели делят лояльных на две равные группы, давая каждой свою версию, и большинство не образуется. Практичный асинхронный алгоритм — PBFT (Castro, Liskov, OSDI 1999): три раунда, O(n²) сообщений на операцию; из-за этого квадрата BFT не применяют там, где можно обойтись. В византийской модели реально работают блокчейны (Bitcoin — вероятностный BFT через PoW, Tendermint/CometBFT — классический), межорганизационные системы без общего администратора, авионика. А etcd, Kafka, Cassandra, Spanner — нет: они предполагают, что узлы падают и молчат, но не врут. Это осознанный выбор, и цена BFT именно такова: втрое больше узлов, криптография на каждом сообщении, квадратичный трафик.

Но «византийское» случается и без злоумышленников

Самый известный пример — авария Amazon S3 20 июля 2008 года: один перевёрнутый бит в сообщении gossip-протокола, не покрытый контрольной суммой, привёл к распространению повреждённого состояния между серверами. Лечилось часами, потому что заражённое состояние размножал сам протокол. Это ровно византийский отказ: узел рассылал корректные по формату, но неверные данные — и остальные ему верили. Флипы битов не экзотика: «DRAM Errors in the Wild» (Schroeder, Pinheiro, Weber, SIGMETRICS 2009, парк серверов Google) — более 8 % модулей DIMM видят хотя бы одну исправимую ошибку в год; ECC ловит одиночные ошибки, но не везде: данные проходят через регистры, сетевую карту, PCIe-транзит. Третий источник — сами системы: в «Redundancy Does Not Imply Fault Tolerance» (Ganesan et al., FAST 2017) авторы вносили единичное повреждение блока в Redis, ZooKeeper, Cassandra, Kafka и RethinkDB — во многих случаях это приводило не к локальной ошибке, а к распространению порчи на здоровые реплики или к необнаруженной потере данных. Формулировка, которую стоит запомнить: репликация без проверки целостности размножает мусор.

Что делают вместо полноценного BFT в обычных системах:

  1. Сквозные контрольные суммы — CRC32C в Kafka, чексуммы блоков HDFS, data_checksums в PostgreSQL, контроль блоков и метаданных в ZFS. По сквозному аргументу Зальцера проверять надо на границе приложения: TCP-сумма не защищает от порчи в памяти промежуточного прокси.
  2. Scrubbing — фоновое перечитывание реплик со сверкой; без него порча обнаружится в момент отказа основной копии, то есть в худший возможный момент.
  3. mTLS внутри кластера и валидация инвариантов на приёме. Реплика, получившая запись с эпохой из будущего или с offset, нарушающим монотонность, обязана отказать и поднять тревогу, а не «починить» у себя.
  4. Сверка (reconciliation) — периодическое сравнение хешей диапазонов Merkle-деревьями, как в Dynamo и Cassandra («Cassandra и wide-column»); находит расхождения, которые протокол репликации пропустил. В логах срабатывание этих механизмов выглядит как Record is corrupt, stored crc = 812739919, computed crc = 43991002 у Kafka или WARNING: page verification failed, calculated checksum 24998 but expected 63301 у PostgreSQL — и это хорошая новость: порчу заметили до того, как её размножили.

8. Отказы носителей: «записано на диск» — не бинарное свойство

Модель crash-recovery целиком опирается на понятие устойчивой записи, а диск отказывает частично, и это ломает допущение изнутри.

  • Латентные секторные ошибки. Bairavasundaram et al., SIGMETRICS 2007, 1,53 млн дисков NetApp: 3,45 % SATA-дисков за 32 месяца показали хотя бы одну латентную ошибку сектора. Сектор читается годами нормально, потом внезапно нет — и обнаруживается это при восстановлении, когда вторая копия уже потеряна.
  • Порванная запись (torn write). Питание пропало посреди записи страницы 8 КБ: часть новая, часть старая. Отсюда full_page_writes в PostgreSQL и double-write buffer в InnoDB.
  • Fsync, потерявший ошибку. Fsyncgate 2018: при ошибке отложенной записи Linux помечал страницу чистой и сообщал об ошибке только одному вызову fsync() — следующий возвращал успех, хотя данные не записаны. PostgreSQL в ответ стал делать PANIC: безопаснее упасть и восстановиться из WAL, чем продолжить с ложной уверенностью. Академический разбор — «Can Applications Recover from fsync Failures?» (ATC 2020); механика страничного кэша — в «Файловых системах».
kernel:   EXT4-fs warning (device nvme0n1p2): ext4_end_bio: I/O error 10 writing to inode
postgres: PANIC: could not flush dirty data: Input/output error

Вывод для модели отказов: «узел записал» ≠ «данные переживут перезапуск». Если протокол считает подтверждение записи фактом устойчивости, добавьте строку «носитель может потерять подтверждённую запись» и решите, что с этим делать: репликация с acks=all, контрольные суммы, регулярная проверка восстановления из бэкапов (бэкап, который никогда не восстанавливали, — это не бэкап, а гипотеза). Смежное — «Репликация: лидер и последователи, кворумы» и «Репликация и шардирование в СУБД».

9. Отказы коррелированы — «независимость» почти всегда ложь

Расчёты доступности исходят из независимости: три реплики по 99,9 % дают девять девяток. Это неверно катастрофически, потому что причины отказа общие для всех реплик сразу: общая инфраструктура (домен питания, стойка, коммутатор, зона доступности); общая конфигурация — плохой конфиг раскатывается на все узлы за секунды, это самый частый источник глобальных аварий в публичных постмортемах; общая версия кода — баг, срабатывающий на определённых данных, сработает во всех репликах одновременно, и репликация тут не помогает вовсе; общее время (истечение сертификата, переполнение счётчика, leap second наступают везде одновременно); общая нагрузка — отказ одной реплики перераспределяет её трафик на остальные, это механизм каскада, а не смягчения. Практические следствия: ячеистая (cell-based) архитектура и shuffle sharding, раскатка конфигурации волнами с автооткатом, канарейка не только для кода, но и для конфигов. И честный подсчёт: доступность системы ограничена сверху доступностью её общих компонентов — DNS, конфигурационного сервиса, аутентификации, панели управления облака. Как всё это собирается в решение о том, какую модель отказов закладывать:

10. Почему «exactly-once» — миф, и что делают вместо

Два генерала на холмах по разные стороны долины должны атаковать одновременно; связь — гонцы, которых могут перехватить. A шлёт «атакуем на рассвете» и ждёт подтверждения. Но B, отправив подтверждение, не знает, дошло ли оно, — значит, ждёт подтверждения подтверждения. И так бесконечно. При возможной потере сообщений консенсуса о совместном действии нельзя достичь за конечное число сообщений (задача сформулирована в 1975 году, название закрепил Джим Грей в 1978-м). Это не абстракция, а ваша ситуация при каждом сетевом вызове:

Отправитель, не получивший подтверждение, имеет ровно два варианта: не повторять (риск потери — at-most-once) или повторить (риск дубля — at-least-once). Третьего не существует, и никакой протокол его не создаёт.

Exactly-once delivery невозможна. Exactly-once processing — наблюдаемый эффект ровно один — достижима, но только за счёт дедупликации и идемпотентности на стороне получателя.

Развёрнутая аргументация — «You Cannot Have Exactly-Once Delivery» (Tyler Treat). А вот что на самом деле делает Kafka под вывеской EOS (KIP-98, разбор Confluent): продюсеру выдаются PID и монотонный sequence number на партицию, брокер помнит последние номера и отбрасывает дубли ретрансмиссии — это дедупликация, а не «доставка один раз»; транзакции делают атомарными запись в несколько партиций и коммит офсетов (шаблон consume-process-produce), а isolation.level=read_committed прячет незакоммиченное. Всё это верно только внутри Kafka: как только обработчик делает внешний побочный эффект — HTTP-вызов, письмо, списание с карты, — гарантия заканчивается, потому что внешняя система не участвует в транзакции.

producer: Got error produce response on topic-partition orders-3, retrying. Error: NETWORK_EXCEPTION
broker:   Detected duplicate sequence number 5518 from producerId 4021,
          returning existing offset 91422 (no append)

Первая строка — тот самый момент неопределённости, вторая — то, что делают вместо exactly-once: распознавание дубля по идентификатору отправителя и порядковому номеру. На прикладном уровне это выглядит так:

  • Ключи идемпотентности плюс таблица обработанных запросов с уникальным индексом: ответ на повтор — сохранённый исходный ответ, а не новое выполнение.
  • Дедупликация по идентификатору сообщения с окном, покрывающим максимальный интервал ретраев. Окно конечно — значит, дедуп по краю вероятностный; это надо записать явно, а не делать вид, что окно бесконечно.
  • Транзакционный outbox: изменение состояния и событие в одной локальной транзакции, отдельный процесс публикует с at-least-once, потребитель обязан быть идемпотентным.
  • Операции, идемпотентные по природе: SET x = 5 вместо x += 1, «перевести заказ в состояние PAID» вместо «списать 500».
  • Компенсации там, где эффект нельзя безопасно повторить: «Распределённые транзакции: 2PC, saga, outbox» и «Saga и распределённые транзакции».

Готовая формулировка для ADR: «Транспорт даёт at-least-once; ровно один наблюдаемый эффект обеспечивается идемпотентностью обработчика и дедупликацией по ключу X с окном Y; за пределами окна возможен повтор, последствия — Z». Детально — «Гарантии доставки и идемпотентность».

11. Как записать модель отказов своей системы

Модель отказов — не эссе на две страницы, а таблица допущений, каждое из которых проверяемо экспериментом.

Что предполагаем Что будет при нарушении Как обнаружим Как проверяем защиту
Узел возвращается, помня зафиксированный WAL потеря подтверждённых записей расхождение офсетов, Truncating partition kill -9 под нагрузкой плюс сверка
Часы расходятся не более чем на 500 мс нарушение порядка, «протухшие» лизы метрика clock offset, clock difference too high инъекция сдвига часов на стенде
Разделения полные, не частичные бесконечные выборы, коммитов нет рост term при нулевой записи iptables-разделение одного узла
Каналы двусторонние узел «жив для себя, мёртв для других» асимметрия в логах heartbeat одностороннее DROP-правило
Диск не возвращает битые данные порча размножается на реплики ошибки контрольных сумм инъекция порчи в блок и чтение
Клиенты повторяют запросы двойное списание дубли по ключу идемпотентности повтор одного запроса N раз
Отказы независимы одновременная потеря кворума инцидент учения с отключением целой AZ

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

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

  1. Считать таймаут измерением. Любое действие по таймауту — снять лидера, освободить лок — должно быть безопасным при ошибочной гипотезе. Отсюда fencing.
  2. Проверять живость там, где нужна полезность. /healthz со статическим 200 зелёный на узле, у которого забит пул соединений к базе.
  3. Строить кворум из чётного числа узлов. Четыре терпят столько же отказов, сколько три, но чаще ловят разделение пополам. Три, пять, семь.
  4. Разворачивать кворум в двух зонах доступности. При разделении между ними одна сторона всегда теряет кворум; нужна третья зона хотя бы под арбитра.
  5. Ретраить всё подряд без бюджета и джиттера — умножение нагрузки ровно тогда, когда её надо снижать; и особенно ретраить неидемпотентную операцию: повторный POST /payments без ключа идемпотентности это второе списание.
  6. Считать отказы независимыми. Три реплики одного образа с одним конфигом в одной стойке — это одна реплика с тремя счетами за электричество.
  7. Не тестировать возврат узла. Тестируют «убили узел» и не тестируют «узел вернулся через 40 минут с устаревшим состоянием и старой эпохой»; второе ломает системы чаще.
  8. Верить слову exactly-once, не читая области действия — почти всегда речь про внутренние границы конкретной системы и никогда про ваш вызов наружу.
  9. Логировать отказ без корреляционного идентификатора и эпохи. Отладить split-brain без term/epoch и trace-id невозможно: см. «Наблюдаемость распределённых систем».
  10. Молча «чинить» расхождение на приёме. Реплика, которая подстраивается под сообщение с эпохой из будущего вместо отказа и алерта, превращает локальную аномалию в распространяющуюся порчу.

13. Как проверяют модель отказов

Модель отказов — это утверждение, а утверждения проверяют экспериментом, а не обсуждением. Инъекция на уровне сети: разделение — iptables, задержки и потери — tc netem. Обязательно попробуйте односторонние правила: они находят то, чего не находят симметричные.

sudo iptables -A INPUT -s 10.0.2.11 -j DROP   # одностороннее разделение: A не слышит B
sudo tc qdisc add dev eth0 root netem loss 2% delay 40ms 300ms distribution pareto  # серый отказ
sudo iptables -D INPUT -s 10.0.2.11 -j DROP && sudo tc qdisc del dev eth0 root      # откат

Jepsen (jepsen.io/analyses) — фреймворк Кайла Кингсбери, проверяющий заявленные гарантии под разделениями, сдвигом часов и перезапусками; его отчёты — лучший в индустрии материал о том, как ломаются знакомые вам базы. Детерминированная симуляция — FoundationDB прогоняет весь кластер в одном процессе с виртуальным временем и управляемым генератором отказов (описание). Chaos engineering в проде — с гипотезой, ограниченным радиусом поражения и кнопкой отмены. Подробно — «Тестирование распределённых систем». Минимальный набор экспериментов, который окупается везде: kill -9 лидера под нагрузкой; пауза процесса дольше TTL лизы (kill -STOP, затем kill -CONT); одностороннее разделение одного узла; сдвиг часов на секунду вперёд и назад; заполнение диска до 100 %; возврат узла с устаревшим состоянием через час.

Итог

  • Модель отказов — это контракт, относительно которого корректен любой ваш протокол: что может сделать процесс, есть ли границы на время, что делает сеть. Без явного контракта нет и корректности.
  • Модели вложены: crash-stop ⊂ crash-recovery ⊂ пропуски ⊂ тайминг ⊂ произвольные. Продакшн живёт в частичной синхронности с crash-recovery, пропусками и тайминговыми отказами — и не в византийской модели.
  • Отличить мёртвый узел от медленного невозможно — отсюда FLP, детекторы класса ◇P, φ-accrual и правило «решение по таймауту должно быть безопасно при ошибке».
  • Опасны не большие разделения, а маленькие и незамеченные: частичные, односторонние, серые. Кворум спасает только от полных.
  • Split-brain лечится не блокировкой, а ограждением на стороне ресурса — монотонным токеном, который ресурс проверяет сам.
  • Византийские отказы случаются без злоумышленников — как порча данных: полноценный BFT почти никому не нужен, сквозные контрольные суммы нужны всем.
  • Отказы коррелированы, умножать «девятки» реплик друг на друга — самообман: общий конфиг и общая версия кода отказывают одновременно везде.
  • Exactly-once доставки не существует — есть at-least-once плюс идемпотентность и дедупликация с конечным окном, и это надо записать явно, а не надеяться на слово в документации брокера.

Источники


Что дальше

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

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

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

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

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