Распределённые системы: карта трека и почему всё ломается иначе
Эта статья — вход в трек. Она не пересказывает остальные четырнадцать материалов, а даёт то, без чего они рассыпаются на список алгоритмов: модель мира, в которой эти алгоритмы имеют смысл. После неё вы будете смотреть на любой распределённый сбой через одну линзу — «какое знание здесь оказалось устаревшим и кто на нём построил решение».
Трек про фундамент: отказы, время, согласованность, координация. Про стили построения систем — микросервисы, event-driven, saga как архитектурный шаблон — есть отдельный трек «Архитектурные паттерны». Здесь мы разбираем не «как нарезать сервисы», а почему любой разрез по сети приносит класс отказов, которого не было в одном процессе.
1. Определение, которое действительно работает
Самая точная формулировка принадлежит Лесли Лампорту (премия Тьюринга 2013 года, автор доброй половины фундамента этой области) — это его письмо 1987 года:
Распределённая система — это такая система, в которой отказ компьютера, о существовании которого вы даже не подозревали, делает ваш собственный компьютер непригодным к работе.
Шутка описывает следствие. Причина — три отсутствующие вещи:
- Нет общей памяти. Любое «общее знание» физически представлено копиями, а копии расходятся. То, что вы прочитали, — снимок прошлого, а не текущее состояние.
- Нет общих часов. У каждого узла свои кварцевые часы со своим дрейфом. Значит, нет объективного «одновременно» и нет объективного порядка событий между узлами.
- Нет общей судьбы. Компоненты падают по отдельности. Это и есть частичный отказ (partial failure) — главный водораздел.
Именно частичный отказ, а не количество машин, делает систему распределённой. Многопоточная программа на одной машине распределённой не является: умер процесс — умерли все потоки сразу, общая память есть, часы одни. Как только между двумя частями появляется сеть, возникает четвёртый, самый неприятный факт:
- Отказ и медленность неразличимы. Узел, который не ответил, мог упасть, мог зависнуть в паузе сборщика мусора, мог ответить — но ответ потерялся. Различить эти случаи принципиально нельзя без верхней границы на задержку сети.
| Свойство | Один процесс | Распределённая система |
|---|---|---|
| Вызов функции | вернул или процесс упал | вернул, упал или неизвестно |
| Время | один монотонный счётчик | N часов с взаимным дрейфом |
| Состояние | одна копия | N копий, каждая устарела по-своему |
| Отказ | всё или ничего | любое подмножество компонентов |
| Отладка | стек-трейс | корреляция логов N сервисов |
| Воспроизводимость теста | детерминированная | зависит от расписания и сети |
Дальше весь трек — способы жить с пунктами 1–4, не разрушив бизнес-инварианты.
2. Карта трека
- 01. Модели отказов — краш, omission, timing, византийские отказы; что каждая модель разрешает предполагать.
- 02. Время и часы — дрейф физических часов, логические часы Лампорта, векторные и гибридные метки.
- 03. CAP и PACELC — что теорема утверждает, почему «AP-база» обычно маркетинг и зачем PACELC.
- 04. Модели согласованности — линеаризуемость, causal, read-your-writes, eventual: иерархия и цена.
- 05. Репликация — лидер и последователи, кворумы, конфликты, CRDT, anti-entropy.
- 06. Партиционирование — хеш и диапазоны, консистентное хеширование, ребалансировка, горячие ключи.
- 07. Консенсус — Paxos, Raft, выбор лидера, репликация лога, смена конфигурации.
- 08. Распределённые транзакции — 2PC и его блокировки, saga, outbox, компенсации.
- 09. Доставка и идемпотентность — at-most-once, at-least-once, миф exactly-once, ключи идемпотентности.
- 10. Координация — распределённые блокировки, лизы, fencing-токены, etcd и ZooKeeper.
- 11. Очереди и потоки — брокеры, порядок, партиции, backpressure, повторы и DLQ.
- 12. Наблюдаемость — трассировка, корреляция, поиск причины в системе без единого стека.
- 13. Тестирование — chaos engineering, Jepsen, детерминированная симуляция.
Маршруты чтения. Разбираетесь с дублями в очереди — 09 → 11 → 08. Проектируете хранилище — 04 → 05 → 06 → 03. Пишете сервис, который держит лидерство или блокировку — 07 → 10 → 02. Пришли на разбор инцидента — 01 → 12 → 13. Готовитесь к system design — 03 → 04 → 05 → 07.
Соседние треки: репликация и шардирование в СУБД и NewSQL и распределённые базы — те же идеи со стороны хранилищ; шаблоны устойчивости — прикладная сторона (circuit breaker, bulkhead, retry); потоковая обработка — окна, watermark, состояние; сетевой стек ОС — что на самом деле делает TCP под вашим RPC. Протоколы уровня «как работает BGP и почему падает DNS» — в треке «Сети» портала.
3. Третий исход: почему catch не спасает
В одном процессе у вызова два исхода. По сети — три, и третий не имеет обработчика в вашем языке.
Списание произошло? Неизвестно. P->>A: retry: charge(card, 4990) A->>D: INSERT charge id=ch_7bd D-->>A: OK, committed A-->>P: 200 OK P-->>C: 200 OK Note over C,D: Клиент оплатил один раз. Списаний — два.
Как это выглядит в логах — почти всегда одинаково:
14:03:22.118 WARN payment-api upstream=acquirer op=charge timeout=3000ms trace=7f3a91c2 attempt=1
14:03:22.119 INFO payment-api retrying op=charge attempt=2 backoff=412ms trace=7f3a91c2
14:03:25.204 INFO payment-api upstream=acquirer op=charge status=200 latency=3085ms trace=7f3a91c2
Со стороны payment-api инцидента нет: одна ошибка, один успешный повтор, метрики зелёные. Инцидент виден только у эквайера — два INSERT с разными идентификаторами, одинаковой картой, суммой и близкими метками времени. Это фундаментальное свойство: отказ регистрируется не там, где он произошёл.
Таймаут не означает «не выполнилось». Таймаут означает «я не знаю».
Отсюда правило для код-ревью: любой сетевой вызов, меняющий состояние, обязан быть либо идемпотентным, либо сопровождаться механизмом выяснения исхода — запросом статуса по бизнес-ключу. Варианта «повторим и понадеемся» не существует, он просто откладывает разбор на бухгалтерию.
4. Тайм-аут — это решение, а не наблюдение
Точное обнаружение отказов в асинхронной системе невозможно. Формально это описано у Чандры и Туэга в «Unreliable Failure Detectors for Reliable Distributed Systems» (JACM, 1996): детектор отказов — оракул, который подозревает узлы, и его качество описывается полнотой (рано или поздно заподозрит всех умерших) и точностью (не будет вечно подозревать живых). Ни один реальный детектор не даёт обоих свойств в сильной форме.
На практике детектор — это таймаут, и его настройка есть торг. Слишком короткий → ложные срабатывания: живой узел объявляется мёртвым, начинаются перевыборы лидера и ребалансировка партиций — под той самой нагрузкой, которая и вызвала задержку. Слишком длинный → система висит и не деградирует, пользователь смотрит на спиннер.
Классика ложного срабатывания — пауза сборщика мусора: узел жив, но не выполняет код полторы секунды.
09:41:02.882 [gc] Pause Full (Allocation Failure) 6144M->1082M(8192M) 1482.117ms
09:41:04.371 WARN zk-client Client session timed out, have not heard from server in 11402ms
for sessionid 0x17b1c9f4a2e0003, closing socket connection
09:41:04.372 INFO worker-7 lock/orders-shard-3 lost, releasing leadership
09:41:04.402 INFO worker-2 acquired lock/orders-shard-3, becoming owner
09:41:05.853 INFO worker-7 resumed, continuing batch flush for orders-shard-3
Последняя строка — баг стоимостью в инцидент: worker-7 очнулся и продолжил писать в шард, которым уже владеет worker-2. Лечится не увеличением таймаута, а fencing-токеном (раздел 10 и статья «Координация»). Канонический разбор — «How to do distributed locking» Мартина Клеппманна.
(ложное срабатывание) Подозреваемый --> ОбъявленМёртвым: истёк порог детектора ОбъявленМёртвым --> Восстанавливается: узел вернулся,
но роль уже отобрана Восстанавливается --> Живой: догнал лог, принял новый term ОбъявленМёртвым --> [*]: узел действительно умер note right of Подозреваемый Это состояние существует только в голове наблюдателя. Сам узел о нём не знает. end note
Более умный вариант — phi accrual failure detector (Hayashibara et al., 2004): вместо булева «жив/мёртв» он выдаёт непрерывную величину подозрения на основе распределения исторических интервалов между heartbeat’ами. Используется в Cassandra (phi_convict_threshold) и Akka Cluster. Проблему он не решает, но адаптируется к тому, что «нормальная» задержка в 03:00 и в час пик — разные величины.
Отдельно стоит запомнить категорию серых отказов (gray failure): узел отвечает на heartbeat, отдаёт 200 OK на health-check, но при этом обрабатывает 3% запросов или пишет на диск с задержкой в секунды. Детектор его не убивает, балансировщик продолжает слать трафик, а SLO деградирует у всех. Разбор феномена — «Gray Failure: The Achilles’ Heel of Cloud-Scale Systems» (HotOS, 2017). Практический вывод: health-check обязан проверять тот же путь, что и рабочий трафик (включая доступ к БД), иначе он проверяет только то, что процесс жив.
5. Шкала задержек: почему средние бесполезны
Граница «локально / по сети» — это смена порядка величины, а не удобства. Обращение в память ~100 нс, RPC внутри дата-центра ~0.5 мс: разница в 5000 раз. Рефакторинг, «просто выносящий модуль в сервис», умножает стоимость обращения на четыре порядка; если модуль вызывался в цикле, вы получите N+1 по сети — самый частый способ уронить latency после распила монолита.
У сетевого вызова нет «типичного» времени. Есть распределение с тяжёлым хвостом, и хвост формируют очереди: на сетевой карте, в пуле потоков, на диске, пауза GC, ретрансмиссия TCP. Джефф Дин описал следствие в «The Tail at Scale» (CACM, 2013): если запрос пользователя порождает 100 параллельных обращений, p99 каждого становится типичным временем ответа страницы. Практический минимум: измеряйте p50/p95/p99/p99.9, а не среднее; ставьте таймаут от измеренного p99 с запасом; следите, чтобы клиентский таймаут был больше суммы серверных — иначе клиент отваливается, когда сервер ещё честно работает, и вы получаете нагрузку без единого успешного ответа.
6. Восемь заблуждений — и что ломается за каждым
Список сформулирован Питером Дойчем и Джеймсом Гослингом в Sun; разбор — The Eight Fallacies of Distributed Computing.
- Сеть надёжна. Бэйлис и Кингсбери в «The Network is Reliable» (ACM Queue, 2014) собрали односторонние разделения (A видит B, B не видит A) и «серые» отказы с потерей 2% пакетов. Сценарий: продюсер Kafka с
acks=allна полудохлом соединении получаетTimeoutException: Expiring 47 record(s) for orders-3: 30000 ms has passed since batch creation, но часть предыдущих батчей прошла — приmax.in.flight.requests.per.connection > 1без идемпотентного продюсера порядок в партиции нарушен. - Задержка равна нулю.
GET /order/{id}последовательно ходит в profile, pricing, inventory, loyalty; у каждого p99 = 300 мс — p99 страницы больше секунды, хотя «все сервисы в SLA». - Пропускная способность бесконечна. Ответ вырос с 40 КБ до 400 КБ после добавления вложенной коллекции. Внутри ДЦ незаметно; на межрегиональном канале упирается в полосу, и растёт p99 у всех соседей по каналу.
- Сеть безопасна. Отказ выглядит не как ошибка, а как успешный запрос из неожиданного места: без взаимной аутентификации «внутренний» вызов становится внешним в момент компрометации одного пода.
- Топология не меняется. Сервис резолвит DNS при старте и кеширует IP навсегда. Поды переезжают, старые адреса достаются чужому приложению — в логах
connection refusedвперемешку с404от совершенно другого сервиса. - Администратор один. Команда А поднимает свой таймаут с 1 с до 10 с, «чтобы убрать ошибки»; пул соединений команды Б выедается ожиданием, и падает уже она.
- Транспорт бесплатен. Перевод внутреннего API с protobuf на JSON «для удобства отладки»: +30% CPU и +8 мс на вызов, что на четырёх уровнях глубины даёт +32 мс.
- Сеть однородна. Между подами всё работает, а через сервис-меш ломается gRPC-стриминг: прокси не поддерживает трейлеры HTTP/2 нужной версии.
7. Карта невозможностей
Задача двух генералов. Два генерала на разных холмах, гонца могут перехватить. Договориться о времени атаки со стопроцентной уверенностью нельзя: каждое подтверждение само требует подтверждения. Прямое следствие — невозможность exactly-once доставки по ненадёжному каналу (раздел 9).
FLP. Fischer, Lynch, Paterson, 1985: в полностью асинхронной системе не существует детерминированного алгоритма, гарантированно достигающего консенсуса при отказе даже одного процесса. Не «трудно», а «невозможно». Обход — не отменить теорему, а сменить модель: частичная синхронность (Dwork, Lynch, Stockmeyer) и рандомизация. Raft и Paxos гарантируют безопасность всегда, а живость — только в периоды нормальной сети. Ровно это вы видите в проде: etcd никогда не отдаст две разные истории лога, но может не отдавать ничего, пока идут выборы.
2026-07-12 09:41:03.114 I | raft: 8e9e05c52164694d became candidate at term 7
2026-07-12 09:41:03.115 I | raft: 8e9e05c52164694d received MsgVoteResp rejection from b2c4d7...
2026-07-12 09:41:06.902 W | etcdserver: failed to send out heartbeat on time (exceeded 100ms for 214.3ms)
2026-07-12 09:41:06.902 W | etcdserver: server is likely overloaded
2026-07-12 09:41:12.447 E | etcdserver: publish error: etcdserver: request timed out
Это безопасность в действии: кластер предпочёл не ответить, чем ответить неправдой. Приложение, хранящее в etcd лидерство, обязано быть готово к нескольким секундам «не знаю, кто лидер», и не считать это исключительной ситуацией. Отдельно отметьте строку server is likely overloaded: у etcd каждая запись — это fsync в WAL на всех членах кворума, и медленный диск ломает не пропускную способность, а именно живость выборов. Классический прод-инцидент выглядит как «Kubernetes API перестал отвечать» и лечится переносом etcd на отдельные NVMe.
CAP. Формулировка Брюера (2000) и доказательство Gilbert и Lynch (2002): при сетевом разделении система не может быть одновременно линеаризуемой и полностью доступной. Чего CAP не говорит: что нужно «выбрать два из трёх» (P не выбирают, разделения случаются сами), что вся система относится к одному классу (гарантии выбираются на уровне операции) и что вне разделений выбора нет. Последнее закрывает PACELC Абади: если Partition — то A или C, иначе (Else) — L (latency) или C. Разбор — «CAP и PACELC».
8. Кворум, split-brain и цена координации
Кворум — не «голосование ради демократии», а свойство пересечения множеств. Если запись подтверждена W узлами, чтение опрашивает R узлов и R + W > N, множества обязательно пересекутся хотя бы в одном узле, который видел последнюю запись. Отсюда конфигурации: N=3, W=2, R=2 — переживаем отказ одного узла; W=N, R=1 — быстрые чтения, недоступная запись при отказе любого узла; W=1, R=1 — быстро и без гарантий (это Cassandra с CL=ONE, ровно та настройка, при которой Jepsen находит потерянные записи).
Второе свойство кворума важнее первого: большинство может быть только одно. Поэтому при разделении 3/2 меньшая половина обязана перестать принимать записи. Системы, делающие failover по heartbeat без кворума, получают split-brain — нижний блок на схеме выше. Реальный пример — инцидент GitHub 21 октября 2018: 43 секунды потери связности между побережьями США, автоматика перевела MySQL-кластер на другой регион, но часть записей уже была принята прежним лидером. Итог — сутки деградации и ручной сверки. Причина не в 43 секундах, а в том, что решение о лидерстве принималось без данных, оставшихся по другую сторону разрыва.
| Механизм | Сообщений на операцию | RTT | Что даёт | Что ломается |
|---|---|---|---|---|
| Локальная запись | 0 | 0 | скорость | не переживает падение узла |
| Асинхронная репликация | O(N) фоном | 0 | долговечность «потом» | потеря хвоста при отказе лидера |
| Кворумная запись | O(N) | 1 | R+W>N, чтение видит запись |
недоступность при потере кворума |
| Консенсус (Raft, steady state) | 2(N−1) | 1 + fsync | линеаризуемый лог | простой на время выборов |
| 2PC | 4(N−1) | 2 | атомарность между ресурсами | блокировки при отказе координатора |
Строка про 2PC — ключевая для практики: между фазами prepare и commit участники держат блокировки, и если координатор умер, держат их до его возвращения. Это причина, по которой XA-коммит почти вымер в интернет-системах и заменён saga и outbox — см. «Распределённые транзакции» и saga как архитектурный шаблон.
Отдельный сюжет — время как способ сэкономить на координации. Spanner (Corbett et al., OSDI 2012) использует TrueTime: часы с известной границей неопределённости (атомные часы и GPS в каждом ДЦ, интервал — единицы миллисекунд). Транзакция после коммита ждёт, пока неопределённость не истечёт (commit-wait), и это даёт внешнюю согласованность без лишнего раунда сообщений. Мораль: Spanner не «победил CAP», он купил верхнюю границу на рассинхронизацию часов за деньги на железо — подробности в статье «Время в распределённой системе».
9. Почему exactly-once — миф, и что делают вместо
Exactly-once delivery невозможна. Exactly-once processing достижима — и это разные вещи.
Доказательство — задача двух генералов. Отправитель не может отличить «сообщение не дошло» от «дошло, но ответ потерялся». Не повторяет — возможна потеря (at-most-once); повторяет — возможен дубль (at-least-once). Третьего не даёт никакой протокол: любой уровень подтверждений требует подтверждения подтверждения. Развёрнутый разбор — «You Cannot Have Exactly-Once Delivery».
Что тогда означает «exactly-once» в документации Kafka и Flink? Строго ограниченную вещь: атомарность цикла read-process-write внутри границы одной системы. Kafka (KIP-98: идемпотентный продюсер плюс транзакции) гарантирует, что при чтении из топика, изменении состояния и записи в другой топик потребитель с isolation.level=read_committed увидит либо всё, либо ничего. Работает, пока источник, состояние и приёмник — внутри Kafka. Как только эффект уходит наружу (списание с карты, письмо, вызов внешнего API), гарантия обрывается ровно на границе. У Flink то же самое: exactly-once там про согласованность состояния через чекпойнты и two-phase commit sinks, а не про однократную доставку в мир.
Граница видна и в логах. Транзакционный продюсер Kafka защищается от «зомби» — старого экземпляра приложения, который очнулся после паузы и пытается дописать в ту же транзакцию:
11:07:41.220 INFO o.a.k.c.p.i.TransactionManager ProducerId set to 4021 with epoch 7
11:07:58.914 ERROR o.a.k.c.p.i.Sender [Producer clientId=billing-1, transactionalId=billing-tx-3]
Aborting producer batches due to fatal error: ProducerFencedException:
There is a newer producer with the same transactionalId which fences the current one
Это ровно fencing-токен из раздела 10, только внутри брокера: epoch монотонно растёт, старый экземпляр отсекается. Гарантия действует до send() в Kafka — и ни на шаг дальше: если тот же обработчик успел дёрнуть платёжный шлюз, транзакция Kafka откатит запись в топик, но не списание.
в ТОЙ ЖЕ транзакции, что и эффект"] B -->|"Во внешнем API"| D{"API принимает
ключ идемпотентности?"} B -->|"В реальном мире:
письмо, SMS, отгрузка"| E["Дубль неустраним.
Сужаем окно, добавляем подтверждение"] D -->|"Да, например Stripe"| F["Передаём стабильный ключ,
провайдер дедуплицирует сам"] D -->|"Нет"| G["Двухшаговый протокол:
сначала статус по бизнес-ключу,
потом действие"] C --> H["Effectively-once:
эффект применён один раз"] F --> H G --> I["Окно дубля сузилось,
но не исчезло"] E --> I H --> J["Дубли логируем как applied=false —
это норма, а не ошибка"] I --> J
Ключевая техника — дедупликация в одной транзакции с эффектом. Если запись в таблицу обработанных и сам эффект коммитятся раздельно, вы переместили окно дубля, а не убрали его.
CREATE TABLE processed_messages (
idempotency_key TEXT PRIMARY KEY, -- стабильный ключ от продюсера, НЕ uuid4() в момент ретрая
handler TEXT NOT NULL, -- один ключ может обрабатываться разными обработчиками
result_hash TEXT, -- чтобы заметить «тот же ключ, другой payload»
processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX ON processed_messages (processed_at); -- для чистки: ключи живут дольше окна ретраев, но не вечно
import hashlib, json, psycopg
from psycopg.errors import UniqueViolation
def handle_payment_captured(conn: psycopg.Connection, msg: dict) -> str:
"""Идемпотентная обработка события: 'applied' либо 'duplicate'.
Инвариант: запись ключа и бизнес-эффект попадают в ОДНУ транзакцию.
Откатится она — откатится и то, и другое, и повтор пройдёт честно.
"""
key = msg["idempotency_key"] # стабилен между ретраями продюсера
payload_hash = hashlib.sha256(json.dumps(msg["data"], sort_keys=True).encode()).hexdigest()
with conn.transaction(): # BEGIN ... COMMIT/ROLLBACK
try:
conn.execute(
"INSERT INTO processed_messages (idempotency_key, handler, result_hash)"
" VALUES (%s, %s, %s)",
(key, "payment_captured", payload_hash),
)
except UniqueViolation:
return "duplicate" # уже обрабатывали — это НОРМА, а не ошибка
conn.execute( # бизнес-эффект, атомарный с записью ключа выше
"UPDATE orders SET status = 'paid', paid_at = now()"
" WHERE id = %s AND status = 'pending'",
(msg["data"]["order_id"],),
)
return "applied"
Сложность: O(1) по времени на сообщение — одна вставка по первичному ключу, то есть O(log n) обращений к страницам B-дерева, на практике 3–4. По памяти O(K), где K — число ключей в окне хранения. Поэтому политика очистки обязательна: без неё таблица растёт линейно по трафику и однажды перестаёт помещаться в кэш, после чего дедупликация становится самой дорогой операцией в системе.
Три вида эффектов и что с ними делать:
- Эффект в моей БД — дедуп в транзакции, как выше. Полная победа.
- Эффект во внешнем API — ключ идемпотентности провайдера (
Idempotency-Keyу Stripe). Ключ обязан быть детерминированной функцией от бизнес-события, а неuuid4()в момент вызова: иначе каждый ретрай приносит новый ключ и дедупликации не происходит. - Эффект в реальном мире — письмо, SMS, отгрузка. Дубль неустраним в принципе. Остаётся сужать окно (записать «отправлено» до отправки, приняв риск потери вместо риска дубля, если потеря дешевле) и делать эффект видимо повторяемым: «мы уже отправляли вам этот код» вместо второго письма.
В резюме архитектурного решения пишите не «мы обеспечиваем exactly-once», а «доставка at-least-once, обработка effectively-once за счёт ключей идемпотентности; для нотификаций допускается дубль с окном до N секунд». Такая формулировка проверяема. Детали — «Гарантии доставки и идемпотентность».
10. Пять инструментов, закрывающих большую часть практики
1. Таймаут + повтор с джиттером + бюджет. Экспоненциальная задержка без случайности синхронизирует клиентов: все повторяют одновременно и добивают сервис. AWS в Exponential Backoff And Jitter показал, что лучший практический вариант — full jitter.
import random, time
from typing import Callable, TypeVar
T = TypeVar("T")
class Retryable(Exception):
"""Таймаут, 503, обрыв соединения — повтор имеет смысл."""
class Unknown(Retryable):
"""Исход неизвестен: запрос мог примениться. Повтор безопасен ТОЛЬКО для идемпотентных операций."""
def call_with_retry(fn: Callable[[], T], *, idempotent: bool, deadline: float,
attempts: int = 3, base: float = 0.1, cap: float = 2.0) -> T:
"""Повтор с full jitter и общим дедлайном (time.monotonic()).
Бюджет времени важнее числа попыток: без дедлайна три попытки по 30 с
превращают быстрый отказ в полутораминутное зависание клиента.
"""
last: Exception | None = None
for attempt in range(attempts):
if time.monotonic() >= deadline:
raise TimeoutError("бюджет вызова исчерпан") from last
try:
return fn()
except Unknown as e:
if not idempotent: # неизвестный исход неидемпотентной операции — не повторяем
raise
last = e
except Retryable as e:
last = e
backoff = random.uniform(0, min(cap, base * (2**attempt))) # full jitter
if backoff >= deadline - time.monotonic():
raise TimeoutError("бюджет вызова исчерпан") from last
time.sleep(backoff)
raise last # type: ignore[misc]
2. Ключ идемпотентности — разобран в разделе 9. 3. Кворум вместо «спросим лидера» — раздел 8 и статья «Репликация».
4. Fencing-токен вместо доверия к блокировке. Блокировка с TTL не гарантирует, что владелец жив: он мог зависнуть. Гарантию даёт монотонный номер, который проверяет сам ресурс.
// Владелец получает лизу с монотонным номером: в etcd это ревизия, в ZooKeeper — zxid, в Raft — term.
type Lease struct {
Owner string
Token int64 // строго возрастает при каждой смене владельца
}
// Ресурс — единственное место, где решается, чья запись легитимна.
func (s *ShardStore) Write(ctx context.Context, l Lease, data []byte) error {
s.mu.Lock()
defer s.mu.Unlock()
if l.Token < s.lastToken { // владелец сменился, пока вызывающий был в паузе GC
return fmt.Errorf("устаревший токен %d, актуальный %d", l.Token, s.lastToken)
}
s.lastToken = l.Token
return s.apply(ctx, data)
}
Без этой проверки сценарий из раздела 4 (worker-7 очнулся после паузы) заканчивается порчей данных, и никакой TTL блокировки его не предотвратит.
5. Наблюдаемость с корреляцией. В распределённой системе нет стек-трейса; его роль играет trace_id, протянутый через все вызовы, включая асинхронные — то есть записанный в заголовки сообщения, а не только в HTTP. Без этого разбор инцидента превращается в сопоставление меток времени с машин, часы которых разъезжаются. Практика — «Наблюдаемость распределённых систем» и observability со стороны эксплуатации.
11. Метастабильные отказы: когда система не восстанавливается сама
Отдельный класс аварий: система под нагрузкой переходит в состояние, где её собственные механизмы восстановления поддерживают отказ. Формализовано в «Metastable Failures in Distributed Systems» (HotOS, 2021). Канонический пример — шторм повторов: сервис замедляется, клиенты ретраят, нагрузка растёт втрое, сервис замедляется сильнее, ретраев становится больше. Даже когда исходная причина устранена, система не выходит из состояния — её держит собственный трафик повторов, и требуется принудительный сброс нагрузки.
Арифметика опасности: если каждый из d уровней стека делает k попыток, один пользовательский запрос порождает до k^d обращений к нижнему сервису. При k=3 и четырёх уровнях — 81 запрос. Это O(k^d), экспонента по глубине, и она превращает «безобидные» ретраи на каждом уровне в механизм самоуничтожения. Что помогает: ретраи только на одном уровне (обычно самом внешнем, где известен бизнес-смысл); retry budget — token bucket, разрешающий повторы не более чем для ~10% трафика (gRPC, подход Google SRE); circuit breaker; сброс нагрузки по длине очереди, а не по CPU (очередь растёт раньше); передача оставшегося бюджета времени вниз по цепочке (deadline propagation). Источники: AWS Builders’ Library и глава про перегрузку в Google SRE Book.
Как метастабильность выглядит в мониторинге: RPS входящий — плоский, RPS исходящий на нижний сервис — вырос втрое, доля успешных ответов — 4%, средняя latency упёрлась в таймаут, а после перезапуска нижнего сервиса всё повторяется через 30 секунд. Последнее и есть диагностический признак: отказ переживает устранение своей причины. Единственный работающий приём — снять нагрузку (закрыть трафик на балансировщике, включить сброс очередей) и вернуть её ступенчато.
12. Типичные ошибки, которые видно на код-ревью
- Повтор неидемпотентной операции.
POST /transferс ретраем и без ключа — прямой путь к двойному списанию. - Сортировка событий по
created_atс разных машин. Часы разъезжаются, NTP делает скачки назад, «более поздний» ивент оказывается раньше. Нужны логические или гибридные часы — «Время в распределённой системе». - Чтение с реплики сразу после записи в лидер. Пользователь сохранил профиль и увидел старый: нарушен read-your-writes. Лечится закреплением сессии за лидером на окно репликации или чтением по версии — «Модели согласованности».
- Клиентский таймаут меньше серверного. Клиент отваливается, сервер продолжает работу, пользователь ретраит — умножение нагрузки без единого полезного ответа.
- Блокировка с TTL без fencing-токена — раздел 10.
- Тесты, проходящие только при строгом порядке сообщений. В проде такой порядок не гарантирован; инструменты — «Тестирование распределённых систем».
- Логи без
trace_idв асинхронных путях. После брокера след теряется, и половина инцидента остаётся невидимой — «Очереди и потоки». - Кворум из чётного числа узлов. N=4 переживает отказ одного узла ровно как N=3, но стоит дороже и даёт больше сценариев разделения пополам — «Консенсус».
- Ключ шардирования, совпадающий с самым горячим полем. Один клиент-гигант кладёт партицию целиком — «Партиционирование».
- Вера в «у нас маленькая нагрузка, разделений не будет». Перегруженный узел, длинная пауза GC и переполненная очередь выглядят для остальных точно так же, как обрыв кабеля — «Модели отказов».
13. Мини-итог
- Распределённой систему делает частичный отказ, а не количество серверов.
- У сетевого вызова три исхода, и третий — «неизвестно» — главный источник багов в деньгах.
- Тайм-аут — решение наблюдателя, а не факт о мире; точных детекторов отказов не бывает, а серые отказы вообще не срабатывают на health-check.
- FLP запрещает надёжный консенсус в асинхронной модели: реальные алгоритмы жертвуют живостью, но никогда — безопасностью.
- CAP не про «выбери два из трёх», а про поведение во время разделения; PACELC добавляет выбор в нормальном режиме.
- Кворум даёт пересечение множеств и единственность большинства — отсюда защита от split-brain.
- Exactly-once доставки не существует; существует at-least-once доставка плюс идемпотентная обработка, атомарная с эффектом.
- Ретраи без бюджета и джиттера превращают замедление в аварию, которая переживает устранение своей причины.
Источники
Первоисточники, к которым трек возвращается постоянно:
- Lamport. Time, Clocks, and the Ordering of Events (CACM, 1978) и The Byzantine Generals Problem (TOPLAS, 1982).
- Fischer, Lynch, Paterson. Impossibility of Distributed Consensus with One Faulty Process (JACM, 1985).
- Lamport. Paxos Made Simple (2001); Ongaro, Ousterhout. Raft (USENIX ATC, 2014).
- DeCandia et al. Dynamo (SOSP, 2007); Corbett et al. Spanner (OSDI, 2012); Burrows. Chubby (OSDI, 2006).
- Gilbert, Lynch. Brewer’s Conjecture (доказательство CAP) (2002); Abadi. PACELC (IEEE Computer, 2012).
- Chandra, Toueg. Unreliable Failure Detectors (JACM, 1996).
Книги и практика: Martin Kleppmann, Designing Data-Intensive Applications (O’Reilly, 2017) — обязательная книга трека, главы 5–9 покрывают половину материала; Jepsen analyses — читайте отчёт по своей системе до того, как поверите её документации; AWS Builders’ Library и Google SRE Book. Постмортемы как учебник: GitHub, октябрь 2018, AWS S3, февраль 2017, Roblox, октябрь 2021 — 73 часа простоя из-за взаимодействия Consul, Raft и бага в BoltDB.
Что дальше
Модели отказов: сбои узлов, сеть, разделение, византийские отказы — разберём, какие именно предположения об отказах делает алгоритм, чем краш отличается от omission и timing, почему византийская модель почти никогда не нужна внутри вашего ДЦ и всегда нужна на границе доверия, и как выбор модели напрямую определяет число узлов, которое придётся содержать.