Пока программа работает на одном компьютере, всё предсказуемо: одна и та же операция даёт один и тот же результат, а если что-то ломается — обычно ломается всё сразу и целиком. Но как только процессы начинают общаться по сети, эта определённость исчезает. Появляется частичный отказ: одни узлы работают, другие нет, а третьи вроде работают, но точно никто не знает. Это и есть главная особенность распределённых систем.
Почти вся эта статья — про то, чему в такой системе нельзя доверять: сети, часам и даже тому, что твой собственный процесс не поставили на паузу. Отсюда прямая дорога к алгоритмам консенсуса, но сначала надо понять, от чего именно они защищают.
Частичный отказ: тайм-аут — единственный инструмент
Узел отправил запрос и не получил ответа. Что случилось? Вариантов много: потерялся сам запрос; узел упал; узел жив, но отвечает медленно; узел всё обработал, а потерялся уже ответ. И вот главная неприятность — различить эти случаи невозможно. У отправителя есть только один факт: «ответа нет». Единственный способ вообще принять хоть какое-то решение — тайм-аут: подождать сколько-то времени и, если ответа так и нет, объявить узел неработающим.
Но какой тайм-аут выбрать? Короткий — быстро замечаешь сбои, но рискуешь объявить мёртвым узел, который просто притормозил под нагрузкой. Тогда его работу передадут другим — и добавят нагрузки уже перегруженной системе, а это прямой путь к каскадному отказу. Длинный тайм-аут — наоборот, долго ждёшь, прежде чем среагировать. «Правильного» значения тут нет, потому что в обычной сети задержки ничем не ограничены: пакет может застрять в очереди перегруженного коммутатора, у занятого процессора, у гипервизора, который приостановил виртуальную машину. (Сеть с гарантированной задержкой в принципе возможна — так работает телефония с выделенным каналом на звонок, — но интернет и сети дата-центров оптимизированы под пиковый трафик через очереди, и платят за это предсказуемостью.) Поэтому тайм-ауты подбирают опытным путём, а ещё лучше — измеряют разброс задержек и подстраивают динамически (так делают Cassandra и Akka).
Три вещи, которым нельзя доверять
Сеть. Пакеты теряются и задерживаются на непредсказуемое время; соединение может работать в одну сторону и молчать в обратную; сетевые сбои случаются даже в аккуратно управляемых дата-центрах чаще, чем кажется. Вывод простой: любой обмен по сети может сорваться, и обработку этих сбоев надо специально проектировать и тестировать (в духе Chaos Monkey — намеренно рвать сеть прямо в проде и смотреть, выживет ли система).
Часы. У каждой машины свои внутренние часы, и они потихоньку расходятся друг с другом (это называют clock drift). Важно: часов бывает два вида, и путать их опасно.
- Часы «настенного времени» (
System.currentTimeMillis()) синхронизируются по сети (по протоколу NTP) и могут прыгнуть назад при коррекции, спотыкаются о «секунды координации», зависят от неизвестной точности сервера времени. Для измерения интервалов они не годятся. - Монотонные часы (
System.nanoTime()) гарантированно идут только вперёд, но их абсолютное значение само по себе бессмысленно (это просто «сколько-то от какого-то момента»). Именно и только ими меряют длительности.
Отсюда главная ловушка — упорядочивать события по меткам настенного времени. Ровно так работает разрешение конфликтов «выигрывает последний» (LWW): у двух конкурирующих записей сравнивают время, побеждает та, что «позже». Но если часы двух узлов разошлись на 100 мс, то «более поздняя» по метке запись могла на самом деле произойти раньше — и та запись, которую клиент реально сделал последней, молча пропадёт. Метки времени не гарантируют причинно-следственный порядок; для него нужны логические часы и версии. (Google Spanner ставит GPS-приёмники и атомные часы в каждый дата-центр как раз для того, чтобы сжать погрешность до пары миллисекунд.)
Паузы процессов. Твой поток может быть остановлен в любой момент на непредсказуемое время. Причин масса: всеобъемлющая пауза сборщика мусора (stop-the-world, иногда на минуты), приостановка виртуальной машины при переезде на другой хост, «украденное» время процессора, подкачка страницы памяти с диска, даже Ctrl+Z. И самое коварное — узел этого не замечает: для него между двумя соседними строками кода прошло «мгновение», а на самом деле — минута, и его уже давно все считают мёртвым.
Истину определяет большинство, а не сам узел
Из паузы вырастает один из самых коварных багов. Представь узел, который держит распределённую блокировку (лок) или считает себя ведущим. Он проверяет: «лок ещё мой?» — да, — и идёт писать. Но между проверкой и записью его заморозила пауза сборщика мусора. За эту минуту лок протух, его перехватил другой узел, а очнувшийся первый — всё ещё уверен, что владеет им — и пишет, портя данные. Два владельца одновременно.
Мораль: узел не может доверять собственному мнению о своём статусе. В распределённой системе истину определяет кворум — решение большинства узлов (обычно больше половины). Если большинство объявило узел мёртвым, он считается мёртвым, даже если на самом деле прекрасно работает. Решение по большинству безопасно, потому что два разных большинства не могут не пересечься — а значит, не будет двух противоречащих решений. Именно так узлы выбирают ведущего и не допускают split brain (когда ведущих сразу два).
Но кворум решает, кто должен писать, — и не мешает опоздавшему всё испортить. Практическая защита — ограждающий маркер (fencing token). Сервис блокировок при каждой выдаче лока возвращает постоянно растущий номер; клиент прикладывает этот номер к каждой своей записи; а хранилище отклоняет запись, если её номер меньше уже виденного. Узел, очнувшийся после паузы, приходит со старым номером — и его запись отвергают. Ключевой момент: проверять маркер должен сам ресурс (хранилище), а не клиент — ведь клиент, который считает себя «в порядке», как раз и есть источник проблемы. Подробный разбор с Redis и кодом — в задаче про лок без fencing.
Византийские сбои и модели системы
До сих пор мы считали узлы «честными»: они могут молчать, тормозить, отдавать устаревшее — но если уж отвечают, то не врут. Если же узел способен врать — слать произвольные или поддельные сообщения (из-за порчи памяти, бага или злого умысла), — это называют византийским сбоем, а согласование в такой среде — задачей о византийских генералах. Защита от неё дорогая и нужна там, где нет доверия: авиакосмос (радиация портит память), блокчейны (участники не доверяют друг другу). Внутри собственного дата-центра византийских сбоев обычно нет, и защищаться от них не окупается — но проверять данные, приходящие от внешних клиентов (валидация, защита от инъекций), нужно всегда.
Чтобы рассуждать о корректности, алгоритмы описывают через модель системы — набор допущений. По времени: синхронная (задержки ограничены — нереалистично), частично синхронная (обычно ведёт себя хорошо, иногда нет — реалистичный дефолт), асинхронная (вообще никаких допущений о времени). По отказам: отказ-остановка, отказ-восстановление (узел падает и снова поднимается, а надёжное хранилище переживает сбой), византийская. И два вида свойств, которые алгоритм должен обеспечивать: безопасность (safety — «ничего плохого не случится»; нарушение имеет конкретный момент и необратимо) и живучесть (liveness — «со временем случится что-то хорошее»; пример — та самая конечная согласованность). Хорошие алгоритмы держат безопасность всегда, а живучесть — при разумных допущениях.
CAP без мифов — и что добавляет PACELC
Популярная формулировка CAP-теоремы «выбери два из трёх: согласованность, доступность, устойчивость к разделению» — неточна и приводит к неверным выводам. Разделение сети не выбирают — оно случается: сеть ненадёжна по природе. Точная формулировка такая: в момент разделения сети система вынуждена выбрать между согласованностью (все узлы отвечают одинаково) и доступностью (каждый запрос получает ответ). В нормальном режиме, без разделения, доступны обе.
Из-за этой оговорки CAP отвечает только на вопрос «что будет при сбое сети» и молчит о повседневной работе. Пробел закрывает PACELC: при разделении (P) — выбор между A и C, иначе (Else) — выбор между задержкой (Latency) и согласованностью (Consistency). Даже в идеально здоровой сети синхронная репликация покупает согласованность ценой задержки на каждую запись, асинхронная — наоборот. Примеры для ориентира: PostgreSQL с синхронной репликой — PC/EC, с асинхронной — PC/EL; Cassandra с настройкой кворумов сдвигается к PA/EL.
Где это применяется
Как только у вас больше одного сервиса и они ходят друг к другу по сети — вы уже в распределённой системе, даже если не планировали. Любой сетевой вызов может потеряться, зависнуть или упереться в тайм-аут; любой узел может замолчать посреди операции. Практическая рамка: не тяните распределённость раньше времени (три условия, когда она реально нужна), но если она уже есть — проектируйте под частичный отказ: тайм-ауты и повторы (обязательно идемпотентные!), никакого упорядочивания по настенным часам, решения через кворум, fencing на общих ресурсах.
Где спотыкаются начинающие:
- Считают, что «нет ответа» значит «узел упал». Это неотличимо от потери ответа или медленного узла. Отсюда двойные списания и потерянные данные — без единой ошибки в логах.
- Упорядочивают события по
currentTimeMillis(). Часы узлов разошлись — и LWW молча теряет запись, которую клиент считал сохранённой. - Проверяют лок и сразу пишут. Пауза сборщика мусора между проверкой и записью даёт двух владельцев. Нужен fencing token, который проверяет сам ресурс.
- Узел верит собственному «я ещё ведущий». А за время его паузы кворум уже выбрал другого. Истину определяет большинство.
- Тянут микросервисы «ради масштаба». И получают все проблемы распределённых систем там, где хватило бы одного узла.
Что почитать дальше: модели репликации — кворумы w+r>n и конфликты записи; строительные блоки — где в системе живут очереди, репликация и координация; распределённый лок без fencing — механизм ограждающего маркера с кодом; когда нужны распределённые паттерны — три условия и три альтернативы.