open atlas
↑ К треку
Разборы System Design SDC · 01 · 04

Спроектируйте распределённую очередь сообщений

Проектируем Kafka-подобный лог: топики делятся на партиции для параллелизма, реплицируются через ISR для надёжности, оффсеты для реплея и компромисс семантики доставки — at-most/at-least/exactly-once — который выбирает любая async-система.

SDC Senior ◷ 32 min
Уровень
ОсновыJuniorMiddleSenior

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

Требования

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

Функциональные

  • Произвести (produce): дописать сообщение в именованный топик.
  • Потребить (consume): читать сообщения из топика, по порядку, начиная с выбранной позиции.
  • Множество независимых консьюмеров: несколько разных систем читают один топик, не мешая друг другу (скоринг фрода, бухгалтерия, аналитический пайплайн — все читают поток платежей).
  • Реплей: консьюмер может перемотать назад и перечитать прошлые сообщения (пропавшая фича из хука).
  • Удержание (retention): сообщения хранятся настроенное окно (время или размер), а не удаляются при доставке.

Нефункциональные

  • Высокая пропускная способность. Миллионы сообщений/с — это пожарный шланг, не почтовый ящик.
  • Надёжность. Подтверждённое сообщение обязано пережить крах брокера; потеря платёжного события недопустима.
  • Горизонтальный масштаб. Добавь брокеров — добавь мощности; ни один брокер не потолок.
  • Гарантия порядка — хотя бы внутри партиции.
  • Настраиваемая семантика доставки. Разные консьюмеры терпят разные режимы отказа (пайплайн метрик может уронить сообщение; бухгалтерия — нет).

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

Оценка

  • Цель — 1 миллион сообщений/с на пике, среднее сообщение 1 КБ~1 ГБ/с пропускной способности записи. Это заголовочное число, и оно исключает любой дизайн на одном брокере.
  • Удержание 7 дней: 1 ГБ/с × 86 400 с × 7~600 ТБ лога на диске (до репликации). При факторе репликации 3 — ~1,8 ПБ. Лог живёт на диске, не в RAM — так что последовательный дисковый I/O и есть вся игра.
  • Пропускная способность на партицию упирается, может, в ~10 МБ/с для упорядоченных реплицированных записей, так что 1 ГБ/с требует порядка ~100+ партиций, разбросанных по брокерам — число партиций и есть то, чем покупают пропускную способность.
  • Консьюмеры: если чтения в основном последовательны с хвоста лога, page cache ОС обслуживает свежие сообщения из RAM, так что здоровый кластер почти не делает случайных дисковых чтений — паттерн доступа и есть то, что заставляет числа работать.

Ключевая мысль из чисел: последовательная дозапись на диск плюс партиционирование — то, что превращает «медленный» диск в пожарный шланг на 1 ГБ/с, ведь последовательная пропускная способность диска соперничает с памятью, тогда как случайный I/O был бы в ~1000 раз медленнее.

Высокоуровневый дизайн

Топик делится на партиции; каждая партиция — упорядоченный лог только на дозапись, живущий на брокере-лидере и реплицируемый на брокеры-фолловеры. Продюсеры дописывают в партиции (по ключу или round-robin). Консьюмеры в группе консьюмеров делят партиции между собой и отслеживают оффсет на партицию.

Глубокое погружение

Партиции: единица параллелизма и порядка

Топик делится на партиции, и партиция — атом всей системы. Из неё следуют три вещи:

  • Параллелизм. Разные партиции живут на разных брокерах и пишутся/читаются независимо, так что пропускная способность масштабируется с числом партиций — так доходишь от лимита одного брокера до 1 ГБ/с.
  • Порядок — на партицию, не на топик. Сообщения внутри партиции строго упорядочены по оффсету; между партициями нет глобального порядка. Так что если нужна упорядоченная обработка событий пользователя, надо маршрутизировать все его сообщения в одну и ту же партицию — делается хешированием ключа партиции (например user_id). Нет ключа → round-robin → максимальный разброс, но нет гарантии порядка.
  • Параллелизм консьюмеров ограничен партициями. Внутри группы консьюмеров каждую партицию потребляет ровно один консьюмер, так что активных консьюмеров может быть не больше, чем партиций. Число партиций — это и ручка пропускной способности, и максимум параллелизма консьюмеров — выбирается наперёд и неловко меняется позже.

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

Почему это работает

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

Репликация и ISR

У каждой партиции один лидер (берёт все чтения и записи) и несколько фолловеров, реплицирующих лог лидера. Надёжность держится на множестве in-sync реплик (ISR): фолловеров, что сейчас догнали лидера. Запись считается закоммиченной лишь когда все члены ISR её имеют, а acks управляет строгостью продюсера:

  • acks=0 — выстрелил-и-забыл; быстрее всего, может потерять данные.
  • acks=1 — подтверждение лишь лидера; теряет данные, если лидер умер до репликации на фолловеров.
  • acks=all — каждый член ISR должен иметь сообщение до подтверждения; переживает отказ лидера без потерь, пока ISR не пуст.

Когда лидер умирает, новый лидер избирается из ISR, так что у него уже есть каждое закоммиченное сообщение — ни одна закоммиченная запись не теряется. Настраиваемый min.insync.replicas задаёт, сколько реплик должны быть in-sync для прохода записи acks=all; поставь 2 (при факторе репликации 3) — и переживёшь отказ одного брокера, всё ещё отклоняя записи, если надёжность была бы скомпрометирована.

Граничные случаи

Что будет, если все фолловеры отстанут, ISR сожмётся до одного лидера, а потом лидер умрёт? Сталкиваешься с дилеммой unclean-leader-election: либо отказаться избирать любого лидера (партиция уходит офлайн — недоступна, но никогда не отдаёт потерянные данные), либо позволить отставшему фолловеру стать лидером (партиция возвращается доступной, но молча роняет каждое сообщение, что имел мёртвый лидер и которого нет у нового). Это компромисс CAP, овеществлённый на уровне партиции: unclean.leader.election.enable=false выбирает согласованность/надёжность (лучше лежать, чем потерять данные), true выбирает доступность (вернуться, приняв потерю данных). Платёжный лог ставит false; пожарный шланг метрик может поставить true. Третьего нет — когда единственная in-sync копия исчезла, нельзя иметь и живую партицию, и полный лог.

Оффсеты, группы консьюмеров и семантика доставки

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

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

Но когда ты коммитишь оффсет относительно того, когда обрабатываешь сообщение, и определяет семантику доставки — важнейшее решение корректности в любой async-системе:

  • At-most-once: коммитить оффсет до обработки. Если консьюмер упал посреди обработки, сообщение пропускается при перезапуске — возможна потеря, нет дублей. (Приемлемо для метрик.)
  • At-least-once: коммитить оффсет после обработки. Если консьюмер упал после обработки, но до коммита, сообщение перечитывается и переобрабатывается — возможны дубли, нет потерь. (Обычный дефолт.)
  • Exactly-once: ни дублей, ни потерь. По-настоящему сложно — требует либо идемпотентной обработки (консьюмер дедуплицирует по ID сообщения, так что переобработка безвредна), либо транзакции, атомарно коммитящей оффсет и побочный эффект вместе. Kafka предлагает транзакционных продюсеров/консьюмеров для этого, но большинство систем добиваются «фактически раз» через at-least-once доставку плюс идемпотентный консьюмер.
Выбери лучший вариант

Консьюмер платёжных событий читает сообщения из Kafka-партиции и вызывает внешний сервис бухгалтерии для каждого. Консьюмер падает в середине батча. Какую семантику доставки и стратегию коммита смещения следует выбрать?

Узкие места и компромиссы

  • Число партиций — почти-односторонее решение. Слишком мало — упрёшься в пропускную способность и параллелизм консьюмеров; слишком много — платишь координацией, файловыми дескрипторами и ребалансами. Увеличение партиций позже ломает порядок по ключу (ключ может хешироваться в новую партицию), поэтому размеряешь щедро наперёд.
  • Горячая партиция / перекошенный ключ. Плохой ключ партиции (например по стране, где одна доминирует) шлёт почти весь трафик в одну партицию и одному консьюмеру — версия горячего шарда для лога, и весь кластер помочь не может.
  • At-least-once значит, что консьюмеры обязаны быть идемпотентными. Обычный дефолт отдаёт дубли при ретрае, так что каждый консьюмер с побочными эффектами (списать карту, отправить письмо) обязан дедуплицировать по ID сообщения — требование, о котором команды забывают, пока кого-то дважды не спишут.
  • Штормы ребалансов. Когда консьюмер входит или выходит из группы, партиции переназначаются и потребление паузится; частая текучка членства (мигающие консьюмеры, долгие паузы GC) вызывает повторные ребалансы, останавливающие группу — так что тюнинг живости консьюмеров операционно важен.
  • Удержание против стоимости. Долгое удержание включает реплей, но стоит петабайтов; окно — прямой размен деньги-против-восстановимости, и «бесконечное удержание ради реплея» — реальный счёт.
Вспомните перед уходом
  1. 01
    Почему очередь — лог, а не почтовый ящик, и что это открывает?
  2. 02
    Что даёт партиция и почему порядок лишь на партицию?
  3. 03
    Как ISR и время коммита оффсета каждое определяют гарантию?
Итог

Распределённая очередь сообщений формы Kafka стоит на одном переосмыслении из хука: очередь — это лог, а не почтовый ящик — только на дозапись, хранимый окно и читаемый по отслеживаемому консьюмером оффсету, а не удаляемый при доставке. Это открывает реплей, множество независимых групп консьюмеров и восстановление после сбоя из одной структуры. Оценка (1 ГБ/с, ~600 ТБ за 7 дней) диктует архитектуру: последовательная дозапись на диск плюс партиционирование превращают медленный диск в пожарный шланг, а партиции — атом: они дают параллелизм (пропускная способность растёт с числом), порядок на партицию (поэтому стабильный ключ партиции решает, где живёт семантика порядка) и потолок параллелизма консьюмеров. Надёжность идёт от репликации и ISR: при acks=all запись коммитится, лишь когда все in-sync реплики её имеют, а новый лидер избирается из ISR, так что ни одна закоммиченная запись не теряется — а дилемма unclean-leader-election — это CAP овеществлённый (лежать офлайн или вернуться с потерями). Сложнейшее решение — семантика доставки, заданная тем, когда ты коммитишь оффсет: at-most-once (потеря), at-least-once (дубли) или exactly-once (идемпотентный консьюмер или атомарная транзакция) — а прагматичный ответ почти везде — at-least-once плюс идемпотентный консьюмер. Ловушки — число партиций, что нелегко поменять, горячая партиция от плохого ключа, штормы ребалансов и петабайтный счёт, который покупает долгое удержание. Теперь, когда видишь задачу на проектирование очереди — или отчёт об инциденте, где консьюмер дважды списал деньги — знаешь, за какой рычаг тянуть первым: какова семантика доставки и идемпотентен ли консьюмер?

Практика

Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.

вспомнитьприменитьуглубить0 из 7 завершено

Что-то непонятно?

Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.

хоткеи развернуть
поиск
K
пред. пьеса
k
след. пьеса
j
тиры
t
это меню
?
sources3
expand
  1. 01
  2. 02
  3. 03

Trademarks belong to their respective owners. Editorial reference only.