Перейти к содержимому

Apache Kafka

Kafka — это не «очередь сообщений» в привычном смысле, а распределённый реплицируемый лог с сохранением на диск. Единица публикации — топик, но физически топик существует только как набор партиций. Партиция — это упорядоченный append-only файл (точнее, каталог на диске, разбитый на сегменты), в который продюсер только дописывает в конец, а каждой записи присваивается монотонно растущий номер — offset. Ничего не «удаляется при чтении»: сообщение живёт ровно столько, сколько разрешает retention (по времени retention.ms, по размеру retention.bytes) или пока его не вытеснит компакция по ключу. Отсюда сразу вытекает всё остальное: чтение неразрушающее, читателей может быть сколько угодно и они не мешают друг другу, а «прочитано или нет» — это состояние клиента, а не брокера.

Вторая половина модели — параллелизм и порядок. Порядок гарантируется только внутри одной партиции, потому что только там есть общий монотонный счётчик. Партиция — это же и единица параллелизма: одну партицию в рамках одной consumer group в один момент времени читает ровно один консьюмер. Значит, число партиций — это потолок параллелизма группы, а роутинг сообщения в партицию (по хешу ключа murmur2(key) % partitions, либо явно указанным номером, либо sticky-батчами при key == nil) — это одновременно и выбор «полосы упорядоченности». Классический приём: ключ = user_id / order_id, тогда все события одной сущности гарантированно попадают в одну партицию и обрабатываются по порядку, а между разными сущностями порядок никого не волнует.

Третья половина — надёжность. Каждая партиция реплицируется на replication.factor брокеров: один лидер (принимает все чтения/записи) и фолловеры, которые тянут данные с лидера. Множество реплик, догнавших лидера, называется ISR (in-sync replicas). Продюсер с acks=all получает подтверждение, только когда запись попала в лог всех ISR, а требуемый минимум задаётся брокерным min.insync.replicas — эта пара и есть настоящая настройка durability. Координацией кластера с Kafka 3.3 занимается встроенный KRaft-кворум контроллеров (в Kafka 4.0 ZooKeeper убран окончательно), а координацией групп потребителей — брокер-group coordinator, который держит членство группы и коммиты оффсетов во внутреннем компактящемся топике __consumer_offsets.

Наконец, производительность. Kafka быстрая не из-за магии, а из-за того, что она делает ровно то, что дёшево для железа: последовательная запись в конец файла, батчинг и сжатие на стороне продюсера (linger.ms, batch.size, compression.type), page cache вместо собственного кеша в heap и sendfile (zero-copy) при отдаче консьюмерам, если не включён TLS. Broker при этом «глупый»: он не роутит, не фильтрует, не помнит, кто что подтвердил по каждому сообщению, — вся логика на клиенте.

Коротко. Kafka — персистентный лог с pull-моделью: сообщения не удаляются после чтения, позиция чтения (offset) принадлежит консьюмеру, порядок гарантирован в партиции, параллелизм ограничен числом партиций. RabbitMQ — классический брокер очередей с push-моделью и умной маршрутизацией: exchange по правилам раскладывает сообщение по очередям, брокер отслеживает ack по каждому сообщению и удаляет его после подтверждения, можно делать retry/DLQ/TTL/приоритеты из коробки.

Глубже. Практическая разница в трёх местах. Первое — реплей: в Kafka можно перечитать историю (сбросить оффсет на неделю назад, поднять новый сервис и прогнать весь поток), в RabbitMQ подтверждённое сообщение физически исчезло. Второе — гранулярность подтверждения: RabbitMQ подтверждает отдельное сообщение и умеет вернуть конкретное в очередь (basic.nack requeue), Kafka коммитит только позицию в партиции, поэтому «пропустить одно битое и продолжить» нужно реализовывать самому (DLQ-топик). Третье — пропускная способность и масштаб: Kafka на порядки лучше держит сотни МБ/с и долгое хранение, RabbitMQ выигрывает в сценариях «задачи с разным приоритетом, сложный роутинг, десятки тысяч эфемерных очередей, RPC request/reply». Стоит упомянуть, что грань размылась: у RabbitMQ есть quorum queues (Raft) и Streams (лог-подобная структура с offset и реплеем, с 3.9), а у Kafka в 4.0 в early access появились share groups (KIP-932) с очередной семантикой. Но по умолчанию выбор такой: поток событий и аналитика — Kafka, задачи и командная маршрутизация — RabbitMQ.

Коротко. Плюсы: очень высокая пропускная способность, горизонтальное масштабирование партициями, durable-хранение с реплеем, порядок в партиции, зрелая экосистема (Connect, Streams, ksqlDB, Schema Registry). Минусы: тяжёлая эксплуатация, нет per-message ack и повторной доставки конкретного сообщения, нет приоритетов и отложенной доставки, параллелизм жёстко упирается в число партиций, latency выше, чем у in-memory брокеров, «одно ядовитое сообщение» может застопорить партицию.

Глубже. Из неочевидного в минусах: изменение числа партиций ломает привязку ключ→партиция; много партиций — это много открытых файлов, дольше выборы лидеров и больше нагрузка на контроллер; ребалансировки группы (stop-the-world при eager-протоколе) дают паузы обработки; exactly-once работает только внутри Kafka, а на границе с внешней БД его всё равно нужно достраивать идемпотентностью. Из неочевидного в плюсах: log compaction превращает топик в «журнал изменений с сохранением последнего значения по ключу», что позволяет держать в Kafka снапшот состояния (CDC, конфигурация, materialized view), а не только поток событий.

кафка - 1 партишн и 3 консьюмера, 3 партишена - 1 консьюмер

Заголовок раздела «кафка - 1 партишн и 3 консьюмера, 3 партишена - 1 консьюмер»

Коротко. Если партиция одна, а консьюмеров в группе три — работает один, два простаивают без назначенных партиций (масштабирования нет). Если партиций три, а консьюмер один — он получает все три партиции и читает их последовательно в одном poll; работать будет, но параллелизма нет и лаг может расти.

Глубже. Формула простая: одна партиция в один момент времени принадлежит ровно одному консьюмеру группы, но один консьюмер может держать много партиций. Поэтому «лишние» консьюмеры — это не ошибка, а горячий резерв: при падении активного произойдёт ребаланс, и партиция перейдёт к простаивавшему. Если хочется, чтобы все три экземпляра действительно работали, надо увеличивать число партиций, а не число процессов. Обратная ситуация (3 партиции — 1 консьюмер) лечится либо запуском ещё двух экземпляров с тем же group.id, либо внутренним пулом воркеров, но тогда сам следишь за порядком и коммитами.

Что позволяет не читать каждый раз сообщения из кафки заново?

Заголовок раздела «Что позволяет не читать каждый раз сообщения из кафки заново?»

Коротко. Закоммиченный offset консьюмер-группы: брокер-координатор хранит для пары (group.id, topic, partition) номер следующего сообщения к чтению во внутреннем топике __consumer_offsets. При старте или ребалансе консьюмер запрашивает этот оффсет и продолжает с него.

Глубже. Если коммита нет вообще (новая группа), поведение определяет auto.offset.reset: latest — начать с конца (по умолчанию), earliest — с начала, none — ошибка. Коммит бывает автоматическим (enable.auto.commit=true, раз в auto.commit.interval.ms, фактически во время poll) и ручным (commitSync/commitAsync, в Go — CommitMessages/StoreOffsets). Оффсеты группы имеют собственный retention (offsets.retention.minutes, по умолчанию 7 дней): если группа не читала дольше, коммиты вычищаются и группа стартует по auto.offset.reset. Кроме коммита в Kafka, оффсет можно хранить где угодно — например, в той же транзакции в PostgreSQL, куда пишется результат обработки, и на старте делать Seek на сохранённую позицию; это стандартный способ получить exactly-once относительно внешней БД.

Коротко. Оффсет хранит сам кластер Kafka — в служебном топике __consumer_offsets (по умолчанию 50 партиций, cleanup.policy=compact). Записывает туда консьюмер (коммитом), а обслуживает конкретный брокер — group coordinator, выбираемый как лидер партиции hash(group.id) % 50.

Глубже. Важно разделить три вещи: позиция чтения (position) — состояние в памяти клиента, продвигается каждым poll; закоммиченный оффсет — то, что записано в __consumer_offsets и переживёт перезапуск; high watermark / log end offset — состояние партиции на брокере. Ключ записи в __consumer_offsets(group, topic, partition), значение — оффсет + метаданные + таймстемп, поэтому компакция оставляет только последний коммит. Именно из-за того, что оффсет — это состояние группы, а не сообщения, две разные группы читают один топик полностью независимо. Исторически (до Kafka 0.9) оффсеты лежали в ZooKeeper — если спросят «где хранились раньше», это правильный ответ.

Можно ли читать один топик несколькими консьюмерами параллельно?

Заголовок раздела «Можно ли читать один топик несколькими консьюмерами параллельно?»

Коротко. Да, двумя способами: несколько консьюмеров в одной группе (партиции делятся между ними, каждое сообщение обрабатывается один раз — это масштабирование) и несколько разных групп (каждая группа получает полную копию потока — это fan-out).

Глубже. Внутри одной группы параллелизм ограничен числом партиций: 8 партиций — максимум 8 полезных консьюмеров. Между группами ограничений нет, кроме сети и диска брокеров: сто групп прочитают топик сто раз, при этом данные, скорее всего, отдадутся из page cache без обращения к диску. Третий вариант — вручную назначить партиции (Assign вместо Subscribe), тогда нет группы, нет ребалансов и нет автоматических коммитов; так делают, когда нужен полный контроль (например, консьюмер привязан к шардированному локальному стейту).

Коротко. Партиция — упорядоченная неизменяемая append-only последовательность сообщений, физически лежащая на диске брокера как каталог с сегментами. Это единица параллелизма, репликации и порядка: внутри партиции есть сквозной offset и строгий порядок, между партициями порядка нет.

Глубже. На диске партиция orders-0 — это каталог orders-0/ с файлами сегментов <baseOffset>.log, разреженными индексами .index (offset → позиция в файле) и .timeindex (timestamp → offset). Активный сегмент один, он ротируется по segment.bytes/segment.ms, а retention/компакция работают посегментно — поэтому Kafka никогда не «удаляет сообщение», она удаляет или переписывает целый сегмент. У партиции есть лидер и фолловеры на других брокерах, high watermark (граница видимого консьюмерам) и leader epoch для корректного усечения логов после смены лидера.

Коротко. Это не вопрос, а пометка интервьюера «спрашивал стандартный набор по Kafka». Стандартный набор: топик/партиция/оффсет/consumer group, гарантии доставки, порядок, репликация и ISR, acks, ребалансировка, consumer lag, отличия от RabbitMQ.

Глубже. На такой пункт полезно иметь заготовленный десятиминутный рассказ «Kafka сверху вниз»: лог и офсеты → партиции и порядок → группы и ребаланс → репликация, ISR и acks → гарантии доставки и идемпотентность → эксплуатация (lag, DLQ, ретраи). Если интервьюер начинает с «расскажите про Kafka», такой каркас закрывает 80% последующих уточнений.

Гарантирована ли доставка сообщений в кафка?

Заголовок раздела «Гарантирована ли доставка сообщений в кафка?»

Коротко. Да, при правильной конфигурации Kafka даёт as-least-once: продюсер с acks=all + ретраями и min.insync.replicas >= 2 гарантирует, что запись сохранена и реплицирована; консьюмер, коммитящий оффсет только после успешной обработки, гарантирует, что сообщение не потеряется. Но «из коробки, без настроек и без правильного порядка коммита» гарантий нет — можно и потерять, и продублировать.

Глубже. Потери возможны на трёх стыках. Продюсер: acks=0/1 — падение лидера до репликации теряет батч; лечится acks=all + min.insync.replicas=2 при RF=3 и enable.idempotence=true (с Kafka 3.0 идемпотентность и acks=all включены по умолчанию). Брокер: unclean.leader.election.enable=true разрешит выбрать лидером отставшую реплику и потерять хвост лога — по умолчанию выключено, включать нельзя. Консьюмер: авто-коммит или коммит до обработки превращает падение в потерю. Отдельно помним, что «сообщение записано в лог» ≠ «сообщение на диске»: Kafka по умолчанию полагается на репликацию, а не на fsync каждой записи (flush.messages/flush.ms не трогают), поэтому одновременная потеря питания на всех репликах теоретически может стоить хвоста.

За счет чего обеспечивается отказоустойчивость Kafka?

Заголовок раздела «За счет чего обеспечивается отказоустойчивость Kafka?»

Коротко. За счёт репликации партиций между брокерами (replication.factor), механизма ISR с автоматическим выбором нового лидера при падении текущего, кворумного контроллера KRaft (раньше — ZooKeeper) и того, что состояние потребителей вынесено в реплицируемый топик __consumer_offsets, а падение консьюмера лечится ребалансом группы.

Глубже. Механика: у каждой партиции есть лидер и N-1 фолловеров, фолловеры постоянно делают fetch с лидера; реплика считается in-sync, пока успевает за replica.lag.time.max.ms (по умолчанию 30 с). Консьюмерам видно только то, что ниже high watermark — то есть реплицировано во все ISR, поэтому смена лидера не «откатывает» уже прочитанное. Если лидер умирает, контроллер назначает нового из ISR; leader.epoch не даёт старой реплике вернуться с расходящимся хвостом. При RF=3 и min.insync.replicas=2 кластер переживает падение одного брокера без потерь и без остановки записи, при падении двух — запись с acks=all начинает отдавать NotEnoughReplicas (сознательный выбор consistency вместо availability). Дополнительно: rack awareness (broker.rack) раскладывает реплики по стойкам/зонам, а MirrorMaker 2 реплицирует между кластерами.

Какие гарантии доставки предоставляет Kafka?

Заголовок раздела «Какие гарантии доставки предоставляет Kafka?»

Коротко. Три: at-most-once (коммитим оффсет до обработки — возможна потеря, дублей нет), at-least-once (коммитим после обработки — дубли возможны, потерь нет; это дефолт для здравой конфигурации) и exactly-once (транзакции producer’а + isolation.level=read_committed, работает для сценария read-process-write внутри Kafka).

Глубже. Гарантия — это не галочка в конфиге, а комбинация настроек продюсера, брокера и порядка действий в консьюмере. Продюсер: acks, enable.idempotence, retries, max.in.flight.requests.per.connection <= 5 (иначе идемпотентность не сохраняет порядок). Брокер: min.insync.replicas, unclean.leader.election.enable=false. Консьюмер: где именно вызывается коммит и что делает код при ошибке. На собеседовании выигрышно сказать: «Kafka даёт at-least-once по умолчанию, а exactly-once end-to-end во внешнюю систему всё равно достигается идемпотентностью приёмника или транзакционным outbox, транзакции Kafka сами по себе внешнюю БД не покрывают».

Коротко. См. выше «Standard Kafka questions» — это та же пометка о стандартном наборе, только по-русски. Отличий по содержанию нет.

Коротко. См. выше «Какие гарантии доставки предоставляет Kafka?»: at-most-once / at-least-once / exactly-once, и разница между ними определяется acks+идемпотентностью на продюсере и моментом коммита оффсета на консьюмере.

Глубже. Единственное, что стоит добавить именно к формулировке «расскажи»: отвечать лучше не списком терминов, а сценарием отказа. «Продюсер отправил, лидер упал до ответа → продюсер ретраит → без идемпотентности получаем дубль, с идемпотентностью брокер отбрасывает по (PID, epoch, sequence). Консьюмер обработал и упал до коммита → после ребаланса прочитает снова → нужен идемпотентный обработчик». Такой ответ показывает, что вы понимаете механику, а не заучили три слова.

Коротко. Это вопрос про личный опыт — отвечать надо конкретикой: какой сценарий, какие объёмы, какая библиотека, какие настройки и какую проблему пришлось решать.

Глубже. Хороший каркас ответа на 60–90 секунд: (1) роль Kafka в системе — «шина событий между сервисом заказов и биллингом», «CDC из Postgres через Debezium», «поток метрик»; (2) масштаб — сколько топиков, партиций, сообщений в секунду, какой размер сообщения, какой retention; (3) клиент — segmentio/kafka-go / confluent-kafka-go / sarama / franz-go и почему; (4) осознанные решения — ключ партиционирования, acks=all, ручной коммит после обработки, DLQ, идемпотентность потребителя по message_id; (5) одна реальная проблема и как чинили — растущий lag, ребалансы из-за max.poll.interval.ms, «ядовитое» сообщение, перекос по ключу. Типичные ошибки: отвечать «да, работал, читал и писал сообщения» без цифр и решений; называть настройки, которых не трогал; заявлять exactly-once, не сумев объяснить транзакции.

Коротко. Топик — именованная категория сообщений (логический канал). Партиция — физическая упорядоченная часть топика, единица параллелизма и порядка. Producer — клиент, который пишет сообщения в топик, выбирая партицию по ключу или явно. Consumer — клиент, который читает из партиций, обычно в составе consumer group, и коммитит оффсеты.

Глубже. Полезно добавить, что сообщение (record) — это (key, value, headers, timestamp, offset, partition); ключ нужен не для идентификации, а для выбора партиции и для компакции. Продюсер асинхронный: он копит записи в аккумуляторе, батчит по batch.size/linger.ms, сжимает и отправляет sender-потоком, поэтому «отправил» без ожидания ack ничего не гарантирует. Консьюмер — однопоточный по контракту (KafkaConsumer не потокобезопасен; в Go у kafka-go Reader тоже рассчитан на одну горутину чтения), цикл всегда poll → обработка → commit, причём poll заодно шлёт heartbeat и участвует в ребалансе.

Коротко. Consumer lag — разница между log end offset партиции (последним записанным оффсетом) и закоммиченным оффсетом группы, то есть сколько сообщений консьюмер ещё не обработал. Возникает, когда скорость производства выше скорости потребления.

Глубже. Причины по частоте: медленная обработка одного сообщения (синхронный поход в БД/внешний API на каждое), слишком мало партиций/консьюмеров для входящего потока, всплеск трафика или backfill, «залипание» на ошибке с бесконечным ретраем одного сообщения, длинные ребалансировки (во время ребаланса группа не читает вообще), перекос по ключу — одна партиция получает 80% трафика и её консьюмер не справляется, а остальные простаивают, а также сетевые/GC-паузы и слишком маленькие fetch.min.bytes/max.poll.records при большом RTT. Отдельно стоит различать lag в сообщениях и lag во времени: 100k сообщений могут быть 2 секундами отставания, а могут быть часом.

Коротко. Это пометка интервьюера о наборе тем, а не вопрос. По Kafka «дефолт» — топик/партиция/группа/оффсет, гарантии доставки, порядок, lag, отличия от RabbitMQ.

Глубже. На стыке Go + БД + Kafka чаще всего спрашивают одну связку: «как не потерять и не задублировать при записи в БД по событию из Kafka». Правильный ответ — идемпотентность на стороне БД (INSERT ... ON CONFLICT DO NOTHING по message_id или natural key), либо хранение оффсета в той же транзакции, что и результат; для обратного направления — transactional outbox + Debezium вместо «записал в БД, потом отправил в Kafka».

Коротко. Record (сообщение с key/value/headers/timestamp), topic, partition, offset, producer, consumer, consumer group, broker, cluster/controller, replica и ISR. Поверх — Kafka Connect (интеграции), Kafka Streams (обработка), Schema Registry (схемы).

Глубже. Стоит выделить «логическое vs физическое»: логические абстракции — топик, группа, оффсет как позиция; физические — брокер, партиция как каталог сегментов, реплика. И третий слой — служебные топики: __consumer_offsets (позиции групп) и __transaction_state (состояние транзакций продюсеров). То, что и метаданные, и координация выражены через те же самые логи, — важная идея архитектуры Kafka.

Что будет, если партиций больше, чем консюмеров в одной группе потребителей?

Заголовок раздела «Что будет, если партиций больше, чем консюмеров в одной группе потребителей?»

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

Глубже. Как именно распределятся — зависит от assignor’а. RangeAssignor (исторический дефолт) делит по каждому топику отдельно и при 5 партициях на 2 консьюмеров даст 3/2, а при нескольких топиках систематически перегружает первого консьюмера. RoundRobinAssignor раскладывает все партиции всех топиков по кругу — ровнее. StickyAssignor/CooperativeStickyAssignor дополнительно минимизируют переезды партиций при ребалансе, а cooperative-версия не останавливает всю группу (incremental rebalance). В Kafka 4.0 доступен новый протокол ребалансировки на брокере (KIP-848), где назначение считает координатор, а не лидер группы. Практическое следствие: держите число партиций кратным ожидаемому максимуму консьюмеров, иначе получите перекос.

На уровне чего поддерживается в Kafka порядок сообщений?

Заголовок раздела «На уровне чего поддерживается в Kafka порядок сообщений?»

Коротко. На уровне партиции. Внутри одной партиции порядок записи = порядок оффсетов = порядок чтения. Между партициями одного топика порядка нет вообще.

Глубже. Чтобы порядок был осмысленным для бизнеса, нужно, чтобы связанные события попадали в одну партицию, — это делается ключом (key = order_id). Порядок могут сломать: смена числа партиций (ключ начнёт хешироваться в другую партицию), ретраи продюсера при max.in.flight.requests.per.connection > 1 без идемпотентности (перестановка батчей), собственная многопоточная обработка внутри консьюмера и отправка «проблемных» сообщений в retry-топик с последующим возвратом. Идемпотентный продюсер (по умолчанию с 3.0) сохраняет порядок в пределах партиции при max.in.flight <= 5, потому что брокер отслеживает sequence number и отвергает записи не по порядку.

Как в Kafka реализуются семантики доставки «at most once» и «at least once» на уровне механизма работы системы?

Заголовок раздела «Как в Kafka реализуются семантики доставки «at most once» и «at least once» на уровне механизма работы системы?»

Коротко. Обе семантики определяются не флагом, а моментом коммита оффсета относительно обработки и настройками ретраев продюсера. At-most-once: коммитим оффсет сразу после poll, до обработки (и/или acks=0/1 без ретраев) — падение между коммитом и обработкой теряет сообщение. At-least-once: коммитим только после успешной обработки (и acks=all + ретраи на продюсере) — падение приводит к повторному чтению того же оффсета, то есть к дублю.

Глубже. На уровне механики at-least-once держится на том, что оффсет — это просто число в __consumer_offsets: после ребаланса/рестарта новый владелец партиции читает последний закоммиченный оффсет и продолжает с него, поэтому весь диапазон «обработано, но не закоммичено» проигрывается повторно. На стороне продюсера дубли рождаются из ретрая после таймаута ответа: запись уже в логе, ack потерян, продюсер шлёт ещё раз. Идемпотентный продюсер убирает именно этот класс дублей: он получает PID и epoch, нумерует записи в партиции sequence number’ами, брокер помнит последние 5 и отбрасывает повтор (DUPLICATE_SEQUENCE_NUMBER). Важно: enable.auto.commit=true — это не at-most-once, а «неуправляемый at-least-once с окном потери»: коммит выполняется внутри poll для записей предыдущего poll, поэтому при падении в середине пачки часть будет переобработана, а при асинхронной обработке в отдельных горутинах — потеряна.

// at-least-once на segmentio/kafka-go: коммит строго после успешной обработки
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"kafka:9092"},
GroupID: "billing",
Topic: "orders",
})
defer r.Close()
for {
m, err := r.FetchMessage(ctx) // читает, но НЕ коммитит
if err != nil {
return err
}
if err := handle(ctx, m); err != nil {
continue // не коммитим — сообщение прочитается снова
}
if err := r.CommitMessages(ctx, m); err != nil {
return err
}
}

Коротко. Exactly-once semantics (EOS) — гарантия, что каждое сообщение отражается в результате ровно один раз. В Kafka она строится из двух частей: идемпотентный продюсер (нет дублей от ретраев) и транзакции — атомарная запись в несколько партиций/топиков вместе с коммитом оффсетов потребления, при чтении с isolation.level=read_committed.

Глубже. Транзакционный продюсер задаёт transactional.id, вызывает InitTransactions (получает PID и увеличивает epoch, фенся старые зомби-инстансы), затем BeginTransaction → отправки → SendOffsetsToTransactionCommitTransaction. Состояние живёт в служебном топике __transaction_state, а в целевые партиции пишутся control-записи (commit/abort markers). Консьюмер с read_committed не отдаёт приложению записи выше LSO (last stable offset) и пропускает аборченные. Ключевые оговорки, которые ждут на собеседовании: EOS покрывает только паттерн read-process-write внутри Kafka (для этого он и встроен в Kafka Streams через processing.guarantee=exactly_once_v2); если результат уходит во внешнюю систему — БД, HTTP API — транзакция Kafka её не охватывает, и нужен либо идемпотентный приёмник (уникальный ключ, upsert), либо хранение оффсета в той же транзакции БД. Плюс EOS стоит latency и throughput: транзакции добавляют раунды и control-записи.

Что произойдет с уже существующими данными в топике Kafka, если увеличить количество партиций - например, с 10 до 20 - пока часть сообщений еще не была потреблена? Что произойдет, и можно ли уменьшить количество партиций обратно с 20 до 10?

Заголовок раздела «Что произойдет с уже существующими данными в топике Kafka, если увеличить количество партиций - например, с 10 до 20 - пока часть сообщений еще не была потреблена? Что произойдет, и можно ли уменьшить количество партиций обратно с 20 до 10?»

Коротко. Существующие данные никуда не переезжают: они остаются в исходных 10 партициях и будут дочитаны как обычно, ничего не теряется. Ломается другое — привязка ключа к партиции: с этого момента murmur2(key) % 20 даёт другой номер, поэтому новые события того же ключа могут попасть в новую партицию, и глобальный порядок «по ключу» нарушится (старые события в P3, новые в P13, обрабатываются разными консьюмерами). Уменьшить число партиций нельзя — Kafka поддерживает только увеличение; чтобы «уменьшить», создают новый топик с нужным числом партиций и перекладывают данные.

Глубже. Механически при --alter --partitions 20 контроллер просто создаёт 10 новых пустых партиций с новыми лидерами; группа получает ребаланс и начинает читать 20 партиций. Риск переупорядочивания реален для stateful-обработки (агрегаты по пользователю, машины состояний): в переходный момент возможно, что событие «order.updated» обработается раньше «order.created». Безопасный порядок действий: либо останавливать продюсеров и дожидаться нулевого лага перед расширением, либо сразу закладывать запас партиций при создании топика, либо использовать кастомный партиционер с consistent hashing / явной маппинг-таблицей, который не зависит от общего количества партиций. Причина запрета уменьшения проста: удаление партиции означало бы удаление её данных и разрыв монотонности оффсетов; вместо этого делают новый топик + Connect/MirrorMaker/свой перекладчик, а потребителей переключают по alias’у.

Коротко. См. выше «Что такое consumer lag и из-за чего он бывает?» — это то же самое: количество непрочитанных сообщений в партиции, logEndOffset - committedOffset, суммируется по всем партициям группы.

Глубже. Отличие только в терминологии инструментов: в JMX-метриках консьюмера это records-lag / records-lag-max (по каждой партиции и максимум), в kafka-consumer-groups.sh --describe — колонка LAG, в Burrow/kafka-exporter — kafka_consumergroup_lag. Метрика клиента показывает лаг относительно текущей позиции чтения, а серверная (по коммитам) — относительно закоммиченного оффсета; при редких коммитах они расходятся.

Коротко. Вопрос про опыт; см. каркас ответа выше в «Работал с кафкой?»: сценарий → масштаб → библиотека → сознательные настройки → решённая проблема.

Глубже. Отличие этой формулировки — она чаще звучит в начале секции как «фильтр»: ответ определит глубину дальнейших вопросов. Поэтому честно калибруйте уровень: «использовал как консьюмер и продюсер в проде, топики настраивал не я» — нормальный ответ, который избавит от вопросов по тюнингу брокера. Заявлять опыт эксплуатации кластера, если вы его не эксплуатировали, — быстрый способ провалиться на min.insync.replicas и unclean leader election.

Коротко. Сначала измерить, где узкое место, потом: увеличить число партиций и консьюмеров, ускорить обработку (батчинг записи в БД, убрать синхронные вызовы, кеш), распараллелить обработку внутри консьюмера с сохранением порядка по ключу, поднять max.poll.records/fetch.min.bytes, вынести тяжёлое в отдельный топик, а при аварии — временно пропустить/сбросить оффсет.

Глубже. Порядок действий на практике: 1) посмотреть, лаг во всех партициях или в одной — если в одной, это перекос по ключу, и добавление консьюмеров не поможет, надо менять ключ партиционирования или досыпать соли к ключу; 2) проверить, не крутится ли консьюмер на ретрае одного сообщения (лаг стоит, оффсет не движется — «ядовитое» сообщение, лечится DLQ); 3) профилировать обработчик — почти всегда это I/O на сообщение, лечится батчингом (INSERT ... VALUES (...), (...) на 500 записей и один коммит); 4) масштабировать: партиции ↑, поды ↑ (партиции только увеличиваются, помним про переупорядочивание); 5) параллелить внутри процесса: пул воркеров с шардированием по hash(key) % N, чтобы порядок по ключу сохранился, и коммит только по «сплошному» префиксу обработанных оффсетов; 6) проверить ребалансы (max.poll.interval.ms, session.timeout.ms, static membership через group.instance.id) — если группа половину времени ребалансится, лаг не разгрести. И отдельно: если данные устарели и не нужны, допустимо осознанно перескочить оффсет на конец (--reset-offsets --to-latest), но это потеря данных, решение бизнес-уровня.

Гарантирует ли Kafka уникальность и последовательность сообщений?

Заголовок раздела «Гарантирует ли Kafka уникальность и последовательность сообщений?»

Коротко. Последовательность — да, но только внутри партиции. Уникальность — нет: Kafka не дедуплицирует по содержимому. Идемпотентный продюсер убирает дубли, порождённые ретраями в рамках сессии продюсера и партиции, но повторная отправка того же события бизнес-логикой или переобработка после падения консьюмера даст дубль.

Глубже. Отсюда стандартное правило: обработчик должен быть идемпотентным. Практически — уникальный message_id в headers и таблица обработанных ключей / ON CONFLICT DO NOTHING, или естественная идемпотентность операции (upsert состояния вместо инкремента). Частичная «уникальность» есть у log compaction: топик с cleanup.policy=compact со временем оставляет по одному последнему значению на ключ, но это eventual-свойство хранения, а не гарантия доставки — до компакции консьюмер увидит все версии.

По Kafka: как работает консьюмер группа? Что такое партиция? Если партиций больше чем потребителей, как поделятся потребители и партиции? А наоборот потрателей больше партиций? За счет чего обеспечивается каждый потребитель будет читать только свои данные, а не другие, где хранятся инфо о том что-то прочитано или нет? При отказе приложения, когда перестартает за счет чего он продолжит читать с того сообщения на котором он упал? Где хранятся смещения, что в них?

Заголовок раздела «По Kafka: как работает консьюмер группа? Что такое партиция? Если партиций больше чем потребителей, как поделятся потребители и партиции? А наоборот потрателей больше партиций? За счет чего обеспечивается каждый потребитель будет читать только свои данные, а не другие, где хранятся инфо о том что-то прочитано или нет? При отказе приложения, когда перестартает за счет чего он продолжит читать с того сообщения на котором он упал? Где хранятся смещения, что в них?»

Коротко. Consumer group — набор консьюмеров с общим group.id, между которыми координатор эксклюзивно распределяет партиции топика: каждая партиция принадлежит ровно одному члену группы, поэтому пересечения по данным нет. Партиций больше потребителей — кто-то возьмёт несколько; потребителей больше партиций — лишние простаивают. Информация о прочитанном — закоммиченный оффсет в топике __consumer_offsets с ключом (group, topic, partition); после рестарта консьюмер получает от координатора этот оффсет и продолжает с него.

Глубже. Полный жизненный цикл: консьюмер шлёт FindCoordinator (координатор = лидер партиции hash(group.id) % offsets.topic.num.partitions), затем JoinGroup — координатор собирает участников, назначает одного лидером группы; лидер считает assignment выбранным assignor’ом и возвращает его в SyncGroup, координатор рассылает каждому его партиции. Живость поддерживается heartbeat-потоком (heartbeat.interval.ms ≈ 3 с, session.timeout.ms 45 с по умолчанию с 3.0) и ограничением на длительность обработки max.poll.interval.ms (5 минут): превысил — тебя выкинут и начнётся ребаланс. При падении процесса координатор по истечении session timeout инициирует ребаланс и передаёт партиции живым членам, которые начинают с последнего закоммиченного оффсета — отсюда возможная переобработка «хвоста». В самой записи оффсета лежат: сам оффсет (номер следующего сообщения к чтению, то есть processed + 1), leader epoch, опциональные строковые метаданные и таймстемп коммита; записи компактятся, так что хранится только последняя. Static membership (group.instance.id) позволяет при рестарте пода не ронять ребаланс: координатор ждёт возвращения того же участника.

Условная задача: какие-то сообщения поступают через Kafka и мы должны эти сообщения перекидывать в другую систему, но есть ограничьение, что тa внешняя система (куда мы перекидываем) не может обрабатывать больше чем 10 одновременных запросов, как реализовать это ограничьение используя примитивы языка Go?

Заголовок раздела «Условная задача: какие-то сообщения поступают через Kafka и мы должны эти сообщения перекидывать в другую систему, но есть ограничьение, что тa внешняя система (куда мы перекидываем) не может обрабатывать больше чем 10 одновременных запросов, как реализовать это ограничьение используя примитивы языка Go?»

Коротко. Семафор на буферизованном канале ёмкостью 10 (или golang.org/x/sync/semaphore, или errgroup.Group.SetLimit(10)): перед запросом захватываем слот, после — освобождаем. Коммитить оффсет при этом надо только после того, как сообщение реально доставлено, поэтому либо ограничиваем «окно» необработанных сообщений, либо коммитим по сплошному префиксу завершённых оффсетов.

Глубже. Простейший корректный вариант — errgroup с лимитом и коммит батча после Wait(): это даёт параллелизм 10 и at-least-once без сложного трекинга оффсетов, ценой того, что при падении переобработается весь текущий батч (значит, приёмник должен быть идемпотентным). Если нужен порядок по ключу — вместо общего пула делаем N шардов-горутин с очередями и распределяем по hash(key) % N. Отдельно упомяните rate limit по RPS (golang.org/x/time/rate) — это другое ограничение, чем concurrency, и интервьюеры любят разницу.

import (
"context"
"github.com/segmentio/kafka-go"
"golang.org/x/sync/errgroup"
)
func run(ctx context.Context, r *kafka.Reader, send func(context.Context, kafka.Message) error) error {
const batchSize = 100
for {
batch := make([]kafka.Message, 0, batchSize)
for len(batch) < batchSize {
m, err := r.FetchMessage(ctx)
if err != nil {
return err
}
batch = append(batch, m)
}
g, gctx := errgroup.WithContext(ctx)
g.SetLimit(10) // не более 10 одновременных запросов во внешнюю систему
for _, m := range batch {
m := m
g.Go(func() error { return send(gctx, m) })
}
if err := g.Wait(); err != nil {
return err // не коммитим: батч переобработается
}
if err := r.CommitMessages(ctx, batch...); err != nil {
return err
}
}
}

Кафка: что такое и зачем нужны партиции, топики? как обеспечить гарантию отправки сообщений?

Заголовок раздела «Кафка: что такое и зачем нужны партиции, топики? как обеспечить гарантию отправки сообщений?»

Коротко. Топик — логическое имя потока, партиции — его физические шарды, дающие параллелизм, масштабирование по брокерам и порядок внутри шарда. Гарантия отправки: acks=all + enable.idempotence=true + retries (по умолчанию максимум) + delivery.timeout.ms, на брокере min.insync.replicas>=2 при RF>=3 и unclean.leader.election.enable=false, а в коде — обязательно проверять ошибку доставки (для асинхронной отправки — в колбэке/канале ошибок), а не «выстрелил и забыл».

Глубже. Зачем партиции, кроме параллелизма: они позволяют топику быть больше одного диска и одного брокера, и они — единица восстановления (реплицируется и переезжает партиция, а не топик). Зачем топики: изоляция схемы данных, retention и прав доступа (ACL выдаются на топик). Про гарантию отправки стоит добавить два практических пункта: (1) при acks=all и падении ниже min.insync.replicas продюсер получит NotEnoughReplicasException — приложение должно уметь это пережить (буфер, backpressure, алерт), а не молча дропнуть; (2) если источник истины — ваша БД, то самая надёжная схема отправки не «пишем в БД и в Kafka», а transactional outbox: в одной транзакции пишем бизнес-данные и строку в outbox, а отдельный процесс (или Debezium) публикует её в Kafka — так исчезает класс проблем «в БД записалось, в Kafka нет».

Коротко. Вопрос про опыт; отвечайте по каркасу из «Работал с кафкой?», делая упор на «как» — то есть на роль Kafka в архитектуре и на конкретные решения.

Глубже. «Как» лучше раскрывать через паттерн, а не через API: event-driven интеграция между сервисами (события домена), CDC/репликация данных из БД, буфер перед медленным приёмником (сглаживание пиков), fan-out одного потока в несколько потребителей (биллинг + аналитика + аудит), лог для реплея и пересчёта. Дальше — одно техническое решение с обоснованием: «ключ = customer_id, чтобы сохранить порядок операций по клиенту; 24 партиции под 8 подов с запасом; ручной коммит после записи в БД; DLQ-топик с headers об ошибке и retention 30 дней».

Коротко. Пометка о наборе тем, а не вопрос. Базовый набор: что такое Kafka и зачем, топик/партиция/оффсет, продюсер/консьюмер/группа, порядок, гарантии доставки, отличия от RabbitMQ.

Глубже. См. выше «Standard Kafka questions» — там каркас рассказа «сверху вниз», которым удобно закрыть весь базовый блок за один заход.

Коротко. См. выше «Чем Kafka отличается от RabbitMQ?». Ключевые оси: лог с retention и реплеем vs очередь с удалением после ack; pull vs push; оффсет группы vs per-message ack; порядок в партиции vs порядок в очереди (и его потеря при нескольких конкурирующих консьюмерах); тупой брокер/умный клиент vs умный брокер с exchange-роутингом.

Глубже. Если просят именно «ключевые», удобно свести к вопросу «кому принадлежит состояние доставки»: в Kafka — потребителю (оффсет), в RabbitMQ — брокеру (unacked-сообщения). Всё остальное — следствия: реплей, множественные независимые подписчики и высокая пропускная способность у Kafka; точечный retry, DLX, TTL, приоритеты и справедливое распределение задач по prefetch — у RabbitMQ.

В Kafka куда отправляем сообщения и откуда забираем?

Заголовок раздела «В Kafka куда отправляем сообщения и откуда забираем?»

Коротко. Отправляем в топик (продюсер выбирает конкретную партицию — по ключу, явно или sticky-батчами), забираем из партиций топика, подписавшись на топик в рамках consumer group. Физически и запись, и чтение идут к брокеру-лидеру нужной партиции.

Глубже. Клиент сам находит, куда идти: сначала Metadata-запрос к любому брокеру из bootstrap.servers, в ответ — список брокеров, партиций и их лидеров; дальше Produce/Fetch отправляются напрямую лидеру. То есть в Kafka нет «точки входа»-прокси, балансировка распределена по клиентам, а bootstrap.servers нужен только для первого запроса метаданных. С Kafka 2.4 появилось чтение с фолловера (client.rack + replica.selector.class) для экономии межзонного трафика, но запись всегда идёт в лидера.

Можем записать сообщение в отдельную партицию?

Заголовок раздела «Можем записать сообщение в отдельную партицию?»

Коротко. Да. Партицию можно задать явно (в Java — конструктор ProducerRecord(topic, partition, key, value), в Go у kafka-go — поле Partition в kafka.Message при записи через Conn/Writer с Balancer), либо косвенно через ключ (одинаковый ключ → всегда одна и та же партиция), либо реализовав собственный partitioner/balancer.

Глубже. Явное указание партиции — мощный, но опасный инструмент: код становится зависимым от текущего числа партиций и от их расположения, при расширении топика логика ломается. Поэтому в 95% случаев правильный способ — ключ. Если key == nil и партиция не указана, современные клиенты используют sticky-партиционирование: пишут в одну партицию, пока не наберётся батч, потом переключаются — так меньше мелких запросов, чем при честном round-robin (KIP-480, с 2.4; в 3.3 добавили partitioner.adaptive.partitioning.enable, учитывающий скорость брокеров).

Как замерить пропускную способность на кафке?

Заголовок раздела «Как замерить пропускную способность на кафке?»

Коротко. Синтетически — штатными kafka-producer-perf-test.sh и kafka-consumer-perf-test.sh (или kafka-e2e-latency), которые дают records/sec, MB/sec и перцентили latency. В проде — по метрикам брокера через JMX: BytesInPerSec, BytesOutPerSec, MessagesInPerSec (per-topic и общие), плюс клиентские метрики продюсера (record-send-rate, request-latency-avg, batch-size-avg) и консьюмера (records-consumed-rate, fetch-latency-avg, records-lag-max).

Глубже. Корректный бенчмарк требует фиксировать условия: размер сообщения, acks, compression.type, linger.ms/batch.size, число партиций, RF, число клиентов и --throughput (при -1 меряется потолок, при фиксированном — latency под нагрузкой). Типичные ошибки: мерить с одним продюсером и одной партицией и делать вывод о кластере; мерить с acks=1, а в проде жить с acks=all; забыть, что при первом прогоне данные ещё в page cache, и «чтение» на самом деле не трогает диск. Для реального ответа на вопрос «выдержим ли пик» полезнее не абсолютный потолок, а headroom: текущий BytesInPerSec против измеренного максимума при тех же настройках, плюс наблюдение за RequestHandlerAvgIdlePercent и UnderReplicatedPartitions — если они деградируют, потолок уже близко.

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

Глубже. Полезно различать несколько «оффсетов» одной партиции: log start offset (начало после retention), log end offset (следующий свободный номер), high watermark (граница реплицированного во все ISR — до неё видят консьюмеры), last stable offset (для read_committed — граница без незавершённых транзакций) и committed offset группы. Оффсеты можно искать не только по номеру: OffsetsForTimes/--to-datetime находит первый оффсет с таймстемпом не меньше заданного, используя .timeindex — это штатный способ «перечитать с 3 часов ночи».

В потоке синхронно приходит количество списаний по пользователю. Как организовать асинхронную работу в данном случае? На что надо обращать внимание в случае использования кафки настройки подписчика и настройки продюсера?

Заголовок раздела «В потоке синхронно приходит количество списаний по пользователю. Как организовать асинхронную работу в данном случае? На что надо обращать внимание в случае использования кафки настройки подписчика и настройки продюсера?»

Коротко. Синхронный API принимает списание, атомарно фиксирует его (БД + outbox) и сразу отвечает клиенту; дальше событие уходит в Kafka с ключом user_id, а обработка идёт асинхронно консьюмером. Ключ user_id критичен: он гарантирует, что все списания одного пользователя лягут в одну партицию и обработаются строго по порядку одним консьюмером.

Глубже. Настройки продюсера: acks=all, enable.idempotence=true (дефолт с 3.0), max.in.flight.requests.per.connection<=5, разумные delivery.timeout.ms/request.timeout.ms, linger.ms 5–20 мс и compression.type=lz4/zstd для батчинга, обязательная обработка ошибки доставки (в синхронном хендлере — либо ждать ack перед ответом 200, либо использовать outbox и отвечать сразу). Настройки консьюмера: enable.auto.commit=false и ручной коммит после успешной записи в БД; max.poll.records и max.poll.interval.ms согласованы со временем обработки пачки (иначе ребаланс посреди работы); isolation.level=read_committed, если продюсер транзакционный; auto.offset.reset=earliest для финансовых данных. И главное — идемпотентность: при at-least-once списание может прийти дважды, поэтому в БД должен быть уникальный ключ операции (transaction_id) и INSERT ... ON CONFLICT DO NOTHING, иначе двойное списание. Плюс контроль баланса нельзя делать «прочитал-посчитал-записал» без блокировки или условного апдейта (UPDATE ... WHERE balance >= amount).

Если consumer упал до того, как сообщение было обработано, как гарантировать повторную обработку этого сообщения? Как Kafka в этом случае управляет offset?

Заголовок раздела «Если consumer упал до того, как сообщение было обработано, как гарантировать повторную обработку этого сообщения? Как Kafka в этом случае управляет offset?»

Коротко. Гарантия одна: не коммитить оффсет до успешной обработки (enable.auto.commit=false + ручной коммит после). Тогда после падения координатор по session timeout отдаст партицию другому члену группы (или тому же после рестарта), тот прочитает последний закоммиченный оффсет из __consumer_offsets и начнёт ровно с необработанного сообщения.

Глубже. Kafka не «возвращает сообщение в очередь» — она вообще ничего не знает про то, обработали вы запись или нет; повтор возникает автоматически просто потому, что позиция не сдвинулась. Из этого следуют три практических вывода. Первый: обработчик обязан быть идемпотентным, потому что переобработается весь диапазон между последним коммитом и падением, а не одно сообщение. Второй: с автокоммитом гарантии нет — оффсет мог уехать вперёд по таймеру. Третий: если сообщение стабильно валит консьюмера, повторная обработка превращается в бесконечный цикл — нужен счётчик попыток и отправка в DLQ/retry-топик с последующим коммитом. Ускорить обнаружение падения помогают session.timeout.ms и корректный graceful shutdown: при аккуратном закрытии консьюмер шлёт LeaveGroup, и ребаланс происходит немедленно, а не через 45 секунд.

Что такое consumer lag в Kafka? Чем он опасен и как его мониторить?

Заголовок раздела «Что такое consumer lag в Kafka? Чем он опасен и как его мониторить?»

Коротко. См. выше определение: logEndOffset - committedOffset. Опасен тем, что данные обрабатываются с растущей задержкой (страдает SLA бизнес-процессов), а при достижении retention сообщения удаляются раньше, чем прочитаны, — это уже безвозвратная потеря. Мониторят через kafka-consumer-groups.sh --describe, JMX-метрику records-lag-max у консьюмера, экспортеры (kafka_exporter, JMX exporter) в Prometheus и Burrow, который оценивает не абсолютный лаг, а тренд.

Глубже. Правильный алерт — не на абсолютное число сообщений (порог зависит от трафика и в пике даст ложные срабатывания), а на производную и на время: «лаг растёт монотонно N минут» и «оценочное время разгребания (lag / скорость потребления) > SLA». Обязательно смотреть лаг по партициям, а не только сумму: ровный рост во всех — не хватает мощности, рост в одной — перекос по ключу или залипший консьюмер. Полезные соседние метрики: скорость коммитов (стоит ли оффсет вообще), частота ребалансов, records-consumed-rate, время обработки. И отдельно — предохранитель: retention топика должен быть заметно больше максимально допустимого времени разгребания лага.

Kafka. Консьюмерная группа одна. Консьюмеров больше партиций. Чем это грозит?

Заголовок раздела «Kafka. Консьюмерная группа одна. Консьюмеров больше партиций. Чем это грозит?»

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

Глубже. Иногда это делают намеренно: «горячий резерв», чтобы при падении активного консьюмера партиция мгновенно перешла к готовому процессу без ожидания старта нового пода. Но чаще это следствие автоскейлинга по CPU, который бесконтрольно поднимает реплики: поды растут, лаг не падает, деньги тратятся. Правильная реакция — масштабировать HPA по лагу (KEDA с Kafka-скейлером), ограничив maxReplicas числом партиций, и/или увеличить число партиций. Ещё один нюанс: простаивающий консьюмер всё равно шлёт heartbeat и участвует в каждом ребалансе, поэтому большая группа «пустых» участников удлиняет ребалансы.

Коротко. Топик — именованный логический поток сообщений, категория, в которую пишут продюсеры и на которую подписываются консьюмеры. Физически состоит из одной или нескольких партиций, каждая из которых реплицируется; топик — единица настройки retention, компакции и прав доступа.

Глубже. Ключевые свойства топика на уровне конфигурации: partitions, replication.factor, retention.ms/retention.bytes, cleanup.policy (delete или compact, можно вместе), min.insync.replicas, max.message.bytes, compression.type. Топик мультиподписной: сколько угодно групп читают его независимо, и чтение не изменяет данные. Отдельно стоит помнить про auto.create.topics.enable: в проде его выключают, потому что опечатка в имени топика иначе создаст топик с дефолтными (обычно неподходящими) настройками.

Чем натс отличается от кафки? Расскажи как работал с кафкой

Заголовок раздела «Чем натс отличается от кафки? Расскажи как работал с кафкой»

Коротко. NATS Core — лёгкий in-memory pub/sub с fire-and-forget (at-most-once), без хранения: нет подписчика — сообщение потеряно; зато микросекундные задержки, wildcard-сабджекты, request-reply и крошечный сервер. NATS JetStream добавляет персистентность, стримы, потребителей с ack по каждому сообщению, дедупликацию по Nats-Msg-Id и work-queue-семантику. Kafka — тяжелее, но заточена под большие объёмы, долгое хранение, реплей и экосистему обработки. Вторая часть вопроса — про личный опыт, отвечайте по каркасу «сценарий → масштаб → библиотека → настройки → решённая проблема».

Глубже. Существенные технические отличия: в Kafka параллелизм и порядок определяются партициями, в JetStream — потребителями и опционально filter_subject, а порядок гарантируется в рамках стрима; в JetStream ack — по каждому сообщению (ack, nak, term, inProgress) с автоматическим redelivery и max_deliver+DLQ-семантикой, чего в Kafka нет; NATS маршрутизирует по иерархическим сабджектам с wildcards (orders.*.created), Kafka — только по имени топика; NATS умеет кластеризацию/суперкластеры и leaf-ноды для edge, что удобно в IoT/микросервисной сетке. Что до эксплуатации — NATS-сервер это один Go-бинарь без внешних зависимостей, Kafka требует существенно больше внимания. Выбор: «шина RPC/сигналов с низкой задержкой» — NATS; «event log с реплеем, аналитикой и большим retention» — Kafka.

Представь что есть в кафке топик с тремя партициями, туда пишутся данные. Ты пишешь приложение со стороны консюмера. Сколько консюмеров нужно задействовать, чтобы наиболее оптимально вычитывать этот топик?

Заголовок раздела «Представь что есть в кафке топик с тремя партициями, туда пишутся данные. Ты пишешь приложение со стороны консюмера. Сколько консюмеров нужно задействовать, чтобы наиболее оптимально вычитывать этот топик?»

Коротко. Три — по одному на партицию: это максимум полезного параллелизма в одной группе. Больше трёх бесполезно (лишние простаивают), меньше — часть партиций достанется одному консьюмеру и он станет узким местом.

Глубже. Оговорки, которые интервьюер хочет услышать. Во-первых, «три» оптимально, если один консьюмер справляется с потоком одной партиции; если нет — надо увеличивать число партиций, а не консьюмеров, либо параллелить обработку внутри процесса пулом воркеров с шардированием по ключу (тогда порядок по ключу сохраняется, но коммит усложняется до отслеживания сплошного префикса). Во-вторых, «консьюмер» — это не обязательно отдельный процесс: три экземпляра Reader в трёх горутинах одного пода тоже займут три партиции. В-третьих, если данные нужны нескольким независимым подсистемам, это разные группы, и счёт для каждой отдельный. И стоит упомянуть резерв: 3 активных + 1 запасной уменьшают время восстановления при падении пода ценой простаивающего процесса.

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

Заголовок раздела «Представь что ты написал консюмер, он читает из кафки какие-то данные, один топик, одна партиция, один консюмер. Нам сказали что копится лаг. Ты зашел в логи консюмера и увидел что там ошибка при анмаршалинге данных (например нам прислали битый json). Нам нужно оперативно пофиксить эту проблему. Как ее решить?»

Коротко. Это классическое «ядовитое сообщение»: консьюмер падает на нём, не коммитит оффсет, перечитывает то же самое и стоит на месте. Оперативное решение — перестать блокироваться на неразбираемом сообщении: ловить ошибку десериализации, логировать (топик/партиция/оффсет/заголовки/сырые байты), отправлять запись в DLQ-топик и коммитить оффсет, продолжая обработку. Как срочный workaround в проде до выката фикса — вручную сдвинуть оффсет группы за проблемную запись (kafka-consumer-groups.sh --reset-offsets --to-offset N+1 --execute при остановленной группе).

Глубже. Важно отделить ошибку данных от ошибки инфраструктуры: битый JSON — невосстановимая ошибка, ретраить её бессмысленно, надо в DLQ; таймаут БД — восстановимая, её надо ретраить с backoff и не коммитить. Смешивать их в одном catch — типичная причина либо бесконечных циклов, либо тихой потери данных. Долгосрочные меры: Schema Registry с Avro/Protobuf и валидацией на стороне продюсера, чтобы битые сообщения не появлялись вовсе; в DLQ класть исходные байты плюс headers с причиной, стектрейсом, временем и оригинальным оффсетом, чтобы можно было переиграть после фикса; метрика и алерт на количество сообщений в DLQ. И не забыть: сдвиг оффсета вручную возможен только когда группа неактивна (иначе координатор откажет), а сама операция — потеря данных, поэтому сначала сохраните проблемную запись (kafka-console-consumer --partition 0 --offset N --max-messages 1).

Допустим у нас сервис читает из топика сообщения, и на каждое сообщение пишет в базу что-то. У нас стало слишком много сообщений, и база не справляется. Как решить проблему на уровне сервиса?

Заголовок раздела «Допустим у нас сервис читает из топика сообщения, и на каждое сообщение пишет в базу что-то. У нас стало слишком много сообщений, и база не справляется. Как решить проблему на уровне сервиса?»

Коротко. Перестать писать по одной строке на сообщение: копить пачку из poll и делать батчевый INSERT/COPY на 500–5000 строк в одной транзакции, коммитя оффсет после успешной записи. Это обычно даёт кратный выигрыш без изменения БД. Дополнительно — дедупликация и агрегация в памяти перед записью (если несколько сообщений по одному ключу схлопываются в одно состояние), ограничение конкурентности пишущих горутин, чтобы не выбивать пул соединений, и backpressure вместо неограниченных буферов.

Глубже. Порядок оптимизаций от дешёвого к дорогому: (1) батчинг + pgx.CopyFrom / multi-values INSERT ... ON CONFLICT; (2) уменьшить работу на строку — убрать лишние индексы, триггеры, вынести аналитические поля; (3) агрегировать в окне (например, писать не каждое изменение баланса, а итог за 200 мс по ключу) — работает, если бизнес допускает потерю промежуточных состояний; (4) ограничить параллелизм так, чтобы число одновременных транзакций совпадало с числом соединений в пуле (больше — только рост latency и блокировок); (5) шардировать запись по hash(key) между воркерами, чтобы уменьшить конкуренцию за одни и те же строки и дедлоки; (6) если поток фундаментально больше, чем БД может принять, — менять хранилище/паттерн (ClickHouse для аналитики, отдельная таблица-очередь с последующей асинхронной свёрткой) или явно вводить сэмплирование. Важно проговорить, что «увеличить число консьюмеров» здесь — плохой первый ответ: узкое место в БД, и больше консьюмеров сделает только хуже.

// батчевая запись: одна транзакция на пачку, коммит оффсета после успеха
msgs := make([]kafka.Message, 0, 1000)
flushAt := time.Now().Add(200 * time.Millisecond)
for {
ctxFetch, cancel := context.WithDeadline(ctx, flushAt)
m, err := r.FetchMessage(ctxFetch)
cancel()
if err == nil {
msgs = append(msgs, m)
if len(msgs) < cap(msgs) {
continue
}
} else if !errors.Is(err, context.DeadlineExceeded) {
return err
}
if len(msgs) > 0 {
if err := insertBatch(ctx, db, msgs); err != nil { // один COPY/INSERT
return err
}
if err := r.CommitMessages(ctx, msgs...); err != nil {
return err
}
msgs = msgs[:0]
}
flushAt = time.Now().Add(200 * time.Millisecond)
}

Какие библиотеки использоавл для работы с Kafka?

Заголовок раздела «Какие библиотеки использоавл для работы с Kafka?»

Коротко. В Go четыре живых варианта: github.com/segmentio/kafka-go (чистый Go, простой API Reader/Writer), github.com/confluentinc/confluent-kafka-go (обёртка над librdkafka через cgo, самая полная поддержка фич и транзакций), github.com/IBM/sarama (бывший Shopify/sarama, чистый Go, много кода в проде, API многословный) и github.com/twmb/franz-go (чистый Go, современный, поддерживает транзакции и новые протоколы).

Глубже. Практические критерии выбора: нужен cgo или нет (в Alpine/дистролесс и при кросс-компиляции librdkafka — боль); нужны ли транзакции и EOS (confluent-kafka-go и franz-go — да; kafka-go транзакции продюсера не поддерживает); насколько важна производительность батчинга (librdkafka и franz-go традиционно быстрее). Если вопрос про экосистему — упомяните Kafka Connect для интеграций без кода, Debezium для CDC, Schema Registry + srclient, и testcontainers-go (или kafka в docker-compose / встроенный kfake у franz-go) для интеграционных тестов. Отвечать надо честно: назвали библиотеку — будьте готовы к вопросу, как в ней делается ручной коммит и обработка ребаланса.

Коротко. Consumer group решает две задачи: горизонтальное масштабирование обработки (партиции эксклюзивно распределяются между членами группы, каждое сообщение обрабатывается одним из них) и хранение общей позиции чтения с автоматическим перераспределением партиций при падении или добавлении участников. Плюс разные группы дают fan-out: один топик читают независимо биллинг, аналитика и аудит.

Глубже. Технически группа — это ещё и единица отказоустойчивости и «единица идентичности» приложения: group.id определяет, где лежат оффсеты, поэтому смена group.id = чтение с нуля или с конца (по auto.offset.reset) — этим иногда пользуются намеренно для реплея, а иногда случайно ломают прод. Координатор группы отслеживает членство по heartbeat, инициирует ребалансы и валидирует коммиты (коммит от «зомби»-участника с устаревшим generation id будет отвергнут — это защита от двойной обработки при затянувшемся GC-паузе).

Если есть три партиции на один топик и мы отправляем туда сообщение, то в какое кол-во партиций это сообщение попадет?

Заголовок раздела «Если есть три партиции на один топик и мы отправляем туда сообщение, то в какое кол-во партиций это сообщение попадет?»

Коротко. Ровно в одну. Сообщение не копируется по партициям: партиции — это шарды, а не реплики. Копии создаёт репликация (replication.factor), но это копии партиции на других брокерах, а не отдельные партиции.

Глубже. В какую именно — определяет партиционер: если задан ключ, murmur2(key) % 3; если ключа нет, современный клиент использует sticky-батчинг (пишет пачку в одну партицию, потом переключается); если партиция указана явно — в неё. Отсюда типичное продолжение вопроса: «а как сделать, чтобы сообщение получили все консьюмеры?» — не дублированием по партициям, а разными consumer group.

Коротко. См. выше «Чем натс отличается от кафки?»: NATS Core — минималистичный at-most-once pub/sub с сабджект-роутингом и очень низкой задержкой, без хранения; JetStream добавляет персистентность, per-message ack, redelivery и дедупликацию. Kafka — партиционированный durable-лог с оффсетами, реплеем и тяжёлой, но богатой экосистемой.

Глубже. Если нужно одно предложение-отличие: в Kafka потребитель управляет позицией в неизменяемом логе, а в NATS (JetStream) брокер управляет доставкой каждого сообщения и его подтверждением — это ближе к RabbitMQ по модели, чем к Kafka, при том что хранение у JetStream тоже есть.

Коротко. Физические. Партиция — реальный набор файлов на диске брокера: каталог <topic>-<n>/ с сегментами .log, индексами .index/.timeindex. Логическая абстракция — это топик, который сам по себе на диске не существует.

Глубже. Из «физичности» следуют практические ограничения: каждая партиция и каждый её сегмент — это открытые файловые дескрипторы и отдельные структуры в памяти брокера; партиция целиком лежит на одном брокере (и на одном диске/логдире), поэтому она не может быть больше диска и не может параллелиться внутри себя; перебалансировка данных между брокерами — это физическое копирование партиций (kafka-reassign-partitions.sh). Десятки тысяч партиций на кластер — реальная эксплуатационная проблема (время выборов лидеров, метаданные, recovery после рестарта), хотя KRaft заметно поднял потолок по сравнению с ZooKeeper.

Коротко. Брокеры (хранят партиции, обслуживают запросы), контроллер (KRaft-кворум, с 3.3; раньше — ZooKeeper) для метаданных и выборов лидеров, продюсеры и консьюмеры как клиенты, group coordinator и transaction coordinator как роли брокеров, служебные топики __consumer_offsets и __transaction_state. Вокруг ядра — Kafka Connect, Kafka Streams, Schema Registry, MirrorMaker 2.

Глубже. Полезно уметь объяснить, кто за что отвечает при обычной записи: клиент берёт метаданные у любого брокера → пишет лидеру партиции → лидер аппендит в активный сегмент и ждёт репликации от ISR → отвечает ack → консьюмер фетчит у того же лидера до high watermark → коммитит оффсет своему group coordinator. Контроллер в этом пути не участвует вообще — он нужен только при изменениях топологии (создание топиков, падение брокера, смена лидера), и это одна из причин, почему Kafka хорошо масштабируется.

Коротко. Топик — логическое имя потока и единица настройки/подписки; партиция — физическая упорядоченная часть этого потока на конкретном брокере, единица параллелизма, порядка и репликации. Один топик = 1..N партиций; порядок существует в партиции, а не в топике.

Глубже. Различие проявляется во всех рабочих вопросах: подписываются на топик, но назначаются на партиции; retention настраивается на топике, но применяется к сегментам каждой партиции; оффсет уникален в партиции, а не в топике; ACL выдаются на топик, лидерство — у партиции. Ещё один частый вопрос-следствие: «сколько партиций делать» — ориентир от целевой пропускной способности и максимума консьюмеров с запасом на рост, помня, что уменьшить нельзя, а увеличение ломает привязку ключей.

Коротко. См. выше: at-most-once, at-least-once, exactly-once. По умолчанию правильно настроенная система даёт at-least-once; exactly-once достигается транзакциями Kafka (read-process-write внутри Kafka) или идемпотентностью приёмника при работе с внешними системами.

Глубже. Единственное, что стоит добавить к предыдущим ответам: гарантия — это свойство всей цепочки, а не Kafka. Если продюсер отправляет «fire and forget» и не проверяет ошибку, никакая настройка брокера не спасёт; если консьюмер коммитит до обработки, acks=all бесполезен. Поэтому на собеседовании отвечайте парой «настройка продюсера + место коммита у консьюмера».

Коротко. См. выше «Для чего нужна консьюмер группа в кафке?»: набор консьюмеров с общим group.id, между которыми координатор эксклюзивно делит партиции; группа хранит собственные оффсеты и переживает падения через ребаланс.

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

Ты говоришь, что работал с кафкой и натс. Можешь рассказать какие-то особенности? Почувствовал ли разницу какую-то или одинаково?

Заголовок раздела «Ты говоришь, что работал с кафкой и натс. Можешь рассказать какие-то особенности? Почувствовал ли разницу какую-то или одинаково?»

Коротко. Вопрос-проверка: интервьюер хочет услышать, что вы понимаете разницу моделей, а не просто «использовал оба клиента». Ответ строится как «разные модели → разные последствия в коде и эксплуатации», с конкретными примерами из вашего опыта.

Глубже. Каркас хорошего ответа: (1) разница в модели — лог с оффсетами против доставки с ack, отсюда реплей и множественные подписчики у Kafka против точечного redelivery и max_deliver у JetStream; (2) разница в коде — в Kafka вы думаете про ключ партиционирования, ребалансы и момент коммита; в NATS — про сабджекты, ack/nak и durable-консьюмеров; (3) разница в эксплуатации — NATS ставится одним бинарём и почти не требует ухода, Kafka требует мониторинга ISR, лага, ребалансов и планирования партиций; (4) разница в latency и throughput — NATS Core даёт микросекунды на сигналах, Kafka выигрывает на больших объёмах с батчингом; (5) вывод «где что применял и почему». Типичная ошибка — сказать «примерно одинаково, обе очереди»: это моментально снижает доверие ко всему предыдущему рассказу.

Коротко. См. выше «Чем Kafka отличается от RabbitMQ?». Кратко: лог с retention и оффсетами против очереди с удалением после ack; pull против push; порядок в партиции против роутинга и per-message ack; масштаб и реплей против гибкой маршрутизации, приоритетов, TTL и DLX.

Глубже. Если хотят «когда что выбрать»: Kafka — потоки событий, аналитика, интеграция многих потребителей, необходимость перечитать историю, десятки–сотни МБ/с. RabbitMQ — распределение задач воркерам, RPC/request-reply, сложный роутинг по правилам, отложенная доставка и приоритеты, умеренные объёмы. Если и то и другое — часто ставят оба, а не пытаются натянуть один на все сценарии.

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

Глубже. Разворачивайте по каркасу: (1) зачем — развязать продюсеров и потребителей, буферизовать пики, дать нескольким подсистемам один поток; (2) как устроено — топик/партиция/оффсет, сегменты на диске, последовательная запись, page cache и zero-copy; (3) надёжность — RF, ISR, acks, min.insync.replicas, KRaft-контроллер вместо ZooKeeper с 3.3 (в 4.0 ZooKeeper удалён); (4) потребление — группы, координатор, ребалансы, коммиты в __consumer_offsets; (5) гарантии — at-least-once по умолчанию, идемпотентный продюсер по умолчанию с 3.0, транзакции для EOS; (6) эксплуатация — lag, DLQ, выбор ключа и числа партиций. Такой рассказ на 3–4 минуты обычно снимает половину следующих вопросов.

Как настроить Kafka, чтобы все 10 реплик сервиса получали и обрабатывали одни и те же сообщения из топика?

Заголовок раздела «Как настроить Kafka, чтобы все 10 реплик сервиса получали и обрабатывали одни и те же сообщения из топика?»

Коротко. Дать каждой реплике свой уникальный group.id (например, svc-cache-<pod_name>): тогда каждая реплика — отдельная группа, получает все партиции и весь поток целиком. Альтернатива без групп — вручную назначить все партиции через Assign/Partition и самому управлять начальной позицией.

Глубже. Это классический паттерн «broadcast» — обновление локального кеша или конфигурации во всех подах. Практические детали: уникальный group.id на под порождает по группе на реплику, и после пересоздания подов остаются «мусорные» группы с коммитами (истекут по offsets.retention.minutes, но в UI будут висеть); поэтому часто вообще отказываются от коммитов и на старте делают Seek в конец (auto.offset.reset=latest + новый group.id при каждом старте) — для кеша важно свежее состояние, а не история. Если нужен полный снапшот при старте — топик делают компактящимся (cleanup.policy=compact) и читают с начала. Чего делать не надо: пытаться добиться broadcast увеличением партиций или ставить 10 подов в одну группу — в группе каждое сообщение достанется ровно одному.

Как вы обрабатываете ошибки при потреблении сообщений из Kafka?

Заголовок раздела «Как вы обрабатываете ошибки при потреблении сообщений из Kafka?»

Коротко. Классифицируем ошибку: восстановимая (таймаут БД, 5xx внешнего API) — ретраим с экспоненциальным backoff, оффсет не коммитим; невосстановимая (битый формат, невалидные данные, бизнес-отказ) — логируем с контекстом, отправляем в DLQ-топик и коммитим оффсет, чтобы не блокировать партицию. Плюс ограничение числа попыток, метрики и алерты на DLQ, идемпотентный обработчик, потому что повторы неизбежны.

Глубже. Развёрнутая схема, которую хорошо описать на собеседовании: retry-топики с разными задержками (orders.retry.5s, orders.retry.1m, orders.retry.10m) — сообщение уходит туда, консьюмер retry-топика ждёт до нужного времени и возвращает в основной поток; после N попыток — в orders.dlq. В headers кладём x-retry-count, x-original-topic, x-original-offset, x-error. Важные тонкости: ретрай «на месте» без коммита блокирует всю партицию, поэтому долгие ретраи всегда выносятся в отдельный топик; повторная обработка ломает порядок по ключу, и если он критичен — приходится останавливать партицию, а не откладывать сообщение; при длинных ретраях внутри poll-цикла легко превысить max.poll.interval.ms и словить ребаланс, поэтому либо pause()/resume() на партиции, либо короткие итерации. И, разумеется, DLQ без процесса разбора — это просто мусорка: нужны дашборд, алерт и инструмент для переигрывания.

Можно ли настроить Kafka как умного брокера и глупого клиента?

Заголовок раздела «Можно ли настроить Kafka как умного брокера и глупого клиента?»

Коротко. Нет, архитектурно Kafka — «глупый брокер, умный клиент»: брокер не маршрутизирует по содержимому, не фильтрует, не хранит состояние доставки каждого сообщения и не делает redelivery; вся логика (партиционирование, коммиты, ретраи, дедупликация) — на клиенте. Настройками это не переключается.

Глубже. Это сознательный дизайн: минимальная работа на брокере — причина высокой пропускной способности (последовательная запись + zero-copy + отсутствие per-message-состояния). Приблизиться к «умному брокеру» можно только надстройками: Kafka Streams/ksqlDB (обработка и роутинг, но выполняются в вашем приложении или отдельном кластере, не на брокере), Kafka Connect (интеграция без кода), брокерные интерцепторы/плагины и, начиная с Kafka 4.0, share groups (KIP-932, «Queues for Kafka») — они добавляют очередную семантику с per-message-подтверждением и без привязки «одна партиция — один консьюмер», но это молодая функциональность в статусе early access. Если по требованиям нужен именно умный брокер — правильный ответ на собеседовании: «взять RabbitMQ, а не выкручивать Kafka».

Почему использовал кафку, а не более легковесное решение?

Заголовок раздела «Почему использовал кафку, а не более легковесное решение?»

Коротко. Ответ должен быть про требования, а не про моду: нужен ли реплей истории и подключение новых потребителей задним числом, нужен ли fan-out на несколько независимых подсистем, какие объёмы и retention, нужен ли строгий порядок по ключу при горизонтальном масштабировании, есть ли уже экосистема (Connect/Debezium/аналитика) и экспертиза в команде. Если ни одно из этих требований не выполняется — честно признать, что легковесного решения (NATS, RabbitMQ, а иногда просто таблица-очередь в Postgres) хватило бы.

Глубже. Сильный ответ выглядит так: «Нам был нужен реплей: аналитика периодически пересчитывает витрины за 30 дней, а новый сервис при запуске проигрывает историю с начала — это сразу отсекает RabbitMQ и NATS Core. Плюс поток ~50k msg/s с пиками x3 и три независимых потребителя. Отдельно взвешивали Postgres-очередь — не подошла из-за объёма и из-за того, что она бы конкурировала за IO с основной нагрузкой». Типичные ошибки: «Kafka — стандарт индустрии» без требований; «она быстрее» без цифр; отрицать издержки — правильно упомянуть, что за это заплатили эксплуатацией, ребалансами и необходимостью думать о партиционировании.

  • Называть Kafka «очередью» и утверждать, что сообщение удаляется после чтения: удаление происходит только по retention или компакции, а чтение неразрушающее.
  • Говорить, что Kafka гарантирует глобальный порядок в топике. Порядок есть только внутри партиции, и он ломается при смене числа партиций или неаккуратных ретраях.
  • Путать масштабирование внутри группы с fan-out: «добавим ещё консьюмеров в ту же группу, чтобы все получили сообщение» — в группе сообщение достанется ровно одному.
  • Считать, что число консьюмеров можно наращивать бесконечно: потолок — число партиций, лишние простаивают.
  • Заявлять exactly-once как галочку в конфиге, не умея объяснить транзакции, transactional.id, read_committed и то, что EOS не распространяется на внешнюю БД.
  • Считать enable.auto.commit=true семантикой at-most-once: на деле это неуправляемый at-least-once, а при асинхронной обработке — ещё и источник потерь.
  • Не различать восстановимые и невосстановимые ошибки при потреблении: бесконечный ретрай битого JSON останавливает партицию и растит лаг.
  • Обещать уменьшение числа партиций: Kafka умеет только увеличивать, «уменьшение» — это новый топик и перекладывание данных.
  • Забывать про min.insync.replicas: acks=all при RF=3 и min.insync.replicas=1 не защищает от потери данных.
  • Отвечать на вопрос про опыт общими словами («да, читал и писал») без масштаба, настроек и хотя бы одной решённой проблемы.