Источник: Распределенные системы и system design
Когда одной базы уже недостаточно по объёму, чтению и записи, шардирование выглядит неизбежным решением. Но на этом масштабирование не заканчивается — оно только начинается. Как только сообщения распределяются по шардам, приходится отдельно решать проблему порядка, генерации ID, Kafka, WebSocket-соединений и доставки между регионами.
В авторском ролике разбирается system design мессенджера на миллиард пользователей. Требования намеренно ограничены: создать чат, добавить или удалить участника, отправить текстовое сообщение и доставить его в реальном времени пользователям, которые сейчас онлайн. Статусы прочтения, реакции и медиа автор не рассматривает — они могут быть добавлены поверх предложенной архитектуры.
Кому смотреть: backend-разработчикам и тем, кто готовится к system design-интервью, но хочет увидеть не список модных технологий, а цепочку инженерных компромиссов: какое решение принимается первым и какие новые проблемы оно создаёт дальше.
Из этого можно взять в работу: при шардировании сначала зафиксируйте инварианты. В этом дизайне сообщения одного чата должны оставаться в одном логическом потоке, а их ID — возрастать. Именно эти требования определяют и ключ шардирования базы, и распределение сообщений в Kafka.
Сначала ломается не масштаб, а простая схема
В исходной архитектуре есть connection server и база данных. Пользователь отправляет сообщение на connection server, тот сохраняет его в базу, находит участников чата с открытыми WebSocket-соединениями и рассылает сообщение.
Такая схема быстро показывает первую проблему: отправитель может долго ждать ответа, потому что рассылка зависит от количества участников и скорости их соединений. Можно сохранить сообщение синхронно, а рассылку выполнять асинхронно, но тогда падение connection server способно оставить сообщение в базе без выполненной доставки.
Автор рассматривает CDC-подход: Debezium читает журнал операций базы, отправляет изменения в Kafka, а отдельный consumer реагирует на появление нового сообщения и выполняет рассылку. Это сохраняет связь между записью и событием, но добавляет ещё одну инфраструктурную компоненту и ограничивает выбор базы — особенно если у неё нет удобного единого журнала операций.
Более простая схема сразу пишет сообщение в Kafka. Consumer читает его, сохраняет в базу и выполняет рассылку. При сбое он повторяет операцию, поэтому запись в базу должна быть идемпотентной. Подтверждение отправителю приходит не в момент записи в Kafka, а после того, как сообщение обработано и отправлено по WebSocket.
База: шардирование — только первый вопрос
Одна база не выдержит одновременно нагрузку на чтение и запись и не вместит все сообщения. Автор предлагает Cassandra с безлидерной репликацией — ради доступности и пропускной способности — и шардирование по chat ID.
Это важный выбор структуры данных: сообщения одного чата оказываются на одной машине, поэтому историю чата можно читать из одного шарда. Внутри шарда сообщения сортируются по message ID, который должен постоянно возрастать в рамках конкретного чата.
Здесь появляется следующий подводный камень. Автоинкремент неудобен при безлидерной репликации: записи могут приходить на разные реплики в разном порядке. Timestamp тоже не даёт гарантии — после падения consumer новые часы могут оказаться чуть «позади» предыдущего процесса. В предложенной схеме message ID строится на offset сообщения в Kafka: offset внутри партиции возрастает, а сообщения одного чата направляются в одну партицию.
Kafka тоже приходится шардировать
Одной машины Kafka и одного consumer недостаточно. Сообщения распределяются по партициям на основании chat ID, а на каждую партицию назначается отдельный consumer. Так сохраняется порядок: сообщения одного чата всегда идут в одну партицию, а offset внутри неё увеличивается.
Но добавить новые партиции без последствий нельзя. Если хеширование начнёт отправлять сообщения уже существующего чата в новую партицию, один чат окажется в нескольких потоках и нарушится логика генерации ID. Один из предложенных вариантов — заранее заложить большое количество партиций. Другой — заводить отдельные топики для диапазонов chat ID, чтобы масштабировать систему добавлением новых топиков и не разносить существующие чаты по разным потокам.
Connection servers: где сейчас пользователь?
Даже после масштабирования базы и Kafka остаётся connection server. Один сервер не сможет держать соединения со всеми пользователями, поэтому появляется кластер connection servers и балансировщик, который направляет пользователей на менее загруженные узлы.
Теперь consumer должен знать, на каком connection server находится конкретный пользователь. Для этого в Redis хранится mapping user ID → connection server. При отключении пользователя запись удаляется, а на случай падения сервера без корректного удаления используется TTL.
У Redis-схемы есть два недостатка: кластер не забирает автоматически часть ключей при добавлении новой машины, а consumer делает дополнительный round trip за адресом connection server. В качестве более сложной альтернативы автор предлагает собственный кластер с внутренним key-value storage, consistent hashing и координацией через etcd. Но здесь важен практический вывод: кастомная схема гибче, а Redis проще и может быть предпочтительнее.
Когда пользователи оказываются в разных регионах
Для миллиарда пользователей одного региона недостаточно: пользователи по всей планете не должны ходить в один далёкий backend и получать высокую latency. У каждого пользователя есть home region, к которому он подключается по WebSocket. Чат при этом привязан к региону — например, к региону пользователя, который его создал.
Сначала можно попробовать подключать участника из другого региона прямо к региону, где живёт чат. Но это создаёт нестабильное дальнее WebSocket-соединение, требует буферизации и приводит к взрывному росту числа соединений между регионами.
Более удачная схема разделяет локальную и межрегиональную доставку. Локальным пользователям сообщение отправляется по WebSocket. Для участников из других регионов оно попадает в outgoing Kafka topic, откуда доезжает в целевой регион и там обрабатывается локальным consumer-ом.
Чтобы не соединять каждый регион с каждым, вводится hub-регион. Регионы отправляют сообщения в хаб, а нужный регион забирает из хаба свой поток. Так вместо полной связанности получается более простая схема: каждый регион знает хаб, а хаб знает регионы. Хаб не становится бутылочным горлышком, потому что предполагается: большинство чатов остаётся внутри одного региона, и их сообщения через хаб не проходят.
Дополнительные варианты
В конце автор разбирает несколько случаев, которые легко забыть в основной схеме:
- миграцию пользователя в другой home region с временным подключением к старому и новому регионам;
- отправку сообщения участником из региона, отличного от региона-владельца чата;
- чтение пропущенных сообщений после возвращения пользователя в сеть;
- альтернативу WebSocket — polling;
- отдельное хранение пользовательских копий сообщений, шардированное по user ID, что упрощает polling и статусы прочтения, но увеличивает объём данных.
Главная мысль
Шардирование — не признак того, что архитектуру усложнили без необходимости. При таком масштабе без него не обойтись. Но шардирование базы сразу требует новых решений: как сохранить порядок сообщений, как распределить Kafka, как найти нужный connection server и как доставить сообщение пользователю в другом регионе.
Хороший system design — это не схема, где заранее выбрали Cassandra, Kafka и Redis. Это последовательное объяснение того, какой инвариант мы сохраняем на каждом шаге и какую новую проблему создаёт очередное масштабирование.