Брокеры сообщений: модель, гарантии доставки, партиции и консьюмер-группы
Кратко о теме
Заголовок раздела «Кратко о теме»Брокер сообщений — это посредник, который принимает сообщение от отправителя, надёжно его сохраняет и отдаёт получателям, разрывая прямую связь между сервисами. Ценность именно в разрыве: продюсер и консьюмер не обязаны быть живыми одновременно, не знают адресов друг друга, могут масштабироваться и релизиться независимо, а всплеск нагрузки поглощается очередью, а не падает на слабое звено. Платой идут потеря немедленной согласованности, дубликаты, ограниченные гарантии порядка и заметно более сложная отладка. Почти любой вопрос из этого блока раскладывается по трём осям: модель доставки (очередь с конкурирующими потребителями против publish/subscribe), гарантии (at-most-once / at-least-once / exactly-once, порядок, durability) и масштабирование потребления (партиции, консьюмер-группы, ребалансировки, лаг).
Полезно держать в голове два разных архетипа брокера, потому что интервьюер почти всегда неявно спрашивает про один из них. Первый — очередь-с-удалением (RabbitMQ, SQS, ActiveMQ): сообщение живёт в очереди, брокер отслеживает состояние каждого сообщения, после подтверждения (ack) оно удаляется. Отсюда богатая маршрутизация, per-message TTL, retry с задержкой, штатный dead-letter, но масштабирование упирается в саму очередь, а порядок ломается при нескольких конкурирующих потребителях. Второй — распределённый лог (Kafka, Pulsar, Redis Streams, отчасти NATS JetStream): сообщения дописываются в конец неизменяемого лога, читатель хранит свою позицию (offset), а брокер удаляет данные по retention, а не по факту прочтения. Отсюда воспроизведение истории, независимые группы читателей, линейное масштабирование через партиции — но нет per-message ack и нет из коробки отложенных ретраев.
В логовой модели вся арифметика параллелизма определяется одной формулой: партиция — единица параллелизма и единица порядка. Внутри партиции сообщения строго упорядочены; между партициями порядка нет. Одну партицию в рамках одной консьюмер-группы читает ровно один консьюмер, поэтому число партиций — это потолок полезного параллелизма группы. Лишние консьюмеры простаивают. Разные группы независимы: каждая имеет собственный набор офсетов и получает полную копию потока, поэтому «одна группа успела, другая потеряет» — невозможно, пока сообщение не удалено по retention.
Гарантии доставки в реальном мире почти всегда сводятся к at-least-once плюс идемпотентность на стороне обработчика. «Exactly once» существует, но означает не «сообщение физически передано по сети один раз», а «эффект обработки применён один раз» — и достигается либо транзакцией внутри одной системы (Kafka read-process-write), либо дедупликацией по ключу в приёмнике. Как только в цепочке появляется внешний побочный эффект (списание денег, отправка письма, вызов чужого HTTP API), никакой брокер за вас exactly-once не сделает — нужен ключ идемпотентности.
Вопросы и ответы
Заголовок раздела «Вопросы и ответы»Использовал ли брокеры сообщений? Если дa то какие?
Заголовок раздела «Использовал ли брокеры сообщений? Если дa то какие?»Коротко. Вопрос на личный опыт: назовите 1–3 брокера, которые реально трогали руками, и сразу привяжите каждый к задаче — «Kafka для событий домена и аналитики, RabbitMQ для отложенных задач с ретраями, Redis Streams для лёгкой внутренней очереди». Не перечисляйте всё, что слышали.
Глубже. Хороший ответ строится по схеме «брокер → задача → почему именно он → на что напоролись». Например: «Kafka, топик order-events, 12 партиций, ключ — order_id, потому что нужен порядок в рамках заказа и три независимых потребителя одного потока: биллинг, поиск и DWH. Напоролись на то, что при долгой обработке консьюмер выпадал по max.poll.interval.ms и группа уходила в бесконечную ребалансировку — вынесли тяжёлую работу в пул воркеров и уменьшили max.poll.records». Такой ответ показывает, что вы работали с продакшеном, а не читали обзор. В Go стек обычно такой: segmentio/kafka-go или twmb/franz-go/IBM/sarama для Kafka, rabbitmq/amqp091-go для RabbitMQ, nats-io/nats.go для NATS. Если опыта мало — честно скажите «в проде только RabbitMQ, Kafka изучал на пет-проекте», и переходите к тому, что знаете хорошо: врать про масштаб опасно, следующим вопросом будет «а сколько партиций и почему столько».
Для чего нужны брокеры? Что за механизм? Как в нем выстраиваются данные?
Заголовок раздела «Для чего нужны брокеры? Что за механизм? Как в нем выстраиваются данные?»Коротко. Брокер нужен для асинхронного, надёжного и слабосвязанного обмена: он развязывает отправителя и получателя во времени (получатель может лежать), в пространстве (никто не знает адресов) и по нагрузке (очередь сглаживает пики). Механически это сервер с durable-хранилищем, который принимает сообщения, persist’ит их и отдаёт подписчикам, отслеживая, что доставлено.
Глубже. Данные внутри выстраиваются по-разному в зависимости от модели. В RabbitMQ продюсер публикует в exchange с routing key; exchange по правилам биндингов раскладывает сообщение по очередям (direct — точное совпадение ключа, topic — по маске order.*.created, fanout — во все, headers — по заголовкам); очередь — это FIFO-структура, сообщение живёт в ней до ack и после подтверждения удаляется. В Kafka продюсер пишет в топик, топик физически нарезан на партиции, партиция — это append-only лог на диске, разбитый на сегменты, где у каждой записи есть монотонный offset. Запись всегда идёт в конец, чтение — последовательное с произвольной позиции; удаление — не по факту прочтения, а по retention.ms/retention.bytes либо через компакцию по ключу (cleanup.policy=compact, остаётся последнее значение на ключ). Именно это даёт Kafka её пропускную способность: последовательный I/O, страничный кеш ОС и zero-copy при отдаче. Порядок гарантируется только внутри партиции, поэтому распределение по партициям (обычно hash(key) % partitions) — это архитектурное решение, а не деталь.
Что подразумевает под собой гарантия доставки “at least once”?
Заголовок раздела «Что подразумевает под собой гарантия доставки “at least once”?»Коротко. At-least-once означает: сообщение точно не будет потеряно, но может быть доставлено и обработано более одного раза. Практически это достигается тем, что отправитель ретраит до получения подтверждения, а получатель подтверждает (коммитит офсет / шлёт ack) после успешной обработки.
Глубже. Дубликаты возникают в двух местах. На стороне продюсера: брокер записал сообщение, но ответ ack потерялся в сети — клиент по таймауту повторяет отправку, в логе оказываются две записи. На стороне консьюмера: обработали, но не успели закоммитить офсет (упал процесс, случилась ребалансировка, истёк max.poll.interval.ms) — после перезапуска чтение начинается с последнего закоммиченного офсета и часть сообщений переигрывается. Второй источник неустраним в принципе, потому что «применить эффект» и «зафиксировать офсет» — это две разные системы, и атомарно связать их можно только транзакцией, охватывающей обе.
Отсюда практический вывод, который и хотят услышать: at-least-once — это дефолт и правильный выбор, а бороться надо не с дубликатами в транспорте, а с неидемпотентностью в обработчике. Порядок действий в консьюмере при at-least-once строго такой: прочитать → обработать → закоммитить. Если поменять местами (коммит до обработки), получится at-most-once с потерями. В Kafka это enable.auto.commit=false и ручной CommitMessages/CommitOffsets; в RabbitMQ — autoAck=false и Ack после работы.
Что делать, если у нас “at least once”, но при этом бизнес-логика чувствительна к дубликатам запросов?
Заголовок раздела «Что делать, если у нас “at least once”, но при этом бизнес-логика чувствительна к дубликатам запросов?»Коротко. Сделать обработчик идемпотентным: завести ключ идемпотентности (id сообщения или бизнес-ключ) и в той же транзакции, что и полезный эффект, записывать факт обработки — повторное сообщение с тем же ключом отбрасывается или возвращает прежний результат.
Глубже. Работающих техник несколько, и лучше назвать их по убыванию надёжности. Первая — таблица обработанных сообщений в той же БД, что и данные: INSERT INTO processed(msg_id) ... ON CONFLICT DO NOTHING внутри той же транзакции, где меняется состояние; если вставка не произошла, откатываемся и просто коммитим офсет. Атомарность обеспечивает СУБД, брокер тут ни при чём. Вторая — естественная идемпотентность: вместо balance = balance - 100 писать UPSERT конечного состояния или использовать версионирование (UPDATE ... WHERE version = $expected), тогда повтор просто не проходит. Третья — дедуп во внешнем сторе вроде Redis SET key val NX EX ttl; быстро, но это не транзакция, между дедупом и эффектом остаётся окно, и TTL надо брать заведомо больше максимального времени переигрывания.
Важные детали, на которых заваливаются: ключ идемпотентности должен генерировать продюсер и класть в сообщение (заголовок Idempotency-Key/message_id), а не консьюмер вычислять хеш тела — тело может отличаться пробелами или временной меткой. TTL записи о дедупе должен покрывать worst case: retention топика, если вы допускаете полное переигрывание, — иначе после reset офсетов дубли пройдут. И если побочный эффект уходит во внешнюю систему (платёжный шлюз, отправка SMS), идемпотентность нужно требовать от неё же — почти все платёжные API принимают ключ идемпотентности именно для этого.
// пример: обработка в одной транзакции с дедупомfunc handle(ctx context.Context, db *sql.DB, msgID string, amount int64, acc string) error { tx, err := db.BeginTx(ctx, nil) if err != nil { return err } defer tx.Rollback()
res, err := tx.ExecContext(ctx, `INSERT INTO processed_messages(msg_id) VALUES ($1) ON CONFLICT DO NOTHING`, msgID) if err != nil { return err } if n, _ := res.RowsAffected(); n == 0 { return nil // дубликат: эффект уже применён, просто коммитим офсет выше }
if _, err := tx.ExecContext(ctx, `UPDATE accounts SET balance = balance - $1 WHERE id = $2`, amount, acc); err != nil { return err } return tx.Commit()}Поддерживает ли из коробки exactly once? Как реализовать?
Заголовок раздела «Поддерживает ли из коробки exactly once? Как реализовать?»Коротко. Kafka поддерживает exactly-once semantics, но только для замкнутого цикла «прочитал из Kafka → обработал → записал в Kafka», через идемпотентный продюсер плюс транзакции. Как только результат уходит во внешнюю систему, EOS из коробки не существует ни у одного брокера — нужна идемпотентность приёмника. RabbitMQ и SQS exactly-once в общем виде не дают (у SQS FIFO есть дедупликация в окне 5 минут, что близко, но не то же самое).
Глубже. Механика Kafka. Идемпотентный продюсер (enable.idempotence=true, включён по умолчанию начиная с Kafka 3.0) присваивает клиенту PID и нумерует записи sequence number в рамках пары (PID, партиция); брокер отбрасывает повторы — это убирает дубликаты от ретраев продюсера, но только в пределах жизни сессии продюсера. Транзакции (KIP-98) добавляют transactional.id, который переживает перезапуск: продюсер открывает транзакцию, пишет и результат, и офсеты потреблённых сообщений (SendOffsetsToTransaction в тот же транзакционный контекст), затем коммитит; читатели с isolation.level=read_committed не видят незакоммиченных записей. Так офсет и результат фиксируются атомарно — это и есть EOS. В Kafka Streams всё это включается одной настройкой processing.guarantee=exactly_once_v2.
Цена: транзакции добавляют латентность (двухфазный коммит с транзакционным координатором), понижают пропускную способность, требуют аккуратного управления transactional.id при масштабировании и делают чтение чуть сложнее (LSO — last stable offset — сдерживает read_committed-читателей до завершения транзакции). Поэтому в большинстве продакшн-систем осознанно выбирают at-least-once + идемпотентный обработчик: проще, дешевле и устойчивее. Альтернатива «EOS без транзакций брокера» — transactional outbox: сервис в одной БД-транзакции меняет состояние и пишет строку в таблицу outbox, а отдельный релей (или Debezium через CDC) публикует её в брокер с at-least-once; дубликаты снимаются идемпотентностью на приёмнике.
Как будет работать комбинация «консьюмеров больше, чем партиций»?
Заголовок раздела «Как будет работать комбинация «консьюмеров больше, чем партиций»?»Коротко. В рамках одной консьюмер-группы лишние консьюмеры не получат ни одной партиции и будут простаивать: партиция назначается ровно одному консьюмеру группы. При 3 партициях и 5 консьюмерах трое работают, двое идут в холостую — но остаются горячим резервом и подхватят партиции при ребалансировке.
Глубже. Это прямое следствие модели: партиция — единица параллелизма, и Kafka не разрешает двум членам одной группы читать одну партицию, иначе рассыпались бы порядок и однозначность офсета. Простаивающие консьюмеры не бесполезны — они шлют heartbeat’ы и при падении активного участника в течение session.timeout.ms (по умолчанию 45 s с Kafka 3.0) получат его партиции без ожидания рестарта пода. Иногда так и делают специально: держат N+1 инстансов ради быстрого failover.
Если параллелизм действительно упёрся, есть три выхода: увеличить число партиций (kafka-topics.sh --alter --partitions, но помните, что при hash-партиционировании это ломает соответствие «ключ → партиция» для новых сообщений и порядок по ключу на границе изменения; уменьшить число партиций нельзя), распараллелить обработку внутри консьюмера пулом горутин с сохранением порядка по ключу и аккуратным коммитом офсетов только до «низкой воды», либо перейти на модель, где несколько потребителей делят одну партицию. Последнее в Kafka появилось относительно недавно: share groups (KIP-932, «Queues for Kafka») — в Kafka 4.0 в статусе preview/early access, там несколько консьюмеров могут читать одну партицию с per-message ack, ценой отказа от гарантий порядка. В RabbitMQ такой проблемы нет вовсе: очередь — competing consumers, сколько подписчиков, столько параллелизма, но и порядок при этом теряется.
Какие есть гарантии доставки?
Заголовок раздела «Какие есть гарантии доставки?»Коротко. Три уровня: at-most-once — не более одного раза, возможны потери, дубликатов нет; at-least-once — не менее одного раза, потерь нет, возможны дубликаты; exactly-once — ровно один эффект, достигается транзакцией или дедупликацией и в общем распределённом случае не бесплатен и не всегда возможен.
Глубже. Уровень определяется не «настройкой брокера», а совокупностью трёх решений. (1) Продюсер: acks=0 — fire-and-forget, at-most-once; acks=1 — ждём лидера, теряем при падении лидера до репликации; acks=all вместе с min.insync.replicas=2 и replication.factor=3 — запись подтверждена только после попадания в кворум ISR, потерь нет. (2) Ретраи и идемпотентность продюсера: без enable.idempotence ретраи дают дубликаты и могут переставить сообщения при max.in.flight.requests.per.connection > 1. (3) Консьюмер: коммит офсета до обработки → at-most-once, после обработки → at-least-once, в одной транзакции с результатом → exactly-once.
Отдельно стоит проговорить, что «гарантии доставки» — не единственные гарантии, о которых спрашивают. Есть ещё durability (сообщение переживает рестарт брокера: flush, репликация, unclean.leader.election.enable=false), порядок (в Kafka — только внутри партиции; в RabbitMQ — только при одном потребителе и без ретраев в хвост очереди) и доставка всем подписчикам (fanout). На собеседовании выигрышно сказать: «на транспортном уровне мы выбираем at-least-once с acks=all, а exactly-once обеспечиваем на прикладном — это дешевле и надёжнее, чем распределённые транзакции».
Есть несколько консьюмеров одного и того же топика, два-три разных сервиса. Каким образом каждый из этих консьюмер-групп будет получать сообщение? Не получится ли так, что первая группа получила некоторое сообщение, а вторая, не успев, потеряет?
Заголовок раздела «Есть несколько консьюмеров одного и того же топика, два-три разных сервиса. Каким образом каждый из этих консьюмер-групп будет получать сообщение? Не получится ли так, что первая группа получила некоторое сообщение, а вторая, не успев, потеряет?»Коротко. Нет, не получится. Каждая консьюмер-группа хранит собственный набор офсетов и читает лог независимо: одно и то же сообщение получат все группы. Чтение одной группой ничего не удаляет — данные живут в партиции до истечения retention, а не до факта прочтения.
Глубже. Офсеты групп хранятся во внутреннем компактящемся топике __consumer_offsets (по умолчанию 50 партиций); ключ — тройка (group.id, topic, partition), и координатор группы выбирается как hash(group.id) % 50. Именно поэтому три сервиса с разными group.id — это три полностью изолированных курсора по одному и тому же логу: биллинг может отставать на час, аналитика — читать в реальном времени, а третий сервис вообще перезапуститься с начала (--reset-offsets --to-earliest), и никто друг другу не помешает.
Единственный реальный способ «потерять» сообщение для отстающей группы — это retention: если группа отстала сильнее, чем retention.ms/retention.bytes, старые сегменты будут удалены, консьюмер получит OffsetOutOfRange и, в зависимости от auto.offset.reset, прыгнет на earliest (пропустив дыру) или на latest (пропустив всё накопленное). Поэтому мониторинг лага плюс retention с запасом — обязательны. Второй, более коварный, случай — совпадение group.id у разных сервисов по недосмотру: тогда это не две группы, а одна, партиции поделятся между ними, и каждый сервис увидит лишь часть потока. Это классическая ошибка в конфигах, и её стоит упомянуть. В RabbitMQ аналогичную семантику даёт fanout/topic exchange с отдельной очередью на каждого потребителя: у каждого сервиса своя очередь, копия сообщения кладётся в каждую.
Как консьюмер-лаг соотносится с консьюмер-группами?
Заголовок раздела «Как консьюмер-лаг соотносится с консьюмер-группами?»Коротко. Лаг считается для пары (группа, партиция) как log_end_offset - committed_offset, то есть сколько сообщений в партиции ещё не подтверждено этой группой. Лаг группы — сумма по её партициям; у разных групп на одном топике лаг совершенно разный и независимый.
Глубже. Смотреть лаг удобно через kafka-consumer-groups.sh --bootstrap-server ... --describe --group billing: он покажет по каждой партиции CURRENT-OFFSET, LOG-END-OFFSET, LAG и хост консьюмера. В продакшене это снимают экспортерами (kafka_exporter, Burrow) или клиентской метрикой JMX records-lag-max. Полезно понимать нюансы: лаг в сообщениях плохо переводится в лаг во времени (одно сообщение может обрабатываться 1 мс, другое 2 с), поэтому важнее алертить на производную — растёт лаг или разгребается, — и на «time lag» (разница между временем последнего обработанного сообщения и now). Второй нюанс: лаг считается по закоммиченному офсету, поэтому при редких коммитах (раз в 30 s) вы увидите пилу, не отражающую реальную обработку. Третий: перекос лага по партициям почти всегда означает горячий ключ — весь трафик валится в одну партицию из-за неудачного key.
Связь с группами практическая: если лаг растёт равномерно по всем партициям — не хватает мощности, добавляйте консьюмеров (но не больше числа партиций) или ускоряйте обработку. Если лаг растёт скачками с провалами до нуля — вероятны ребалансировки: во время eager-ребалансировки вся группа останавливает потребление. Если лаг растёт только у одной группы — проблема в её сервисе, а не в брокере.
Какая разница Pub/Sub и Message Queue?
Заголовок раздела «Какая разница Pub/Sub и Message Queue?»Коротко. Message Queue — модель «один-к-одному»: сообщение забирает ровно один из конкурирующих потребителей, это распределение работы. Pub/Sub — «один-ко-многим»: каждое сообщение получают все подписчики, это распространение события.
Глубже. Разница не только в кардинальности, но и в семантике: очередь передаёт команду/задачу («обработай платёж»), у неё есть адресат и ожидание, что работа будет выполнена один раз; pub/sub передаёт факт («платёж проведён»), отправитель не знает и не хочет знать, кто на него подпишется, и добавление нового подписчика не требует изменений в отправителе. Отсюда следствие для связанности: очередь тянет за собой знание о получателе, событийная шина — нет.
Технически многие брокеры дают обе модели. RabbitMQ: очередь с несколькими подписчиками — это MQ, а fanout-exchange с отдельной очередью на подписчика — это pub/sub. Kafka объединяет их элегантно: внутри одной группы — competing consumers (очередь), между группами — fanout (pub/sub), и это, пожалуй, лучший однострочный ответ на вопрос. Классический pub/sub без персистентности (JMS topic, NATS core, Redis Pub/Sub) имеет важное отличие: подписчик, которого не было в момент публикации, сообщение не получит вовсе; Kafka же за счёт лога позволяет подписчику появиться позже и вычитать историю. Не путайте также с pull/push: Kafka — pull-модель (консьюмер сам делает fetch), RabbitMQ по умолчанию push с prefetch-окном.
С какими базами данных, языками, технологиями, очередями сообщений вы работали?
Заголовок раздела «С какими базами данных, языками, технологиями, очередями сообщений вы работали?»Коротко. Обзорный вопрос на профиль. Отвечайте структурировано, группами, и на каждой группе называйте не список, а глубину: где вы владелец решения, где просто пользователь. Не перечисляйте больше 3–4 пунктов в группе.
Глубже. Рабочий каркас ответа: «Основной язык — Go, ~N лет; до этого <язык>. Хранилища: PostgreSQL как основное — писал миграции, разбирался с планами и блокировками; Redis для кеша и распределённых локов; ClickHouse — только читал аналитические витрины. Брокеры: Kafka в проде, RabbitMQ в предыдущем проекте. Инфраструктура: Docker, Kubernetes, Prometheus/Grafana, CI на GitLab». Дальше добавьте один якорь — конкретный кейс, который вы готовы разворачивать: «самое интересное было, когда мы переносили обмен между сервисами с синхронного gRPC на события в Kafka и разбирались с outbox». Интервьюер почти всегда цепляется за такой якорь, и это выгодно: вы заводите разговор на территорию, где сильны.
Две типичные ошибки. Первая — перечислить двадцать технологий: вас спросят про самую экзотическую, и ответить будет нечего. Вторая — не проградуировать уровень: скажите прямо «с Cassandra работал ограниченно, только через готовый репозиторий», это воспринимается как зрелость, а не как слабость.
Что такое консьюмер группа?
Заголовок раздела «Что такое консьюмер группа?»Коротко. Консьюмер-группа — это набор консьюмеров с одинаковым group.id, которые совместно потребляют топик: брокер делит партиции между членами группы так, что каждая партиция достаётся ровно одному члену, а офсеты хранятся общие для всей группы.
Глубже. Группа — единица масштабирования и единица позиции чтения одновременно. Ею управляет group coordinator — брокер, выбранный по hash(group.id) % количество партиций __consumer_offsets. Члены группы шлют ему heartbeat’ы каждые heartbeat.interval.ms (3 s по умолчанию); если heartbeat пропал дольше session.timeout.ms (45 s), или консьюмер не вызвал poll дольше max.poll.interval.ms (5 min), участник считается мёртвым и запускается ребалансировка — перераспределение партиций.
Стратегия назначения задаётся partition.assignment.strategy: RangeAssignor (диапазонами по каждому топику, склонен к перекосу при нескольких топиках), RoundRobinAssignor (равномернее), StickyAssignor и CooperativeStickyAssignor (стараются сохранить прежние назначения, чтобы минимизировать перетасовку). Начиная с Kafka 2.4 кооперативный протокол (KIP-429) делает ребалансировку инкрементальной — вместо stop-the-world отзываются только те партиции, которые реально переезжают. В Kafka 4.0 по умолчанию используется новый серверный протокол групп (KIP-848), где логика назначения перенесена на координатора, что заметно сокращает паузы. Ещё одна важная настройка — static membership (group.instance.id, KIP-345): участник с постоянным id при рестарте не вызывает ребалансировку, если уложился в session.timeout.ms; очень полезно для деплоев в Kubernetes.
Было две партиции и два консьюмера, потом один консьюмер умер. Что произойдет и что можно сделать?
Заголовок раздела «Было две партиции и два консьюмера, потом один консьюмер умер. Что произойдет и что можно сделать?»Коротко. Координатор обнаружит смерть по пропаже heartbeat через session.timeout.ms, запустит ребалансировку, и выживший консьюмер получит обе партиции. Данные не теряются, но пропускная способность падает вдвое, лаг растёт, а сообщения, обработанные умершим консьюмером без коммита офсета, будут переиграны — то есть возможны дубликаты.
Глубже. По шагам: (1) умерший перестал слать heartbeat; (2) через session.timeout.ms координатор помечает его выбывшим и объявляет ребалансировку; (3) при eager-протоколе выживший тоже на короткое время отдаёт свои партиции и потребление в группе полностью встаёт, при CooperativeStickyAssignor — только добирает осиротевшую партицию без остановки; (4) чтение осиротевшей партиции возобновляется с последнего закоммиченного офсета, поэтому необработанного не потеряется, а частично обработанное повторится.
Что делать. Оперативно — убедиться, что выживший справляется, иначе поднять новый инстанс (оркестратор обычно делает это сам). Системно: держать реплики консьюмера с запасом (N+1, лишний просто простаивает и подхватит быстро), использовать CooperativeStickyAssignor, чтобы падение одного не останавливало всю группу, включить static membership, чтобы плановые рестарты вообще не вызывали ребалансировку, снизить session.timeout.ms для более быстрого обнаружения (не забыв про соотношение heartbeat.interval.ms ≈ session.timeout.ms / 3), и обязательно сделать обработчик идемпотентным — иначе переигранные сообщения дадут двойные эффекты. И заранее заложить число партиций больше числа консьюмеров, чтобы отказ одного давал деградацию, а не потерю половины мощности.
Максимальное количество консьюмеров для партиции в одной группе?
Заголовок раздела «Максимальное количество консьюмеров для партиции в одной группе?»Коротко. Один. Партиция в рамках одной консьюмер-группы назначается ровно одному консьюмеру; больше одного одновременного читателя от одной группы у партиции быть не может. При этом одна и та же партиция может параллельно читаться сколь угодно многими консьюмерами из разных групп.
Глубже. Ограничение вытекает из того, что позиция чтения (офсет) хранится одна на пару (группа, партиция): два читателя от одной группы конкурировали бы за один курсор, порядок и коммиты потеряли бы смысл. Обходов ровно два. Первый — параллелить внутри одного консьюмера: читать батч, раскидывать по горутинам с шардированием по ключу (чтобы порядок в рамках ключа сохранился) и коммитить офсет только до непрерывно завершённого префикса. Второй — новые share groups (KIP-932) в Kafka 4.0 (preview): там модель ближе к классической очереди, несколько консьюмеров читают одну партицию с per-message acknowledgement, но гарантия порядка при этом не даётся. В RabbitMQ вопрос не стоит: очередь по определению допускает произвольное число конкурирующих потребителей, распределение идёт round-robin с учётом prefetch.
Что такое брокер сообщений?
Заголовок раздела «Что такое брокер сообщений?»Коротко. Брокер сообщений — промежуточный сервис (обычно кластер), который принимает сообщения от продюсеров, надёжно их хранит, маршрутизирует и доставляет консьюмерам, беря на себя буферизацию, персистентность, ретраи и учёт доставленного.
Глубже. С точки зрения архитектуры брокер — это реализация паттерна Mediator/Message Broker из Enterprise Integration Patterns: вместо N×M прямых интеграций между сервисами получается N+M связей с шиной, и каждая сторона знает только контракт сообщения. Технически брокер обычно предоставляет: durable-хранилище (лог или очередь на диске), репликацию между узлами для отказоустойчивости, модель адресации (топик/exchange/queue), протокол (собственный бинарный у Kafka, AMQP 0-9-1 у RabbitMQ, MQTT, STOMP, NATS), механизм подтверждений и учёт позиций читателей. Не путайте брокер с шиной ESB: ESB дополнительно берёт на себя трансформацию форматов и оркестрацию бизнес-логики, что сегодня считается антипаттерном («smart endpoints, dumb pipes»). Не путайте и с брокерless-подходом вроде ZeroMQ — там библиотека, а не сервер, и никакой персистентности нет.
Что можешь сказать про топик, партицию, продюсера, консьюмера, группу консьюмеров?
Заголовок раздела «Что можешь сказать про топик, партицию, продюсера, консьюмера, группу консьюмеров?»Коротко. Топик — именованный логический поток сообщений. Партиция — физический упорядоченный append-only лог, из которых состоит топик; единица параллелизма и порядка. Продюсер — тот, кто пишет, выбирая партицию по ключу. Консьюмер — тот, кто читает, продвигая свой офсет. Группа консьюмеров — набор консьюмеров с общим group.id, между которыми делятся партиции и у которых общие офсеты.
Глубже. Стоит добавить связки, которые показывают понимание, а не заучивание. У каждой партиции есть лидер и набор реплик; продюсеры и консьюмеры по умолчанию работают с лидером, а min.insync.replicas вместе с acks=all определяет durability. Записи в партиции адресуются монотонным offset, который никогда не переиспользуется. Продюсер выбирает партицию так: если партиция задана явно — она; если задан ключ — hash(key) % numPartitions (в Java-клиенте murmur2); если ключа нет — «липкий» партиционер (KIP-480, с 2.4; переработан в KIP-794 в 3.3), который набивает батч в одну партицию и переключается, что лучше для пропускной способности. Отсюда правило: ключ = сущность, порядок по которой вам важен (order_id, user_id). Консьюмер работает по pull-модели: fetch с указанием офсета, брокер отдаёт батч. Группа связывает всё вместе: количество полезных консьюмеров ≤ количество партиций, а разные группы читают независимо. Хорошее завершение ответа: «число партиций — это одновременно потолок параллелизма, гранулярность порядка и стоимость ребалансировки, поэтому его выбирают заранее и с запасом: увеличить можно, уменьшить — нет».
Что такое Dead Letter Queue?
Заголовок раздела «Что такое Dead Letter Queue?»Коротко. DLQ — отдельная очередь/топик, куда складываются сообщения, которые не удалось обработать: после исчерпания ретраев, при ошибке десериализации, при истечении TTL или переполнении очереди. Смысл — не блокировать основной поток «ядовитым» сообщением и не потерять его, а отложить для разбора.
Глубже. В RabbitMQ DLQ работает нативно: у очереди задаётся x-dead-letter-exchange (и опционально x-dead-letter-routing-key), и сообщение автоматически перекладывается туда при basic.reject/basic.nack с requeue=false, при истечении per-message или per-queue TTL, при превышении x-max-length и при превышении x-delivery-limit у quorum-очередей. Классическая связка «retry с задержкой» строится через промежуточную очередь с TTL и dead-letter-обратно в основную.
В Kafka DLQ из коробки нет — это прикладной паттерн: обработчик ловит ошибку, публикует сообщение в топик orders.DLQ вместе с диагностическими заголовками (исходный топик, партиция, офсет, текст ошибки, номер попытки, timestamp) и коммитит офсет исходного сообщения. Готовая реализация есть в Kafka Connect (errors.tolerance=all, errors.deadletterqueue.topic.name). Важные практики: DLQ обязательно мониторить (сообщение в DLQ = алерт, иначе очередь превращается в кладбище), хранить в заголовках всё для воспроизведения, предусмотреть процедуру reprocess (перелить из DLQ обратно после фикса), и отделять транзиентные ошибки (таймаут БД — ретраить с backoff) от перманентных (битый JSON, отсутствующая сущность — сразу в DLQ, ретраи бессмысленны и только заблокируют партицию). Ещё одна тонкость Kafka: пока вы ретраите одно сообщение, вся партиция стоит — поэтому либо ограничивайте число ретраев в консьюмере жёстким таймаутом, либо выносите ретраи в отдельные retry-топики с задержкой (подход «retry topics» из Spring Kafka).
Как вы тестируете асинхронные сценарии и обработку очередей сообщений?
Заголовок раздела «Как вы тестируете асинхронные сценарии и обработку очередей сообщений?»Коротко. Тремя слоями: юнит-тесты на чистый обработчик (брокер за интерфейсом, вход — распарсенное сообщение), интеграционные тесты на реальном брокере в Testcontainers (проверяем сериализацию, партиционирование, коммиты, DLQ), и отдельно — тесты на свойства асинхронности: идемпотентность, повторную доставку, порядок, поведение при отравленном сообщении.
Глубже. Практические приёмы, которые стоит назвать. Изоляция транспорта: в Go объявите на стороне приложения узкий интерфейс type Publisher interface { Publish(ctx, key string, payload []byte) error } и подменяйте его в юнит-тестах — тогда 90 % логики тестируется без брокера. Testcontainers (testcontainers-go с модулями kafka/rabbitmq/redpanda) поднимает настоящий брокер на время теста; Redpanda часто берут как более лёгкую Kafka-совместимую замену для CI. Детерминизм ожидания: никаких time.Sleep — используйте polling с дедлайном (require.Eventually) или синхронизацию через канал, который закрывает обработчик; иначе тест будет флакать. Тест на дубликат: отправьте одно и то же сообщение дважды и проверьте, что эффект применён один раз — это прямая проверка идемпотентности. Тест на переигрывание: сбросьте офсет группы на earliest и убедитесь, что состояние не разъехалось. Тест на poison message: подайте невалидный payload и проверьте, что он ушёл в DLQ, а партиция не встала. Контрактные тесты на схему сообщения (Avro/Protobuf + Schema Registry с проверкой совместимости, или Pact для message-контрактов) — они ловят самый частый прод-инцидент: продюсер поменял формат, консьюмер упал. Наконец, для сквозных сценариев полезна корреляция: trace_id в заголовках сообщения и OpenTelemetry-контекст, прокинутый через брокер, — иначе распределённый асинхронный флоу невозможно отладить ни в тесте, ни в проде.
Использовали ли вы message broker? Для каких сценариев?
Заголовок раздела «Использовали ли вы message broker? Для каких сценариев?»Коротко. См. выше про «Использовал ли брокеры сообщений» — здесь акцент смещён на сценарии, поэтому отвечать нужно списком задач, а не списком технологий: развязка сервисов через доменные события, вынос долгих операций из HTTP-запроса, сглаживание пиков, фан-аут одного события нескольким потребителям, интеграция с внешними системами через буфер.
Глубже. Убедительный набор сценариев из практики: (1) доменные события — сервис заказов публикует order.created, на него независимо подписаны биллинг, склад и поиск; добавление четвёртого подписчика не требует изменений в сервисе заказов; (2) офлайн-работа — генерация PDF, отправка email/пуша, пересчёт отчёта: HTTP-хендлер кладёт задачу и отвечает 202, воркеры разгребают; (3) сглаживание пика — импорт файла на миллион строк режется на сообщения, и БД получает ровную нагрузку вместо всплеска; (4) интеграция с медленным партнёром — очередь как буфер и точка ретраев, чтобы недоступность внешнего API не роняла наш сервис; (5) CDC/репликация в аналитику — Debezium читает WAL PostgreSQL и льёт в Kafka, оттуда в ClickHouse. Обязательно назовите и то, чего не делали брокером: синхронный запрос-ответ, где нужен немедленный результат, — там gRPC/HTTP, а RPC поверх очереди почти всегда плохая идея.
Как вы обрабатывали dead-letter сообщения и гарантии доставки?
Заголовок раздела «Как вы обрабатывали dead-letter сообщения и гарантии доставки?»Коротко. См. выше про DLQ и про гарантии доставки. Практический ответ: at-least-once с acks=all и ручным коммитом после обработки, идемпотентность в обработчике, ограниченные ретраи с экспоненциальным backoff для транзиентных ошибок, немедленный DLQ для перманентных, алерт на любое сообщение в DLQ и ручной или скриптовый reprocess после фикса.
Глубже. Что стоит добавить, чтобы ответ звучал как из продакшена. Различайте ошибки по классам ещё в коде: errors.Is(err, ErrTransient) → ретрай, иначе → DLQ; смешивание этих веток приводит к тому, что битый JSON ретраится 50 раз и блокирует партицию. В сообщение-«надгробие» кладите заголовки: x-original-topic, x-original-partition, x-original-offset, x-error, x-retry-count, x-first-failed-at — без них разбор DLQ через неделю невозможен. Ретраи с backoff внутри консьюмера ограничивайте суммарным временем меньше max.poll.interval.ms, иначе вы сами вызовете ребалансировку; для длинных задержек используйте отдельные retry-топики (orders.retry.5s, orders.retry.1m) или, в RabbitMQ, TTL-очередь с dead-letter обратно. Метрики, которые действительно смотрят: размер DLQ, rate попаданий в DLQ, возраст самого старого сообщения в DLQ, consumer lag и rate ребалансировок. И обязательно проговорите обратную сторону: reprocess из DLQ должен быть идемпотентным и в контролируемом темпе, иначе разбор инцидента сам станет инцидентом.
Что такое очередь сообщений в брокере?
Заголовок раздела «Что такое очередь сообщений в брокере?»Коротко. Очередь — именованный буфер, куда брокер складывает сообщения до тех пор, пока потребитель их не заберёт и не подтвердит; по умолчанию это FIFO-структура с семантикой «сообщение получает один потребитель и после ack оно удаляется».
Глубже. В RabbitMQ очередь — первоклассная сущность с собственными свойствами: durable (переживает рестарт брокера), exclusive, auto-delete, x-message-ttl, x-max-length, x-dead-letter-exchange, x-delivery-limit. Продюсер в очередь напрямую не пишет — он публикует в exchange, а связь «exchange → очередь» задаёт binding. Типы очередей в современном RabbitMQ: classic (простые, локальные), quorum (реплицированные через Raft, рекомендованный дефолт для надёжности; в RabbitMQ 4.0 классические mirrored-очереди удалены) и streams (append-only лог, kafka-подобная модель с возможностью перечитать историю). Строгий FIFO в очереди сохраняется только при одном потребителе без ретраев: несколько конкурирующих потребителей и возвраты через nack requeue=true порядок ломают.
В Kafka сущности «очередь» в этом смысле нет — есть партиция лога, из которой чтение не удаляет данные, а лишь двигает офсет. Это принципиальное различие, и его стоит проговорить: в очереди состояние доставки хранит брокер (на каждое сообщение), в логе — читатель (одно число на партицию). Первое даёт per-message ack, приоритеты и отложенную доставку; второе — дешёвое масштабирование, воспроизведение истории и несколько независимых потребителей.
Работал ли с брокерами?
Заголовок раздела «Работал ли с брокерами?»Коротко. См. выше про «Использовал ли брокеры сообщений». Короткая версия ответа: «Да — Kafka в текущем проекте, RabbitMQ в предыдущем», и сразу один конкретный пример задачи, чтобы задать направление дальнейшим вопросам.
Глубже. Такие короткие вопросы почти всегда служат разогревом перед углублением, поэтому отвечайте так, чтобы задать удобную для себя тему. Если вы уверенно знаете партиционирование и ребалансировки — упомяните масштабирование консьюмеров; если сильны в надёжности — скажите про outbox и идемпотентность. Плохая тактика — ответить односложным «да» и ждать: интервьюер тогда сам выберет тему, и она может оказаться неудобной.
С какими брокерами ссобщения работал?
Заголовок раздела «С какими брокерами ссобщения работал?»Коротко. Дубликат предыдущих вопросов о личном опыте. Отличие только в формулировке — назовите конкретные продукты и версии/режимы: «Kafka (в основном 2.8 и 3.x, кластер из 3 брокеров, RF=3), RabbitMQ 3.11 с quorum-очередями, Redis Streams для мелких внутренних задач».
Глубже. Упоминание версий и конфигурации кластера сильно повышает доверие к ответу, потому что эти детали невозможно выучить абстрактно. Если помните — скажите replication.factor, min.insync.replicas, число партиций на ключевом топике и порядок величины трафика («порядка 5 тысяч сообщений в секунду в пике»). Если не помните точных цифр — лучше сказать «порядок такой-то, точные значения не назову», чем придумать.
Сколько консьюмеров можно создать на три партиции?
Заголовок раздела «Сколько консьюмеров можно создать на три партиции?»Коротко. Создать можно сколько угодно, но полезно работающих в одной группе — максимум три: по одному на партицию. Четвёртый и далее будут простаивать как горячий резерв. Если же консьюмеры в разных группах, то каждая группа независимо получает все три партиции и весь поток целиком.
Глубже. Вопрос проверяет, различаете ли вы «внутри группы» и «между группами» — обязательно проговорите оба случая, иначе ответ считается неполным. Практическая рекомендация: число партиций подбирают исходя из целевой пропускной способности (нужный RPS / RPS одного консьюмера) с запасом ×2–3 на будущий рост, потому что увеличить партиции можно, но это перераспределит ключи и разорвёт порядок по ключу на границе, а уменьшить нельзя вовсе. При этом бесконечно наращивать партиции тоже нельзя: каждая партиция — это файловые дескрипторы, память под буферы на брокере и клиенте и время на выборы лидеров при отказе; десятки тысяч партиций на кластер деградируют его (в KRaft-режиме, ставшем единственным в Kafka 4.0, потолок заметно выше, чем был с ZooKeeper, но он есть).
Чем занимается брокер?
Заголовок раздела «Чем занимается брокер?»Коротко. См. выше «Что такое брокер сообщений». Функционально: принимает и валидирует публикации, персистит сообщения на диск и реплицирует их между узлами, маршрутизирует по топикам/очередям, отдаёт консьюмерам, отслеживает подтверждения и позиции чтения, чистит данные по retention/TTL и обеспечивает отказоустойчивость через выбор лидеров.
Глубже. Отличие от предыдущего вопроса — акцент на runtime-обязанностях, поэтому уместно перечислить внутренние подсистемы на примере Kafka: сетевой слой с батчингом и сжатием, лог-менеджер с сегментами и индексами (по офсету и по времени), репликация с понятием ISR и high watermark, контроллер (в KRaft-режиме — кворум контроллеров на Raft, ZooKeeper удалён в Kafka 4.0), координатор групп (назначение партиций, хранение офсетов), транзакционный координатор, фоновые задачи (retention, компакция, перебалансировка реплик). Отдельно подчеркните, чего брокер не делает и делать не должен: не выполняет бизнес-логику, не трансформирует данные, не гарантирует глобальный порядок между партициями и не решает за вас проблему дубликатов.
Можно сделать так чтобы продюсер писал все в одну партицию?
Заголовок раздела «Можно сделать так чтобы продюсер писал все в одну партицию?»Коротко. Да, тремя способами: явно указать номер партиции в записи, использовать один и тот же ключ для всех сообщений, или создать топик с одной партицией. Получите строгий глобальный порядок, но потеряете параллелизм: пропускная способность топика будет ограничена одним брокером-лидером и одним консьюмером в группе.
Глубже. В kafka-go явное указание выглядит так:
package main
import ( "context" "log"
"github.com/segmentio/kafka-go")
func main() { w := &kafka.Writer{ Addr: kafka.TCP("localhost:9092"), Topic: "orders", // Balancer не используется, если Partition задан явно в сообщении RequiredAcks: kafka.RequireAll, } defer w.Close()
err := w.WriteMessages(context.Background(), kafka.Message{ Partition: 0, // жёстко в партицию 0 Key: []byte("k"), Value: []byte("payload"), }) if err != nil { log.Fatal(err) }}Вариант «один ключ на всё» (Key: []byte("global")) даст тот же эффект без привязки к номеру: hash(key) % N всегда попадёт в одну и ту же партицию, и это устойчивее к изменению номеров, но развалится при увеличении числа партиций (ключ переедет). Осознанно так делают редко и только когда глобальный порядок действительно нужен и трафик мал — например, поток команд конфигурации или журнал изменений одной небольшой сущности. Гораздо чаще правильный ответ звучит иначе: не нужен глобальный порядок, нужен порядок по сущности — тогда ключом берут order_id/user_id, и вы получаете и порядок там, где он важен, и параллелизм. Отдельно предупредите про «горячий ключ»: если 80 % трафика идёт под одним tenant_id, одна партиция станет узким местом даже при 50 партициях в топике; лечится составным ключом (tenant_id:bucket) с потерей строгого порядка внутри тенанта.
Действительно ли exactly once обеспечит отправку одного сообщения?
Заголовок раздела «Действительно ли exactly once обеспечит отправку одного сообщения?»Коротко. Нет. Exactly-once — это гарантия однократного эффекта, а не однократной передачи по сети: физически сообщение может быть отправлено и записано несколько раз, просто дубликаты будут отброшены по (PID, sequence) или скрыты транзакционной семантикой. И действует эта гарантия только внутри Kafka; на внешние побочные эффекты она не распространяется.
Глубже. Полезно сослаться на теорию: задача надёжной однократной доставки в асинхронной сети с отказами (проблема «двух генералов») в общем виде неразрешима — отправитель никогда не может достоверно знать, потерялось сообщение или потерялся ответ. Поэтому любая практическая реализация — это at-least-once на транспорте плюс дедупликация в приёмнике; «exactly-once» — маркетинговое название этой связки, честнее говорить effectively-once.
Где именно гарантия перестаёт действовать, стоит перечислить прямо: (1) обработчик после коммита Kafka-транзакции отправляет HTTP-запрос во внешний сервис — при падении между коммитом и отправкой запрос не уйдёт, при падении после отправки, но до коммита — уйдёт дважды; (2) запись в стороннюю БД без общей транзакции с офсетами; (3) transactional.id, сгенерированный случайно при каждом старте — тогда после рестарта старая транзакция не будет корректно зафенсена; (4) потребитель с isolation.level=read_uncommitted увидит откаченные записи. Вывод для собеседования: «включаем EOS там, где цепочка Kafka→Kafka, а на границе с внешним миром всё равно ставим ключ идемпотентности».
100 партиций, 200 консьюмеров, как это будет работать?
Заголовок раздела «100 партиций, 200 консьюмеров, как это будет работать?»Коротко. Если все 200 в одной группе — 100 консьюмеров получат по одной партиции, оставшиеся 100 будут простаивать без назначений. Пропускная способность не вырастет по сравнению со 100 консьюмерами. Если это, скажем, две группы по 100 — обе полностью загружены и каждая независимо обрабатывает весь поток.
Глубже. Кроме «половина простаивает» есть неочевидные издержки, о которых полезно упомянуть. Каждый член группы участвует в ребалансировках, шлёт heartbeat’ы и держит соединения с брокерами; 200 членов в группе делают ребалансировку заметно дороже и дольше, а при eager-протоколе каждая такая пауза останавливает всю группу. Поэтому «лишние» консьюмеры не бесплатны: разумный резерв — единицы процентов, а не 100 %. Если задача действительно требует 200-кратного параллелизма, решения два: увеличить число партиций до 200+ (и заранее заложить это при создании топика) или параллелить обработку внутри консьюмера. Если же 200 инстансов появились не ради параллелизма, а потому что сервис горизонтально масштабируется по HTTP-нагрузке, то правильный ответ — развести роли: HTTP-инстансы отдельно, консьюмер-инстансы отдельно, в количестве, равном числу партиций.
25 партиций, 20 консьюмеров, как это будет работать?
Заголовок раздела «25 партиций, 20 консьюмеров, как это будет работать?»Коротко. Все 20 консьюмеров получат работу, но неравномерно: 25 = 20 × 1 + 5, поэтому пять консьюмеров получат по две партиции, а пятнадцать — по одной. Простаивающих нет, но нагрузка распределена с перекосом до 2× — это нормально и работает, просто пропускная способность группы будет ограничена самыми загруженными участниками.
Глубже. Точное распределение зависит от partition.assignment.strategy. RoundRobinAssignor и CooperativeStickyAssignor дадут именно 5×2 + 15×1. RangeAssignor считает диапазоны отдельно по каждому топику: для одного топика результат будет тем же, но при нескольких топиках перекос накапливается, потому что «лишние» партиции каждого топика достаются одним и тем же первым консьюмерам в лексикографическом порядке — это классическая причина того, что несколько подов работают вдвое тяжелее остальных.
Практический вывод: если хочется ровной нагрузки, берите число партиций, кратное числу консьюмеров (25 партиций и 5 консьюмеров дадут ровно по 5), либо просто держите партиций заметно больше консьюмеров — тогда относительный перекос падает (при 100 партициях на 20 консьюмеров получится ровно по 5, а при 101 — перекос всего 20 %). И помните, что перекос по партициям может быть куда сильнее перекоса по количеству: одна партиция с горячим ключом даст больше нагрузки, чем три холодные.
Всегда ли один и тот же консьюмер будет обрабатывать 2 партиции?
Заголовок раздела «Всегда ли один и тот же консьюмер будет обрабатывать 2 партиции?»Коротко. Нет. Назначение партиций не постоянно: любая ребалансировка (падение или добавление участника, рестарт пода, изменение числа партиций, подписка на новый топик) может перераспределить партиции, и «двойную нагрузку» получит другой инстанс. Стабильность обеспечивают только частично — sticky-стратегиями и static membership.
Глубже. StickyAssignor/CooperativeStickyAssignor при ребалансировке стараются сохранить максимум прежних назначений, перемещая минимально необходимое число партиций, но «стараются» — не гарантия: при изменении состава группы часть партиций обязана переехать. group.instance.id (static membership) делает сильнее: участник с постоянным идентификатором, перезапустившийся в пределах session.timeout.ms, получает те же самые партиции и ребалансировку вообще не вызывает — это то, что обычно и нужно в Kubernetes с rolling update. Если требуется жёсткая привязка «инстанс ↔ партиция» (например, локальный кеш или state store на диске, привязанный к партиции), правильный путь — не полагаться на ассайнер, а использовать assign() вместо subscribe(): ручное назначение партиций без консьюмер-группы, где вы сами отвечаете за распределение и failover. Ценой идёт потеря автоматической ребалансировки. И обязательно проектируйте обработчик так, чтобы переезд партиции был безопасен: сохраняйте состояние вне процесса или восстанавливайте его из лога, реализуйте ConsumerRebalanceListener-подобную логику (коммит офсетов и слив буферов при отзыве партиции).
С какими брокерами сообщений работали и в чем их особенности?
Заголовок раздела «С какими брокерами сообщений работали и в чем их особенности?»Коротко. См. выше вопросы про личный опыт; здесь дополнительно требуется сравнительная характеристика. Короткий каркас: Kafka — распределённый лог, высокая пропускная способность, порядок в партиции, реплей истории, несколько независимых групп; RabbitMQ — брокер очередей AMQP, гибкая маршрутизация, per-message ack, TTL и нативный DLQ, отложенные ретраи; NATS/JetStream — лёгкий и быстрый, простая эксплуатация; SQS/SNS — управляемый сервис без эксплуатации, но с ограничениями (видимость, порядок только в FIFO-очередях).
Глубже. Опорные различия, которые стоит назвать. Модель хранения: Kafka не удаляет прочитанное (retention/компакция) — отсюда реплей и несколько групп; RabbitMQ удаляет после ack — отсюда невозможность перечитать. Гранулярность подтверждения: у RabbitMQ ack на сообщение (можно подтверждать вразнобой, можно nack с requeue), у Kafka — один офсет на партицию, поэтому «пропустить одно и продолжить» требует прикладной логики. Порядок: Kafka гарантирует внутри партиции, RabbitMQ — только при одном потребителе. Маршрутизация: у RabbitMQ богатая (topic-маски, headers-exchange), у Kafka её нет — фильтрация на стороне консьюмера или отдельные топики. Отложенная доставка/приоритеты: есть у RabbitMQ (TTL + DLX, плагин delayed message, x-max-priority), у Kafka нет. Масштаб: Kafka линейно масштабируется партициями и держит сотни тысяч сообщений в секунду на узел за счёт последовательного I/O и zero-copy; RabbitMQ упирается в отдельную очередь. Эксплуатация: Kafka в 4.0 отказалась от ZooKeeper (KRaft), стало проще, но кластер всё равно тяжёлый; RabbitMQ 4.0 сделал quorum-очереди основным механизмом надёжности вместо удалённых mirrored-очередей. Практическое правило выбора: поток событий, аналитика, реплей, много подписчиков → Kafka; распределение задач с ретраями, приоритетами, отложенным выполнением и сложной маршрутизацией → RabbitMQ; хочется минимум эксплуатации и объёмы небольшие → managed-сервис или NATS.
Частые ошибки на собесе
Заголовок раздела «Частые ошибки на собесе»- Утверждать, что «Kafka удаляет сообщение после того, как консьюмер его прочитал». Kafka удаляет по retention или компакции; чтение только двигает офсет и ничего не удаляет.
- Считать, что «exactly-once» решает проблему дубликатов во внешних системах. Он ограничен циклом Kafka→Kafka; на границе с HTTP-API или сторонней БД нужна идемпотентность.
- Не различать «внутри группы» и «между группами» в вопросах про число консьюмеров. Правильный ответ почти всегда состоит из двух частей.
- Обещать глобальный порядок в топике. Порядок есть только внутри партиции, и он ломается при увеличении числа партиций, при отсутствии ключа и при параллельной обработке батча внутри консьюмера.
- Коммитить офсет до обработки (или оставлять
enable.auto.commit=trueпри долгой обработке) и при этом называть свою схему at-least-once. Это at-most-once с потерями. - Забывать, что бесконечные ретраи внутри консьюмера блокируют партицию и выбивают участника по
max.poll.interval.ms, вызывая ребалансировку. Ретраи должны быть ограничены по времени, дальше — DLQ или retry-топик. - Путать
acks(подтверждение записи продюсеру) сackконсьюмера (подтверждение обработки) и считатьacks=allзащитой от дубликатов — это про durability, а не про дедупликацию. - Заводить DLQ и не заводить на неё алерт и процедуру reprocess: очередь молча накапливает потерянные бизнес-события.
Что почитать
Заголовок раздела «Что почитать»- Kafka Documentation — Design и Implementation: https://kafka.apache.org/documentation/#design (разделы про реплики, ISR, гарантии, консьюмер-группы).
- KIP-98 «Exactly Once Delivery and Transactional Messaging» и KIP-129/KIP-447: https://cwiki.apache.org/confluence/display/KAFKA/KIP-98+-+Exactly+Once+Delivery+and+Transactional+Messaging
- KIP-429 (Incremental Cooperative Rebalancing) и KIP-848 (новый протокол консьюмер-групп): https://cwiki.apache.org/confluence/display/KAFKA/KIP-848%3A+The+Next+Generation+of+the+Consumer+Rebalance+Protocol
- RabbitMQ — Dead Letter Exchanges и Quorum Queues: https://www.rabbitmq.com/dlx.html , https://www.rabbitmq.com/quorum-queues.html
- Martin Kleppmann, «Designing Data-Intensive Applications», глава 11 «Stream Processing» — лучшее системное изложение логов, гарантий и дедупликации.