Теперь можно пройти мок, чтобы познакомиться с нами

System design · Распределённые системы

Распределённые системы: как устроены и как их масштабировать

Что такое распределённая система, почему в ней нет общей памяти и общих часов, как узлы общаются, выбирают лидера, реплицируют данные и переживают сбои. CAP, Raft, шардирование, CQRS и circuit breaker простым языком.

Коротко

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

  • порядок событий без общих часов задают логические часы — часы Лэмпорта и векторные часы;
  • единый источник истины дают выбор лидера (Raft) и репликация данных;
  • по теореме CAP при разделении сети система выбирает между согласованностью и доступностью;
  • нагрузку распределяют шардирование, репликация, балансировщики, CQRS и асинхронные очереди;
  • от каскадных отказов защищают таймауты, ретраи с экспоненциальной задержкой, circuit breaker и rate limiting.

Что такое распределённая система?

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

  • Общая память: нет
  • Общие часы: нет
  • Связь: сообщения по сети — gRPC, HTTP, Kafka

Слыша «распределённые системы», обычно вспоминают кластеры, микросервисы или Kubernetes. Но суть проще: несколько машин идут к одной общей цели, и ни одна из них не «знает всё».

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

Распределённую систему определяют три свойства: нет общей памяти, нет общих часов, связь только через сообщения.

Почему в распределённой системе нет общей памяти?

Каждый узел работает независимо и не может напрямую прочитать переменные другой машины. Всё взаимодействие идёт через сообщения по сети, а сообщения теряются и не доходят. Чтобы система оставалась корректной и в этом случае, нужны идемпотентность, логика повторных попыток (retry) и алгоритмы консенсуса.

Почему в распределённой системе нет общих часов?

Единого времени для всех машин не существует: часы расходятся (clock drift), а сетевые задержки делают время доставки непредсказуемым. Из-за этого трудно установить, в каком точно порядке произошли события.

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

Как узлы распределённой системы обмениваются сообщениями?

Узлы общаются по протоколам вроде gRPC и HTTP или через события Kafka. Сообщение может не дойти, прийти с опозданием, продублироваться или прийти не по порядку.

Распределённая система остаётся корректной не потому, что связь идеальна, а потому, что аккуратно обрабатывает сбои.

Зачем нужны распределённые системы?

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

Традиционно программы работали на одной машине. С ростом популярности интернета стало ясно, что в масштабе одной машины может не хватить, — так появились вертикальное и горизонтальное масштабирование.

Что такое вертикальное масштабирование?

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

Что такое горизонтальное масштабирование?

Горизонтальное масштабирование — это добавление новых серверов вместо усиления одного. Вместо одного мощного сервера задаются вопросом: «А что если работу разделят много небольших серверов?» Если эти серверы координируются и действуют как единое целое, получается распределённая система — в этом её ключевая идея.

Горизонтальное и вертикальное масштабирование: слева три одинаковых сервера, которые просто добавляют, справа один сервер, который наращивают всё больше

Где используются распределённые системы?

Распределённые системы стоят за большинством крупных интернет-сервисов: поиском Google, стримингом Netflix, распределёнными базами данных вроде Cassandra, брокером событий Kafka, платёжной платформой Stripe. Везде работа разнесена по множеству узлов и дата-центров, а пользователь видит один сервис.

  • Google Search — огромная сеть краулеров, индексаторов и узлов ранжирования, которые работают во множестве дата-центров.
  • Netflix — сервисы уровня региона для стриминга, рекомендаций, аутентификации, транскодирования видео и доставки контента.
  • Cassandra — распределённое хранилище, которое реплицирует данные между узлами ради доступности и масштаба.
  • Kafka — распределённый отказоустойчивый журнал событий для асинхронного взаимодействия.
  • Stripe — событийно-ориентированная архитектура на надёжной доставке сообщений и идемпотентных операциях.

Чем распределённая система отличается от децентрализованной и параллельной?

Распределённая система состоит из слабо связанных узлов, которые общаются сообщениями и часто подчиняются лидеру или control plane. Децентрализованная система обходится без единого центра управления и договаривается через консенсус равноправных узлов. Параллельная система выполняет вычисления одновременно на одной машине или тесно связанном кластере с общей памятью.

Три термина звучат похоже, но описывают разные устройства систем.

Как устроена распределённая система?

  • Связь: слабая, обмен сообщениями
  • Управление: часто лидер или control plane
  • Примеры: распределённые SQL-базы, Apache Kafka

Группа независимых узлов работает вместе и представляет себя как единую логическую систему. Часто в ней есть «лидер или control plane», который отвечает за согласованность, репликацию и взаимодействие узлов.

Примеры: распределённые SQL-базы вроде CockroachDB (глобально распределённые базы данных) и Apache Kafka (распределённый журнал).

Как устроена децентрализованная система?

  • Связь: peer-to-peer
  • Управление: нет центрального органа
  • Примеры: Bitcoin, Ethereum, BitTorrent

Децентрализованная система убирает единую точку управления или сводит её роль к минимуму. Каждый узел действует автономнее, а координация идёт через peer-to-peer консенсус, а не через центральную оркестрацию.

Примеры: Bitcoin и Ethereum (блокчейн-сети), BitTorrent (обмен файлами между равноправными узлами).

Как устроена параллельная система?

  • Связь: тесная, общая память
  • Управление: одни общие часы
  • Примеры: GPU, многопоточные программы, HPC

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

Примеры: вычисления на GPU, многопоточные программы, кластеры высокопроизводительных вычислений (HPC).

Тип системыСвязанностьКак координируютсяНа что ориентирована
РаспределённаяСлабая, обмен сообщениямиЛидер или control planeНадёжность
ДецентрализованнаяРавноправные узлы (peer-to-peer)Консенсус без центрального органа, без доверия друг к другуОтсутствие единого центра
ПараллельнаяТесная, общая памятьОбщие часы и памятьПроизводительность

Если многим машинам нужно работать вместе, им приходится эффективно общаться и делиться состоянием. Отсюда следующий вопрос: как машинам общаться в масштабе?

Как узлы распределённой системы общаются по сети?

Узлы распределённой системы общаются только по сети, а сеть по умолчанию ненадёжна. Надёжную и упорядоченную доставку даёт TCP, шифрование и проверку подлинности собеседника — TLS, а находить нужный узел по имени помогает DNS.

Какие проблемы создаёт ненадёжная сеть?

По умолчанию сеть ненадёжна. Сообщения могут:

  • потеряться — так и не дойти до адресата;
  • продублироваться — прийти больше одного раза;
  • повредиться — биты изменятся при передаче;
  • прийти не по порядку — поздние сообщения опередят ранние;
  • задержаться — с непредсказуемой задержкой.

Проектировать распределённую систему — значит ожидать ХУДШЕГО и всё равно добиваться, чтобы система работала. Как гласит закон Мёрфи: «Если что-то может пойти не так, оно пойдёт не так».

Как TCP обеспечивает надёжную доставку?

TCP (Transmission Control Protocol) — транспортный протокол, на который распределённые системы сильно опираются. Он даёт:

  • надёжную доставку: данные либо доставляются, либо отправитель узнаёт о сбое;
  • порядок: пакеты приходят в том порядке, в котором были отправлены;
  • проверку ошибок: повреждённые пакеты обнаруживаются и отправляются повторно.
Тройное рукопожатие TCP: клиент отправляет серверу SYN, сервер отвечает SYN-ACK, клиент подтверждает ACK

Соединение TCP устанавливает тройным рукопожатием (3-way handshake):

  1. SYN: клиент говорит: «Хочу подключиться».
  2. SYN-ACK: сервер отвечает: «Принято, готов продолжать».
  3. ACK: клиент подтверждает: «Отлично, начинаем».

Что такое TLS и зачем он нужен?

TLS (Transport Layer Security) — протокол, который защищает соединение поверх TCP. Если TCP — это почтальон, то TLS — конверт с замком и подписью. TLS гарантирует:

  • шифрование: никто посторонний не прочитает данные;
  • аутентификацию: вы знаете, с кем говорите;
  • целостность: сообщение нельзя незаметно изменить.
Рукопожатие TLS: клиент отправляет ClientHello, сервер отвечает ServerHello, сертификатом и ServerHelloDone, затем клиент отправляет ClientKeyExchange, ChangeCipherSpec и Finish, а сервер завершает своими ChangeCipherSpec и Finish

Во время рукопожатия TLS обе стороны:

  • выбирают наборы шифров (cipher suites);
  • договариваются о версии TLS;
  • проверяют TLS-сертификат сервера;
  • генерируют сессионные ключи для зашифрованного обмена.

Надёжность теряет смысл, если данные можно перехватить или подменить. В современных системах каждый RPC и каждый сетевой вызов должен идти через TLS.

Как DNS помогает узлам находить друг друга?

DNS (Domain Name System) превращает понятные человеку имена в IP-адреса и служит простейшим механизмом обнаружения сервисов (service discovery). Когда серверы разбросаны по всему миру, именно так они находят друг друга.

Разрешение имени через DNS: браузер запрашивает example.com у DNS-сервера, получает IP-адрес 104.121.221.0, отправляет запрос на этот адрес и получает ответ сервера по HTTPS

В больших распределённых системах DNS обычно дополняют реестрами сервисов или балансировщиками нагрузки. Но в основе всегда нужен способ превратить имя в сетевой адрес.

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

Как распределённая система обнаруживает отказ узла?

Отказ узла в распределённой системе обнаруживают по пропавшим сигналам жизни: узлы регулярно шлют друг другу heartbeat-сообщения или распространяют сведения друг о друге по gossip-протоколу. Если сигналы от узла слишком долго не приходят, его считают упавшим.

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

Как работает механизм heartbeat?

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

Как работает gossip-протокол?

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

Состояние узлов gossip-протокола в момент t1: пять узлов S1–S5 на кольце, у каждого свой список известных ему сведений, например [1, 2, 4] у S1 и [1, 2, 8, 9] у S5

Эти техники помогают обнаруживать отказы в шумной и ненадёжной сети.

Как часто проверять, жив ли узел?

Приходится искать баланс между скоростью и точностью:

  • проверять часто — ложные срабатывания: вы считаете узел упавшим, хотя он работает;
  • проверять редко — система поздно реагирует на настоящие отказы.

Как упорядочить события в распределённой системе без общих часов?

Порядок событий в распределённой системе задают логические часы, а не системное время: часы разных машин расходятся, и синхронизировать их идеально почти невозможно. Часы Лэмпорта показывают, какое событие произошло раньше, а векторные часы вдобавок распознают события, произошедшие конкурентно.

Когда машины разнесены по миру, трудно понять, какое из двух событий на разных серверах произошло первым. Если решать «что когда произошло» через time.Now(), начинаются проблемы: у каждой машины свои локальные часы, и они расходятся. Поэтому идеально синхронизировать часы на многих машинах почти невозможно.

В распределённых системах точная метка времени нужна редко. На самом деле важно знать, какое событие было первым. Для этого используют алгоритмы, которые называют логическими часами.

Что такое часы Лэмпорта?

Часы Лэмпорта (Lamport clocks) — простой счётчик на каждой машине, который задаёт логический порядок событий. Правила такие:

  • каждый процесс начинает со счётчиком 0;
  • при каждом локальном событии счётчик увеличивается: LC = LC + 1;
  • при отправке сообщения к нему прикладывается текущее значение счётчика;
  • при получении сообщения счётчик обновляется до
    LC = max(local LC, received LC) + 1.

Пример с двумя машинами, M1 и M2:

  • M1 отправляет сообщение с LC = 5;
  • M2 получает его, когда её собственные часы показывают LC = 3;
  • M2 обновляет часы до max(3, 5) + 1 = 6.

Событие отправки получает LC = 5, событие получения — LC = 6, и правильный порядок «отправка < получение» сохраняется. Гарантия работает в одну сторону: если событие a произошло раньше события b, метка a меньше метки b.

Почему часы Лэмпорта не распознают конкурентные события?

Часы Лэмпорта говорят, что произошло раньше чего, но не говорят, были ли два события независимыми, — то есть не распознают конкурентность. Пример:

  • у M1 LC = 5, она выполняет событие A → LC(A) = 6;
  • у M2 LC = 2, она выполняет событие B → LC(B) = 3.

События независимы: машины между собой не общались. Но по часам Лэмпорта 3 < 6, и выглядит так, будто B произошло раньше A, хотя события были конкурентными. Значит, меньшая метка Лэмпорта сама по себе не доказывает, что событие было раньше: часы показывают только порядок, но не конкурентность. Эту проблему решают векторные часы.

Что такое векторные часы?

Векторные часы (vector clocks) расширяют идею часов Лэмпорта: каждая машина хранит не один счётчик, а вектор счётчиков — по одному на каждую машину системы. Идея звучит так: «А что если каждая машина знает порядок событий у всех остальных?»

Если машин N, вектор хранит N элементов, и каждый элемент M[i] считает события своей машины. Это значит:

  • каждое событие несёт снимок того, что машина знает о состоянии остальных;
  • получив сообщение, машина объединяет знания поэлементным максимумом.

Так у каждой машины появляется частичная история всей системы. С векторными часами система может определить:

  • какое событие произошло первым;
  • произошли ли два события одновременно (конкурентно).

Как работает выбор лидера в Raft?

Raft — алгоритм консенсуса, который выбирает лидера голосованием: все узлы стартуют как последователи (followers), не дождавшийся лидера узел становится кандидатом, и кандидат, набравший большинство голосов, становится лидером. Лидер координирует обновления, а последователи реплицируют его журнал.

  • Состояния: follower, candidate, leader
  • Решение: большинство голосов
  • Где: распределённые базы данных, системы консенсуса

Распределённой системе нужен единый источник истины для всех операций, поэтому удобно, чтобы одна машина была лидером. Лидер становится источником истины, а остальные машины (последователи) повторяют его состояние. Но если все серверы устроены одинаково, как решить, какой из них достоин стать лидером?

Как проходят выборы лидера в Raft?

Алгоритмы вроде Raft решают эту задачу чисто и предсказуемо:

  1. все узлы стартуют как последователи;
  2. если лидер не обнаружен, последователь переходит в состояние кандидата, чтобы стать лидером;
  3. проходит голосование, и узел, получивший большинство голосов, становится лидером.

После выборов лидер координирует обновления, а последователи реплицируют его журнал (log).

Схема состояний алгоритма Raft: follower по таймауту становится candidate и начинает выборы, candidate с большинством голосов становится leader, а при обнаружении действующего лидера или сервера с более высоким term узел возвращается в follower

Почему Raft так популярен?

Raft прост, у него чётко определённые состояния (follower, candidate, leader), и он даёт сильные гарантии безопасности. Его используют многие распределённые базы данных и системы консенсуса, например etcd и CockroachDB.

Какие бывают модели согласованности данных?

Модель согласованности определяет, какие значения увидят читатели, пока запись реплицируется между узлами. Основные модели — линеаризуемость (всегда последняя запись), последовательная согласованность (единый порядок операций без привязки к реальному времени) и согласованность в конечном счёте (реплики сойдутся, когда записи прекратятся).

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

Что такое линеаризуемость?

  • Строгость: максимальная
  • Цена: координация на каждую операцию, задержка

Линеаризуемость (linearizability) — самая строгая модель: каждая операция выглядит мгновенной и глобально упорядоченной, и читатель всегда видит последнюю запись. В масштабе она дорога, потому что каждая операция требует координации между машинами, и часто добавляет задержку.

Что такое последовательная согласованность?

  • Строгость: средняя
  • Цена: ниже, чем у линеаризуемости

Последовательная согласованность (sequential consistency) гарантирует, что операции каждого узла идут в его собственном порядке, а все узлы видят один и тот же общий порядок операций. Этот общий порядок не обязан совпадать с реальным временем, поэтому достичь такой модели проще, чем линеаризуемости.

Что такое согласованность в конечном счёте?

  • Строгость: слабая
  • Примеры: Cassandra

Согласованность в конечном счёте (eventual consistency) гарантирует: если новых записей нет, все реплики со временем сойдутся к одному значению. Её используют высокодоступные системы вроде Cassandra. Немедленную корректность здесь меняют на более высокую доступность и лучшую устойчивость к разделению сети.

МодельЧто гарантирует читателюЦенаГде встречается
ЛинеаризуемостьВсегда последняя запись, операции мгновенны и глобально упорядоченыКоординация на каждую операцию, задержкаТам, где устаревшее чтение недопустимо
ПоследовательнаяЕдиный для всех порядок операций, порядок каждого узла сохранёнМеньше координации: порядок не привязан к реальному времениПромежуточный вариант
В конечном счётеРеплики сойдутся, когда записи прекратятсяМожно прочитать устаревшее значениеCassandra

Что такое теорема CAP?

Теорема CAP утверждает, что распределённая система не может одновременно гарантировать согласованность (Consistency), доступность (Availability) и устойчивость к разделению сети (Partition tolerance). Разделения сети неизбежны, поэтому при разделении система выбирает между согласованностью и доступностью.

Выбор модели согласованности напрямую связан с теоремой CAP. Три свойства:

  • Согласованность (C): каждое чтение возвращает последнюю запись или ошибку.
  • Доступность (A): каждый запрос получает ответ, пусть даже устаревший.
  • Устойчивость к разделению (P): система продолжает работать, даже если сеть распалась на части.

Почему нельзя все три? Разделения сети (P) в распределённых системах — обычное дело. Когда разделение происходит, система должна выбрать:

  • согласованность (C) — перестать отдавать устаревшие данные, пока разделение не устранено;
  • доступность (A) — отдавать те данные, что есть, даже если они устарели.

Отсюда деление систем на CP, AP и CA, но все три свойства одновременно не получить никогда: компромисс есть всегда. Оговорка: CA возможна, только если разделений сети не бывает, то есть по сути на одном узле. Для системы, узлы которой связаны сетью, реальный выбор — между CP и AP.

Дальше — техники масштабирования: ради них распределённые системы и существуют.

Как масштабируют распределённые системы?

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

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

Как разделение ответственности помогает масштабировать систему?

Разделение ответственности (separation of concerns) — первый шаг при переходе от монолита к распределённой системе: все обязанности системы выписывают и разносят по отдельным частям. Тогда масштабировать можно только сервисы под высокой нагрузкой, а не всю систему целиком.

Что такое микросервисы?

Микросервисы — небольшие сервисы, на которые разбивают монолит. Каждый отвечает за конкретную бизнес-задачу и деплоится независимо.

Что такое API Gateway?

API Gateway — единая точка входа в систему. Он направляет каждый запрос в нужный микросервис и берёт на себя общие задачи: аутентификацию, rate limiting (ограничение частоты запросов) и мониторинг.

API Gateway маршрутизирует запросы: GET /product/{id} в Catalogue Service, POST /payment/{order_id} в Payment Service, GET /cart/{cart_id} в Cart Service

Аналогия: кухня ресторана

Каждый повар специализируется на своём блюде, а шеф (API Gateway) следит, чтобы заказы были согласованы и доставлены правильно.

Что такое CQRS?

CQRS (Command Query Responsibility Segregation, разделение ответственности команд и запросов) — паттерн, который разносит запись и чтение данных по разным сторонам системы. Сторона команд оптимизирована под запись, сторона запросов — под быстрое чтение, и каждая масштабируется независимо.

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

  • Сторона команд (запись) обрабатывает изменения данных системы и оптимизирована под нагрузку с преобладанием записи.
  • Сторона запросов (чтение) отдаёт данные и оптимизирована под быстрое чтение, в том числе через кэширование.

CQRS особенно полезен в системах с преобладанием чтения: социальных лентах, аналитических дашбордах, каталогах товаров.

Что такое асинхронный обмен сообщениями?

Асинхронный обмен сообщениями — способ выполнять работу в фоне, не заставляя пользователя ждать. Два основных паттерна — очередь, в которой каждое сообщение обрабатывает один потребитель, и Pub/Sub, в котором сообщение получает каждый подписчик.

Пример: при каждой регистрации пользователя нужно отправить приветственное письмо. Отправка занимает время, и если миллион человек зарегистрируется одновременно, немедленная отправка перегрузит систему. Асинхронный обмен выносит эту работу в фон и не тормозит основной пользовательский сценарий.

Как работает очередь сообщений?

В очереди сообщений производители (producers) добавляют сообщения в очередь, а потребители (consumers) обрабатывают их в своём темпе. Очередь гарантирует доставку, даже если потребитель медленный или временно недоступен.

Очередь сообщений: издатель кладёт сообщения M1, M2, M3 в очередь, каждое сообщение забирает и обрабатывает один из трёх потребителей, после чего оно удаляется из очереди

Как работает Pub/Sub?

В модели Pub/Sub (publish/subscribe, «издатель — подписчик») одно сообщение рассылается сразу многим подписчикам. Модель подходит для уведомлений, потоков событий и обновления кэшей.

Pub/Sub: издатель публикует события M1, M2, M3, и каждый из трёх потребителей получает свою копию события и обрабатывает её независимо

Асинхронный обмен развязывает компоненты, держит пользовательские пути быстрыми и помогает системе выдерживать огромные нагрузки.

Что такое шардирование данных?

Шардирование (sharding, или партиционирование) — разбиение данных между многими серверами, когда один сервер уже не может эффективно хранить и обрабатывать их все. Данные делят по диапазонам ключей, по хешу ключа или через консистентное хеширование.

Когда система дорастает до миллионов пользователей, один сервер перестаёт справляться с данными. Способов разделить их несколько.

Как работает шардирование по диапазону?

При шардировании по диапазону (range-based partitioning) данные делят по диапазонам ключей. Например, пользователи A–H на одном сервере, I–Q на втором, R–Z на третьем.

Шардирование по диапазону: балансировщик нагрузки распределяет запросы со всего мира между тремя базами с диапазонами ключей A–H, I–Q и R–Z

Как работает шардирование по хешу?

При шардировании по хешу (hash-based partitioning) хеш-функция определяет, на каком сервере хранится каждый элемент. Данные распределяются равномерно, но запросы по диапазону становятся сложнее.

Что такое консистентное хеширование?

Консистентное хеширование (consistent hashing) уменьшает объём данных, которые приходится перемещать при добавлении или удалении серверов. Серверы и данные размещаются на кольце, и каждый элемент хранится на ближайшем по часовой стрелке сервере. Этот подход используют Cassandra и Memcached.

Консистентное хеширование: хеш от user_id указывает точку на кольце, данные D1–D7 хранятся на ближайшем по часовой стрелке сервере из S1–S5, и для этого user_id данные лежат на S2

Серверы могут отказать или добавиться в любой момент, поэтому система должна уметь перемещать данные без простоя. Эффективная перебалансировка предотвращает горячие точки (hotspots) и сохраняет плавную работу системы.

Что такое балансировка нагрузки?

Балансировка нагрузки — распределение входящих запросов между многими серверами так, чтобы ни один из них не был перегружен. Балансировщик уровня L4 (TCP) распределяет трафик простыми правилами вроде round-robin, а балансировщик уровня L7 (HTTP) маршрутизирует по содержимому запроса.

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

Как работает балансировщик L4?

Балансировщик L4 (TCP) использует простые методы распределения, например round-robin, и менее гибок.

Балансировщик L4: запрос клиента с адреса 192.168.1.10:4567 по round-robin уходит на сервер s2; балансировщик не терминирует SSL и не видит HTTP-заголовки и тело запроса

Как работает балансировщик L7?

Балансировщик L7 (HTTP) маршрутизирует умнее: решает, куда отправить запрос, по результатам health checks, cookies или путям URL.

Балансировщик L7 терминирует SSL и проверяет заголовки, JWT-токены и пути: запросы /api/v1/cart уходят на сервер s1, а /api/v1/orders — на сервер s2

Как балансировщик узнаёт о здоровых серверах?

Через обнаружение на стороне сервера (server-side discovery) и health checks: серверы регистрируются в сервисе обнаружения вроде ZooKeeper, и балансировщик направляет трафик только на здоровые серверы.

Что такое репликация данных?

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

Как работает репликация с одним лидером?

В репликации с одним лидером (single-leader) записи принимает один узел, а реплики обслуживают чтение. Это даёт строгую согласованность записей.

Как работает репликация с несколькими лидерами?

В репликации с несколькими лидерами (multi-leader) записи могут принимать несколько узлов, но тогда нужно разрешать конфликты между ними.

Как работает репликация без лидера?

В репликации без лидера (leaderless) запись принимает любой узел. Чтения и записи подчиняются правилам кворума — так устроены, например, Amazon Dynamo и Cassandra.

Три схемы репликации: multi-leader — записи идут на W1 и W2 с разрешением конфликтов; leaderless — чтение и запись на любой из узлов S1–S4; single-leader — записи идут на P1, который реплицирует данные на S1–S4 для чтения

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

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

Как сделать распределённую систему отказоустойчивой?

Отказоустойчивость распределённой системы строят на паттернах, которые поглощают, изолируют и восстанавливают сбои. От медленных зависимостей защищают таймауты, повторы с экспоненциальной задержкой и circuit breaker, а от перегрузки входящим трафиком — rate limiting и health checks на балансировщике.

Почему распределённые системы отказывают?

Чтобы построить устойчивую систему, сначала нужно понять, откуда берутся сбои:

  • Отказы железа: диски выходят из строя, память повреждается, сетевые интерфейсы отказывают.
  • Программные сбои: баги, утечки памяти или ошибки конфигурации останавливают сервисы.
  • Каскадные отказы: один упавший сервис вызывает отказы сервисов, которые от него зависят.

Распределённые системы отказывают часто: вопрос не в том, «случится ли», а в том, когда. Паттерны отказоустойчивости нужны, чтобы поглощать, изолировать сбои и восстанавливаться после них, а не делать вид, что сбоев не будет.

Как защититься от сбоев зависимых сервисов (downstream)?

Downstream-сбой происходит, когда ваш сервис зависит от другого сервиса, который тормозит или не отвечает. Защищают три паттерна:

  • Таймауты. Не ждать ответа вечно: задать максимальное время ожидания, чтобы запросы и потоки не блокировались.
  • Повторы с экспоненциальной задержкой (retries with exponential backoff). Если запрос не прошёл, повторить его, но с каждой попыткой ждать дольше (1 с, 2 с, 4 с…). Так вы не перегружаете и без того страдающий сервис и предотвращаете каскадные отказы.
  • Circuit breaker (предохранитель). Если сервис продолжает отказывать, на время перестать слать ему запросы — «выбить предохранитель». Это не даёт одному сломанному сервису уронить всю систему. После паузы (cooldown) запросы можно пробовать снова.

Как защититься от перегрузки входящими запросами (upstream)?

Upstream-сбой происходит, когда система получает больше запросов, чем способна обработать. Перегрузку предотвращают две техники:

  • Rate limiting / throttling. Ограничить число запросов, которые может отправить клиент. Это защищает систему и не даёт перегрузить сервисы ниже по цепочке.
  • Health checks на балансировщике. Периодически проверять здоровье каждого сервера. Если сервер тормозит или отказывает, балансировщик перестаёт слать ему трафик, и запросы уходят только туда, где их могут обработать.

Сравнительная таблица техник распределённых систем

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

ТехникаКакую проблему решаетКак работаетГде встречается
Heartbeat, gossipОбнаружение отказавших узловРегулярные сигналы «я жив», тишина — признак отказаКластеры баз данных и брокеров
Логические часыПорядок событий без общих часовСчётчик (Лэмпорт) или вектор счётчиков (векторные часы)Упорядочивание событий, поиск конфликтов
Выбор лидера (Raft)Единый источник истиныГолосование, лидер получает большинствоРаспределённые базы данных, etcd
Микросервисы, API GatewayМасштабировать только нагруженную частьСервисы по бизнес-задачам, единая точка входаПереход от монолита
CQRSРазные нагрузки на чтение и записьОтдельные стороны команд и запросовЛенты, дашборды, каталоги
Очередь, Pub/SubДолгая работа тормозит пользователяБрокер доставляет сообщения в фонеKafka, письма, уведомления
ШардированиеДанные не помещаются на один серверРазбиение по диапазону или хешу ключаCassandra
Консистентное хешированиеМассовый переезд данных при смене числа серверовКольцо, данные на ближайшем по часовой сервереCassandra, Memcached
Балансировка L4 / L7Перегрузка отдельных серверовRound-robin (L4) или маршрутизация по запросу (L7)Вход в любой кластер
РепликацияДоступность и масштабирование чтенияКопии данных: один лидер, несколько или без лидераБазы данных, Kafka
Таймауты, ретраи, circuit breakerМедленные и упавшие зависимостиОграничить ожидание, повторять с паузой, отключать сбойный сервисМежсервисные вызовы
Rate limitingПерегрузка входящим трафикомЛимит запросов на клиентаAPI Gateway, публичные API

Какую технику распределённых систем выбрать под задачу?

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

СимптомЧто применить
Одна машина не справляется с нагрузкойГоризонтальное масштабирование и балансировщик нагрузки
Данные не помещаются на один серверШардирование
При добавлении серверов переезжает слишком много данныхКонсистентное хеширование
Чтений намного больше, чем записейCQRS, реплики для чтения, кэширование
Долгая операция тормозит ответ пользователюОчередь сообщений
Одно событие нужно нескольким сервисамPub/Sub
Зависимый сервис тормозит или падаетТаймауты, ретраи с backoff, circuit breaker
Клиенты шлют больше запросов, чем система выдерживаетRate limiting
Трафик идёт на упавший серверHealth checks на балансировщике
Нужно понять порядок событий на разных машинахЛогические или векторные часы
Нужен единый источник истины для записейЛидер, выбранный через Raft, и репликация с одним лидером

Какие ошибки чаще всего допускают при проектировании распределённых систем?

  • Считают сеть надёжной и не обрабатывают потерю, дублирование и переупорядочивание сообщений.
  • Упорядочивают события по системному времени разных машин, например через time.Now().
  • Повторяют неидемпотентные операции: дубль сообщения выполняет операцию дважды.
  • Ждут ответа зависимого сервиса без таймаута.
  • Повторяют запросы без экспоненциальной задержки и добивают и без того перегруженный сервис.
  • Проверяют живость узлов слишком часто и ловят ложные срабатывания — или слишком редко и поздно замечают реальный отказ.
  • Требуют линеаризуемости там, где хватило бы согласованности в конечном счёте, и платят за это задержкой.
  • Пускают внутренние RPC без TLS.

В чём главная сложность распределённых систем?

Главная сложность распределённых систем в том, что это набор ненадёжных частей, которые пытаются работать как одно целое. Сбои устранить нельзя, поэтому систему проектируют так, чтобы сбои предусматривать, локализовать и переживать незаметно для пользователя.

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

Каждая распределённая система идёт на компромиссы: согласованность против доступности, простота против масштаба, безопасность против скорости. Цель не в том, чтобы устранить хаос, а в том, чтобы управлять им так, чтобы пользователи ничего не заметили.

Проектировать приходится со скромностью: ожидать отказов и опираться на паттерны и протоколы, которые заставляют множество независимых частей работать как одна система. Поэтому распределённые системы одновременно и раздражают, и завораживают: отказать может любая система, но строить нужно такие, что продолжают работать вопреки отказам.

Частые вопросы о распределённых системах

Микросервисы — это распределённая система?

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

Можно ли по теореме CAP выбрать CA?

Только если разделений сети не бывает, то есть по сути в системе на одном узле. Узлы распределённой системы связаны сетью, и разделения сети в ней неизбежны, поэтому на практике при разделении выбирают между согласованностью (CP) и доступностью (AP).

Чем часы Лэмпорта отличаются от векторных часов?

Часы Лэмпорта хранят на каждой машине один счётчик и гарантируют, что если событие a произошло раньше события b, то метка a меньше метки b. Отличить конкурентные события они не могут. Векторные часы хранят по счётчику на каждую машину системы и определяют как порядок событий, так и то, что два события произошли конкурентно.

Чем шардирование отличается от репликации?

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

Чем очередь сообщений отличается от Pub/Sub?

В очереди каждое сообщение забирает и обрабатывает один потребитель, после чего сообщение удаляется из очереди. В Pub/Sub каждое сообщение получает каждый подписчик и обрабатывает свою копию независимо. Очередь подходит для распределения работы, Pub/Sub — для уведомлений, потоков событий и обновления кэшей.

Kafka — это очередь или Pub/Sub?

Kafka — распределённый отказоустойчивый журнал событий, и она поддерживает оба сценария. Разные группы потребителей (consumer groups) читают один топик независимо, как в Pub/Sub, а внутри одной группы партиции топика делятся между потребителями, как в очереди. Сообщения после чтения не удаляются, а хранятся в журнале заданное время.

Зачем нужна идемпотентность в распределённых системах?

Сеть может потерять ответ или доставить сообщение дважды, поэтому клиенты повторяют запросы. Идемпотентная операция при повторном выполнении даёт тот же результат, что и при однократном, и повтор не создаёт двойной платёж или двойной заказ. На идемпотентных операциях и надёжной доставке сообщений построена, например, событийная архитектура Stripe.


Источники