Мессенджер на миллиард пользователей: шардирование неизбежно — но что пойдёт не так дальше?

public

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