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

Распределённые системы: карта трека и почему всё ломается иначе

Распределённые системы: карта трека и почему всё ломается иначе

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

Трек про фундамент: отказы, время, согласованность, координация. Про стили построения систем — микросервисы, event-driven, saga как архитектурный шаблон — есть отдельный трек «Архитектурные паттерны». Здесь мы разбираем не «как нарезать сервисы», а почему любой разрез по сети приносит класс отказов, которого не было в одном процессе.

1. Определение, которое действительно работает

Самая точная формулировка принадлежит Лесли Лампорту (премия Тьюринга 2013 года, автор доброй половины фундамента этой области) — это его письмо 1987 года:

Распределённая система — это такая система, в которой отказ компьютера, о существовании которого вы даже не подозревали, делает ваш собственный компьютер непригодным к работе.

Шутка описывает следствие. Причина — три отсутствующие вещи:

  1. Нет общей памяти. Любое «общее знание» физически представлено копиями, а копии расходятся. То, что вы прочитали, — снимок прошлого, а не текущее состояние.
  2. Нет общих часов. У каждого узла свои кварцевые часы со своим дрейфом. Значит, нет объективного «одновременно» и нет объективного порядка событий между узлами.
  3. Нет общей судьбы. Компоненты падают по отдельности. Это и есть частичный отказ (partial failure) — главный водораздел.

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

  1. Отказ и медленность неразличимы. Узел, который не ответил, мог упасть, мог зависнуть в паузе сборщика мусора, мог ответить — но ответ потерялся. Различить эти случаи принципиально нельзя без верхней границы на задержку сети.
Свойство Один процесс Распределённая система
Вызов функции вернул или процесс упал вернул, упал или неизвестно
Время один монотонный счётчик N часов с взаимным дрейфом
Состояние одна копия N копий, каждая устарела по-своему
Отказ всё или ничего любое подмножество компонентов
Отладка стек-трейс корреляция логов N сервисов
Воспроизводимость теста детерминированная зависит от расписания и сети

Дальше весь трек — способы жить с пунктами 1–4, не разрушив бизнес-инварианты.

2. Карта трека

Маршруты чтения. Разбираетесь с дублями в очереди — 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 не спасает

В одном процессе у вызова два исхода. По сети — три, и третий не имеет обработчика в вашем языке.

Как это выглядит в логах — почти всегда одинаково:

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» Мартина Клеппманна.

Более умный вариант — 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.

  1. Сеть надёжна. Бэйлис и Кингсбери в «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 без идемпотентного продюсера порядок в партиции нарушен.
  2. Задержка равна нулю. GET /order/{id} последовательно ходит в profile, pricing, inventory, loyalty; у каждого p99 = 300 мс — p99 страницы больше секунды, хотя «все сервисы в SLA».
  3. Пропускная способность бесконечна. Ответ вырос с 40 КБ до 400 КБ после добавления вложенной коллекции. Внутри ДЦ незаметно; на межрегиональном канале упирается в полосу, и растёт p99 у всех соседей по каналу.
  4. Сеть безопасна. Отказ выглядит не как ошибка, а как успешный запрос из неожиданного места: без взаимной аутентификации «внутренний» вызов становится внешним в момент компрометации одного пода.
  5. Топология не меняется. Сервис резолвит DNS при старте и кеширует IP навсегда. Поды переезжают, старые адреса достаются чужому приложению — в логах connection refused вперемешку с 404 от совершенно другого сервиса.
  6. Администратор один. Команда А поднимает свой таймаут с 1 с до 10 с, «чтобы убрать ошибки»; пул соединений команды Б выедается ожиданием, и падает уже она.
  7. Транспорт бесплатен. Перевод внутреннего API с protobuf на JSON «для удобства отладки»: +30% CPU и +8 мс на вызов, что на четырёх уровнях глубины даёт +32 мс.
  8. Сеть однородна. Между подами всё работает, а через сервис-меш ломается 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 и цена координации

Кластер из пяти узлов при сетевом разделении: половина с кворумом продолжает запись, половина без кворума уходит в read-only

Кворум — не «голосование ради демократии», а свойство пересечения множеств. Если запись подтверждена 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 откатит запись в топик, но не списание.

Ключевая техника — дедупликация в одной транзакции с эффектом. Если запись в таблицу обработанных и сам эффект коммитятся раздельно, вы переместили окно дубля, а не убрали его.

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

Три вида эффектов и что с ними делать:

  1. Эффект в моей БД — дедуп в транзакции, как выше. Полная победа.
  2. Эффект во внешнем API — ключ идемпотентности провайдера (Idempotency-Key у Stripe). Ключ обязан быть детерминированной функцией от бизнес-события, а не uuid4() в момент вызова: иначе каждый ретрай приносит новый ключ и дедупликации не происходит.
  3. Эффект в реальном мире — письмо, 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. Типичные ошибки, которые видно на код-ревью

  1. Повтор неидемпотентной операции. POST /transfer с ретраем и без ключа — прямой путь к двойному списанию.
  2. Сортировка событий по created_at с разных машин. Часы разъезжаются, NTP делает скачки назад, «более поздний» ивент оказывается раньше. Нужны логические или гибридные часы — «Время в распределённой системе».
  3. Чтение с реплики сразу после записи в лидер. Пользователь сохранил профиль и увидел старый: нарушен read-your-writes. Лечится закреплением сессии за лидером на окно репликации или чтением по версии — «Модели согласованности».
  4. Клиентский таймаут меньше серверного. Клиент отваливается, сервер продолжает работу, пользователь ретраит — умножение нагрузки без единого полезного ответа.
  5. Блокировка с TTL без fencing-токена — раздел 10.
  6. Тесты, проходящие только при строгом порядке сообщений. В проде такой порядок не гарантирован; инструменты — «Тестирование распределённых систем».
  7. Логи без trace_id в асинхронных путях. После брокера след теряется, и половина инцидента остаётся невидимой — «Очереди и потоки».
  8. Кворум из чётного числа узлов. N=4 переживает отказ одного узла ровно как N=3, но стоит дороже и даёт больше сценариев разделения пополам — «Консенсус».
  9. Ключ шардирования, совпадающий с самым горячим полем. Один клиент-гигант кладёт партицию целиком — «Партиционирование».
  10. Вера в «у нас маленькая нагрузка, разделений не будет». Перегруженный узел, длинная пауза GC и переполненная очередь выглядят для остальных точно так же, как обрыв кабеля — «Модели отказов».

13. Мини-итог

  • Распределённой систему делает частичный отказ, а не количество серверов.
  • У сетевого вызова три исхода, и третий — «неизвестно» — главный источник багов в деньгах.
  • Тайм-аут — решение наблюдателя, а не факт о мире; точных детекторов отказов не бывает, а серые отказы вообще не срабатывают на health-check.
  • FLP запрещает надёжный консенсус в асинхронной модели: реальные алгоритмы жертвуют живостью, но никогда — безопасностью.
  • CAP не про «выбери два из трёх», а про поведение во время разделения; PACELC добавляет выбор в нормальном режиме.
  • Кворум даёт пересечение множеств и единственность большинства — отсюда защита от split-brain.
  • Exactly-once доставки не существует; существует at-least-once доставка плюс идемпотентная обработка, атомарная с эффектом.
  • Ретраи без бюджета и джиттера превращают замедление в аварию, которая переживает устранение своей причины.

Источники

Первоисточники, к которым трек возвращается постоянно:

Книги и практика: 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, почему византийская модель почти никогда не нужна внутри вашего ДЦ и всегда нужна на границе доверия, и как выбор модели напрямую определяет число узлов, которое придётся содержать.

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

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

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

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