Раздел 4 · Данные и очереди — глава 4.2
Kafka, RabbitMQ и Redis: зачем нужны и как выбирать
Все три продукта умеют передавать данные между сервисами, но решают разные задачи. Kafka — сохраняемый журнал событий с повторным чтением, RabbitMQ — брокер очередей с гибкой маршрутизацией, Redis — быстрое хранилище структур данных в памяти. На собеседовании ценится не список терминов, а способность выбрать инструмент, назвать гарантию доставки и объяснить, что произойдёт при сбое.
§1Карта выбора за минуту
| Инструмент | Модель | Когда подходит | Главное ограничение |
|---|---|---|---|
| Kafka | Распределённый append-only log | Потоки событий, CDC, аудит, аналитика, replay, большой throughput | Порядок только внутри partition; эксплуатация тяжелее обычной очереди |
| RabbitMQ | Exchange маршрутизирует сообщения в очереди | Фоновые задания, команды сервисам, сложная маршрутизация, request/reply | После ack сообщение обычно удалено; массовый replay не его сильная сторона |
| Redis | In-memory key-value и структуры данных | Кэш, сессии, счётчики, rate limit, leaderboard, быстрые временные данные | Память дорога; durability и async replication требуют осознанного компромисса |
Граница не абсолютна: у RabbitMQ есть Streams, у Redis — Streams, а Kafka можно использовать как очередь. Но на собеседовании сначала называйте основную модель, затем исключение. «Умеет» не означает «лучший выбор».
Спросят
«Почему не положить всё в Redis?» — Потому что быстрый доступ не равен журналу событий или надёжной очереди. Для кэша допустима потеря и восстановление из основной БД. Для бизнес-события нужно определить retention, replay, подтверждение записи, репликацию и поведение при файловере. Выбор идёт от RPO/RTO, нагрузки и семантики, а не от одного показателя latency.
§2Гарантии доставки: где появляются потери и дубли
At-most-once — сообщение обрабатывается ноль или один раз: подтвердили слишком рано или не повторяем после ошибки, поэтому возможна потеря. At-least-once — повторяем до подтверждения: не теряем подтверждённое, но возможны дубли. Exactly-once — удобный термин, но сквозная гарантия сложна: транзакция Kafka не откатит уже отправленное письмо или списание во внешней БД.
Рабочая стратегия — at-least-once плюс идемпотентный consumer. Сообщение получает уникальный event_id; consumer сохраняет обработанные ID или делает операцию естественно идемпотентной, например UPSERT по бизнес-ключу. Повтор не меняет итог второй раз.
Ещё одна классическая проблема — dual write: приложение записало заказ в PostgreSQL, но упало до публикации события. Или событие отправило, а транзакция БД откатилась. Решение — transactional outbox: бизнес-изменение и строка события пишутся одной транзакцией в БД, а отдельный relay/CDC позже гарантированно публикует outbox. Consumer всё равно должен терпеть дубли.
Грабли
Фраза «у нас exactly-once, поэтому дубликатов не бывает» почти всегда вызывает уточняющий вопрос. Скажите границы гарантии: кто подтверждает запись, где хранится offset, входит ли внешняя БД в ту же транзакцию и что будет при таймауте, когда результат неизвестен.
§3Kafka: topic, partition, offset и consumer group
Topic — именованный поток. Он разбит на partitions; каждая partition — упорядоченный append-only log. Запись получает монотонный offset внутри своей partition. Сообщение не удаляется после чтения: Kafka хранит его по времени/размеру retention или оставляет последнюю запись на ключ при log compaction.
Глобального порядка между partitions нет. Если нужен порядок событий одного заказа, используйте order_id как key: одинаковый ключ попадёт в одну partition. Цена — горячий ключ способен создать hotspot.
Consumer group делит работу: одна partition в конкретный момент назначена максимум одному consumer этой группы. Поэтому при 6 partitions группа масштабируется максимум до 6 одновременно читающих consumers; седьмой простаивает. Другая группа читает тот же topic независимо и имеет собственные offsets — например billing и analytics получают всю историю каждый.
Изменение состава группы или числа partitions вызывает rebalance: partitions перераспределяются, обработка может кратко остановиться. Долгая обработка без poll, нестабильная сеть или постоянно перезапускающиеся поды создают rebalance storm. Основной показатель отставания — consumer lag: разница между последним offset в partition и подтверждённым offset группы.
Спросят
«Consumer прочитал сообщение и упал. Что будет?» — Зависит от момента commit offset. Commit до обработки даёт риск потери; после обработки — риск повторной обработки при падении между действием и commit. Типовой выбор — commit после успешной обработки, at-least-once и идемпотентность.
§4Надёжность Kafka и цена настроек
У partition есть leader, через которого идут чтение и запись, и replicas на других brokers. Реплики, успевающие за лидером, входят в ISR (in-sync replicas). При отказе лидера новый выбирается из ISR. replication.factor=3 означает три копии, но сам по себе ещё не гарантирует безопасную запись.
| Настройка producer | Когда запись считается успешной | Компромисс |
|---|---|---|
acks=0 | Producer не ждёт broker | Максимум скорости, возможна незаметная потеря |
acks=1 | Leader записал локально | Leader может умереть до репликации |
acks=all | Подтвердили все требуемые ISR | Надёжнее, но latency/availability зависят от ISR |
Типовая надёжная связка: replication factor 3, acks=all, осмысленный min.insync.replicas (часто 2) и idempotent producer. Тогда кластер предпочитает временно отказать в записи, а не подтвердить запись, оставшуюся в одной копии. Это осознанный выбор consistency против availability.
Idempotent producer не даёт повторной отправке породить дубль в partition в рамках своей сессии. Kafka transactions позволяют атомарно прочитать, записать в Kafka и зафиксировать offsets при Kafka-to-Kafka обработке. Но внешние side effects всё равно требуют идемпотентности или outbox.
Метаданные современного Kafka-кластера хранятся в кворуме KRaft controllers. В Kafka 4.x режима с ZooKeeper уже нет. Контроллеры лучше держать нечётным кворумом, обычно три; потеря большинства мешает изменениям метаданных и выборам лидеров.
$ kafka-topics.sh --bootstrap-server kafka:9092 \
--describe --topic orders
$ kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group billing
$ kafka-metadata-quorum.sh --bootstrap-server kafka:9092 \
describe --status
§5RabbitMQ: exchange, binding, queue и ack
Producer в AMQP обычно публикует не прямо в очередь, а в exchange. Exchange по routing key и bindings решает, в какие очереди положить сообщение. Consumers читают очереди и конкурируют за задания.
Consumer acknowledgement говорит брокеру: «я успешно обработал сообщение, его можно удалить». При manual ack падение consumer до подтверждения приводит к redelivery. При auto-ack брокер считает доставку успехом сразу — быстро, но сообщение потеряется, если процесс упадёт во время работы.
Prefetch ограничивает число доставленных, но ещё не подтверждённых сообщений. Слишком большой prefetch даёт одному consumer забрать много работы и памяти, слишком маленький недогружает быстрый consumer. Начинают с небольшого значения и измеряют throughput/latency.
Спросят
«Publisher confirm и consumer ack — это одно?» — Нет. Publisher confirm подтверждает producer, что broker принял публикацию. Consumer ack подтверждает broker, что consumer завершил обработку. Между этими событиями сообщение может долго лежать в очереди.
§6Надёжный RabbitMQ: quorum, retry и DLQ
Для отказоустойчивой replicated queue в современных версиях используют quorum queue: она реплицируется через Raft и подтверждает безопасную публикацию после записи большинством. Типовой размер — три узла; нечётное число сохраняет полезный кворум. Старые mirrored classic queues удалены из RabbitMQ 4.0.
Для сохранности нужны не только durable queue и persistent message, но и publisher confirms: иначе producer не знает, успел ли broker принять публикацию. Репликация не заменяет резервные копии, а кластер, где все узлы стоят на одном гипервизоре, не защищает от его отказа.
Ошибка обработки не должна создавать бесконечный цикл nack(requeue=true). Делают ограниченное число попыток с задержкой через retry queue/TTL, затем отправляют сообщение в dead-letter exchange/queue (DLX/DLQ). В сообщение кладут причину, число попыток и correlation ID; DLQ мониторят и разбирают, а не используют как кладбище.
Строгий порядок легко нарушается несколькими consumers и повторной доставкой: следующее сообщение может завершиться раньше предыдущего. Если порядок критичен, применяют partitioning по ключу, один consumer или single active consumer — ценой параллелизма.
§7Redis: не только кэш
Redis хранит рабочий набор в памяти и даёт атомарные команды над структурами: strings, hashes, lists, sets, sorted sets, streams. Примеры: INCR для счётчика, SET ... EX для TTL, sorted set для рейтинга, hash для компактного объекта. Большинство команд выполняются последовательно в основном event loop, поэтому долгий O(N) запрос или тяжёлый Lua script задерживает остальных клиентов.
Типовой паттерн кэша — cache-aside: приложение читает Redis; при miss читает БД, кладёт значение с TTL и возвращает его. При изменении сначала обновляет источник истины, затем инвалидирует кэш. Нужно заранее решить, допустимы ли устаревшие данные.
- Stampede: популярный ключ истёк, сотни запросов одновременно пошли в БД. Помогают request coalescing/короткая блокировка и раннее обновление.
- Avalanche: множество ключей истекает одновременно. Добавляют случайный jitter к TTL.
- Penetration: постоянно спрашивают несуществующие ключи. Кратко кэшируют отрицательный результат и валидируют вход.
- Eviction: при
maxmemoryполитика выбирает, какие ключи удалить. Политику иevicted_keysнадо контролировать, иначе Redis сам начнёт менять ваш набор данных.
Грабли
KEYS * обходит всё пространство ключей и может надолго заблокировать рабочий инстанс. Для итерации используют SCAN, но и его нагрузку измеряют. Большие keys и коллекции опасны ещё памятью, сетевым ответом, удалением и временем репликации.
§8Redis: RDB, AOF, Sentinel и Cluster
| Механизм | Что делает | Чего не делает |
|---|---|---|
| RDB | Периодический компактный snapshot, быстрый restore | Между snapshots можно потерять последние изменения |
| AOF | Пишет операции в append-only log; fsync задаёт компромисс | При everysec возможна потеря примерно последней секунды |
| Replication | Асинхронно копирует данные с primary на replicas | Не является backup; недавняя запись при failover может потеряться |
| Sentinel | Мониторинг, обнаружение отказа и автоматический failover primary/replica | Не шардирует данные |
| Cluster | Шардирует 16 384 hash slots и даёт HA каждого shard | Сложнее multi-key операции: ключи должны быть в одном slot |
Persistence, replication и backup — три разные защиты. RDB/AOF помогают пережить перезапуск, replica — отказ узла, а независимый проверенный backup — удаление, повреждение или потерю всего кластера. Асинхронная репликация Redis означает, что failover не обещает нулевой RPO; команда WAIT снижает риск, но не превращает систему в строго синхронную БД.
В Redis Cluster ключ отображается на один из 16 384 slots. Multi-key command работает, только если ключи на одном shard; hash tag в фигурных скобках, например cart:{42} и user:{42}, принудительно помещает их в один slot. Не стоит собирать всё под один tag — получится hotspot и исчезнет смысл шардирования.
§9Pub/Sub, Streams, транзакции и locks в Redis
Pub/Sub — fire-and-forget: подписчик получает только сообщения, опубликованные пока он подключён. Нет сохранения истории, replay и consumer offsets. Подходит для эфемерных уведомлений, но не для критичного задания.
Redis Streams сохраняет записи с ID, поддерживает consumer groups, pending entries и XACK. Это уже очередь/журнал для умеренной нагрузки, но её отказоустойчивость наследует компромиссы Redis persistence и async replication. Если нужны длительный retention, большой поток и много независимых consumers с replay, Kafka обычно естественнее.
MULTI/EXEC выполняет пакет команд последовательно без вмешательства других клиентов, WATCH даёт optimistic locking. Классического rollback как в SQL нет: ошибка одной команды не откатывает уже выполненные. Lua script атомарен относительно других команд, но длинный script блокирует event loop.
Минимальная распределённая блокировка: SET lock value NX PX 10000, где value — уникальный token владельца. Освобождать нужно атомарным compare-and-delete Lua script: удалить, только если token совпадает. TTL обязателен на случай смерти владельца. Для операций, где просрочившийся владелец способен испортить данные, нужен fencing token — монотонный номер, по которому защищаемый ресурс отвергнет старого владельца. Redis-lock сам по себе не заменяет транзакционную гарантию бизнес-системы.
§10Что мониторить и как искать проблему
| Система | Сначала смотреть | Типовые причины |
|---|---|---|
| Kafka | consumer lag, under-replicated/offline partitions, ISR changes, disk, request latency, controller quorum | медленный consumer, rebalance storm, заполненный диск, сеть/GC, мало ISR |
| RabbitMQ | messages ready/unacked, publish/deliver rates, consumers, memory/disk alarms, quorum health | consumer умер/завис, мало prefetch или ресурсов, poison message, backpressure |
| Redis | hit ratio, used memory, evictions, latency/slowlog, blocked clients, replication link/lag, persistence status | big key, O(N) команда, swap, fork/CoW, сетевой лимит, исчерпана память |
Для любой системы сначала отвечайте на четыре вопроса: producer публикует или уже получает ошибки; данные накопились на broker; consumers живы и подтверждают; отказоустойчивый кворум и диски здоровы. Не лечите lag добавлением consumers, пока не проверили число Kafka partitions, внешнюю БД и время обработки одного сообщения.
# RabbitMQ
$ rabbitmqctl list_queues name messages_ready messages_unacknowledged consumers
# Redis
$ redis-cli INFO
$ redis-cli SLOWLOG GET 10
$ redis-cli MEMORY USAGE some:key
§11Что спрашивают на собеседовании
Kafka — очередь или база? Это распределённый сохраняемый log: consumer двигает свой offset, а чтение не удаляет запись. Поэтому доступны replay и несколько независимых групп.
Где Kafka гарантирует порядок? Только внутри partition. Для порядка одной сущности выбирают её ID ключом; глобальный порядок уменьшает параллелизм до одной partition.
Почему consumers больше partitions не ускоряют группу? Одну partition одновременно читает один consumer группы, поэтому лишним consumers нечего назначить.
Что такое lag? Насколько committed offset consumer group отстаёт от конца partition. Большой lag — симптом; надо проверить скорость producer/consumer, ошибки, rebalance и downstream.
RabbitMQ: ready и unacked? Ready ещё ждут выдачи. Unacked уже доставлены consumers, но подтверждения нет. Рост unacked часто означает зависшую/медленную обработку или завышенный prefetch.
Что произойдёт при смерти Rabbit consumer до ack? Сообщение будет повторно доставлено, поэтому обработка должна быть идемпотентной. При auto-ack оно могло потеряться.
Зачем quorum queue? Чтобы очередь переживала отказ узла через репликацию большинством. Нужны минимум три разумно размещённых узла и publisher confirms; кворум уменьшает доступность записи при потере большинства.
RDB или AOF? RDB компактнее и быстрее восстанавливается, но теряет изменения между snapshots. AOF пишет операции и при everysec обычно ограничивает потерю примерно секундой, ценой большего диска/I/O. Часто сочетают оба и отдельно делают backups.
Sentinel или Cluster? Sentinel делает failover одного набора primary/replicas, но не шардирует. Cluster распределяет slots между shards и даёт HA, зато усложняет multi-key операции.
Redis однопоточный? Основное выполнение команд сериализовано, что упрощает атомарность; сетевой I/O и фоновые задачи могут использовать другие threads/processes. Практический вывод: тяжёлая команда блокирует latency остальных.
Готовый короткий ответ
«Kafka выбираю, когда нужен долговечный поток событий: retention, replay, большой throughput и несколько независимых consumer groups. Масштабирование идёт partitions, порядок есть внутри partition, а типовая доставка — at-least-once с идемпотентным consumer.
RabbitMQ беру для рабочих очередей и команд, где важны exchange/routing, manual ack, prefetch, retries и DLQ. Для HA использую quorum queues и не путаю publisher confirms с consumer acknowledgements.
Redis использую прежде всего как быстрое хранилище: кэш, TTL, счётчики, сессии и rate limiting. Заранее выбираю eviction и persistence, понимаю, что replication асинхронная; Sentinel даёт failover, Cluster — шардирование. Pub/Sub не хранит сообщения, Streams хранит, но это не автоматически замена Kafka. Во всех трёх системах сначала формулирую допустимые потери, дубли, порядок и поведение при сетевом разделении.»