Самая важная фраза:
Kafka хранит упорядоченные последовательности записей в partition. Consumer не забирает сообщение из Kafka. Он читает log с некоторой позиции.
Kafka концептуально ближе к database commit log, чем к классическому message broker. Запись после чтения consumer’ом не исчезает. Consumer хранит позицию чтения, то есть offset, и поэтому может вернуться назад и перечитать данные. (Apache Kafka)
Модель:
Topic: payments
Partition 0:
offset 0 -> PaymentCreated
offset 1 -> PaymentConfirmed
offset 2 -> PaymentFailed
Partition 1:
offset 0 -> PaymentCreated
offset 1 -> PaymentConfirmed
То есть:
Kafka
└── Topic
├── Partition 0 -> ordered log
├── Partition 1 -> ordered log
└── Partition 2 -> ordered log
Offset имеет смысл только внутри partition.
Не существует общего:
Kafka offset = 10500
Есть:
payments-0 offset 10500
payments-1 offset 7821
payments-2 offset 9211
Вот это фундамент.
2. Topic и partition
Topic является логическим потоком данных:
payments
orders
users
notifications
Но физическая единица параллелизма Kafka это partition.
Например:
orders
partition 0
partition 1
partition 2
partition 3
Каждая partition является упорядоченным log. Consumer groups масштабируются через распределение partition между consumers. Именно partition одновременно даёт Kafka порядок внутри одного log и возможность параллельного чтения разных log. (Apache Kafka)
Отсюда важный вывод:
Количество partitions ограничивает параллелизм классической consumer group.
Есть:
3 partitions
10 consumers
В классической consumer group эффективно читать будут максимум 3 consumer’а.
P0 -> C1
P1 -> C2
P2 -> C3
C4 idle
C5 idle
...
На собеседовании:
Как масштабировать Kafka consumer?
Ответ:
Добавлять consumer instances внутри одной consumer group можно до уровня полезного параллелизма, определяемого количеством partitions. Если consumers больше partitions, дополнительные consumers не увеличивают параллелизм классической группы.
3. Key определяет partition и тем самым порядок
Producer отправляет:
key
value
headers
Например:
key = userId
value = UserUpdated
При стандартной partitioning logic наличие key приводит к выбору partition на основе hash key. Без явно указанной partition и без key producer использует sticky partitioning для улучшения batching. (Apache Kafka)
Упрощённо:
partition = hash(key) % partitions
Поэтому:
userId = 100
события:
UserCreated
UserUpdated
UserDeleted
при стабильной схеме partitioning попадают в одну partition.
И тогда:
offset 10 -> UserCreated
offset 11 -> UserUpdated
offset 12 -> UserDeleted
Порядок сохраняется.
Ключевой вопрос:
Kafka гарантирует порядок сообщений?
Правильный ответ:
Kafka гарантирует порядок внутри partition. Глобального порядка между partitions нет.
И дальше:
Если нужен порядок событий одной сущности, я выбираю entity ID в качестве message key, чтобы события этой сущности маршрутизировались в одну partition.
Например:
key = accountId
key = orderId
key = userId
Ловушка
Ты увеличил partitions:
4 -> 8
Если partition выбирается через hash modulo partition count, соответствие:
key -> partition
может измениться.
То есть увеличение количества partitions не является совершенно нейтральной операцией для систем, где архитектурный смысл завязан на стабильное распределение ключей.
4. Producer
Producer не обязательно делает:
send message
wait
send message
wait
Kafka producer batching’ует records, причём batches формируются для конкретных partitions. batch.size задаёт верхнюю границу размера batch, а linger.ms позволяет немного подождать накопления записей; начиная с Kafka 4.0 default linger.ms равен 5 ms. (Apache Kafka)
Модель:
Application
|
KafkaProducer
|
Record accumulator
|
+-- batch P0
+-- batch P1
+-- batch P2
|
Broker
Нужно знать:
batch.size
linger.ms
compression.type
acks
retries
delivery.timeout.ms
enable.idempotence
batch.size
Размер batch.
Больше batch:
better throughput
better compression
Но потенциально:
more buffering
more memory
linger.ms
Producer может немного подождать:
record
record
record
record
↓
batch
вместо:
record -> request
record -> request
record -> request
Kafka официально описывает linger.ms как верхнюю границу искусственной задержки для накопления batch. (Apache Kafka)
compression
Обычно нужно понимать:
gzip
snappy
lz4
zstd
Не надо помнить сравнительную таблицу наизусть.
Нужно понимать принцип:
Kafka эффективно работает с batch compression. Компрессия уменьшает network и storage pressure ценой CPU.
5. Consumer работает по pull model
Вот здесь на прошлом разговоре было важное различение.
Kafka не push’ит сообщения consumer’у.
Consumer делает fetch/poll и запрашивает данные после своей текущей позиции. Pull model позволяет медленному consumer’у просто отстать, а потом догнать log; Kafka также использует long polling и batching при чтении. (Apache Kafka)
Условно:
Consumer:
Give me records after offset 105
Kafka:
106
107
108
109
Потом:
Consumer:
Give me records after offset 109
На Java:
consumer.poll(...)
Концептуально:
consumer -> broker: fetch
broker -> consumer: records
Не:
broker -> PUSH -> consumer
Это тебе особенно важно не перепутать.
6. Consumer group
Допустим:
topic orders
4 partitions
Есть сервис:
OrderProcessor
Запускаем 2 instances:
OrderProcessor-1
OrderProcessor-2
Оба:
group.id = order-processor
Kafka распределяет partitions:
C1:
P0
P1
C2:
P2
P3
В классической consumer group каждую partition в конкретный момент обрабатывает один consumer группы. Это и позволяет представлять position группы как offset следующей записи для каждой partition. (Apache Kafka)
Запускаем ещё:
C3
Может произойти:
C1 -> P0
C2 -> P1, P2
C3 -> P3
Это rebalance.
7. Разные consumer groups читают одни и те же события независимо
Например:
Topic: OrderCreated
Есть:
InventoryService
PaymentService
AnalyticsService
NotificationService
Группы:
inventory
payments
analytics
notifications
Каждая group имеет свою позицию чтения.
inventory offset 1000
payments offset 995
analytics offset 800
notifications offset 1000
То есть одно событие:
OrderCreated
может независимо обработать:
Inventory
Payment
Analytics
Notification
Вот почему Kafka прекрасно подходит под event-driven architecture.
Не нужно:
Producer -> Inventory
Producer -> Payment
Producer -> Analytics
Producer -> Notification
Есть:
Producer
|
OrderCreated topic
|
+-- Inventory group
+-- Payment group
+-- Analytics group
+-- Notification group
8. Offset
Это ещё один фундамент.
Consumer имеет две разные по смыслу позиции.
Условно:
current position
committed offset
Consumer уже прочитал:
100
101
102
103
104
Его текущая позиция:
105
Но committed offset может быть:
100
Consumer падает.
После восстановления consumer group ориентируется на committed position и может снова получить ранее прочитанные records. Именно поэтому Kafka позволяет replay и почему commit strategy напрямую определяет delivery semantics. (Apache Kafka)
Принцип:
fetch != processed
processed != committed
Вот три разных события.
И интервьюеры любят проверять, различаешь ли ты их.
9. At-most-once и at-least-once
At-most-once
1. receive
2. commit offset
3. process
Consumer:
received message 100
committed 101
CRASH
Kafka считает:
next = 101
Но 100 не был обработан.
Получаем:
message loss
Kafka docs прямо описывают commit before processing как способ построить at-most-once semantics. (Apache Kafka)
At-least-once
1. receive
2. process
3. commit
Например:
receive PaymentCreated
write DB
CRASH
commit не произошёл
После restart:
PaymentCreated
будет прочитан снова.
Получаем:
duplicate processing
Kafka по умолчанию ориентирована на at-least-once delivery semantics. (Apache Kafka)
Поэтому:
At-least-once требует idempotent consumer или deduplication strategy.
Например:
event_id = UUID
DB:
processed_events
----------------
event_id UNIQUE
Получили событие:
if event_id already processed:
skip
10. Exactly once: здесь нельзя говорить магические слова
Вопрос:
Kafka поддерживает exactly once?
Плохой ответ:
Да.
Нормальный ответ:
Нужно определить границы exactly-once semantics.
Kafka обеспечивает exactly-once processing для сценария:
Kafka
↓
process
↓
Kafka
через transactions, transactional producer, transactional update consumer offsets и read_committed. Kafka может атомарно связать output records и изменение consumer position. (Apache Kafka)
Например:
Topic A
↓
Consumer
↓
transform
↓
Topic B
В transaction:
produce Topic B
commit consumed offsets
атомарно.
Но:
Kafka
↓
Consumer
↓
PostgreSQL
Kafka transaction не превращает PostgreSQL transaction в часть Kafka transaction.
Официальная документация прямо выделяет внешнюю систему как отдельную проблему: нужно координировать consumer position с фактически сохранённым output, а exactly-once для внешних destinations требует сотрудничества с этой системой. (Apache Kafka)
Вот здесь говоришь:
transactional outbox
idempotency key
deduplication
unique constraint
store offset with output
Это сильный ответ.
11. Idempotent producer
Представим:
Producer -> Broker
Broker записал record.
ACK потерялся.
Producer думает:
не записалось
Retry:
send again
Без защиты потенциально:
message
message
Idempotent producer защищает producer retry path от записи дубликатов. В актуальной конфигурации Kafka idempotence включена по умолчанию, если ей не противоречат другие producer settings; она требует совместимых параметров, включая acks=all. (Apache Kafka)
Важно:
Producer idempotence не означает, что твой business consumer стал idempotent.
Это две разные границы.
Producer retry duplication
и:
Consumer processes same event twice
не одно и то же.
12. Replication
Partition имеет replicas.
Например:
Partition 0
Broker 1 -> Leader
Broker 2 -> Follower
Broker 3 -> Follower
Producer пишет leader.
Producer
↓
Leader
↓
Followers replicate
Для partition один broker является leader, а остальные replicas поддерживают копии log. Kafka использует in-sync replicas, ISR, в своей replication model. (Apache Kafka)
Нужно знать:
replication.factor
leader
follower
ISR
acks
min.insync.replicas
13. ISR
ISR:
In-Sync Replicas
Условно:
Leader B1 offset 1000
Follower B2 offset 1000
Follower B3 offset 999
Это не означает автоматически, что B3 моментально исключён.
ISR определяется Kafka по состоянию replication lag и соответствующим временным условиям.
Концептуально тебе надо понимать:
ISR это replicas, считающиеся достаточно синхронизированными с leader для участия в модели надёжности.
14. acks
acks=0
Producer не ждёт acknowledgement.
Producer -> send
Producer -> дальше
Kafka не может сообщить producer’у о большинстве ошибок записи, а retries в обычном смысле не работают, потому что client не знает о failure. (Apache Kafka)
Высокая скорость.
Слабая гарантия.
acks=1
Leader записал:
ACK
Но followers ещё потенциально не replicated.
Leader падает.
Возможна потеря записи.
acks=all
Leader ждёт acknowledgement от полного текущего набора ISR. Apache Kafka называет это strongest available producer acknowledgement guarantee. (Apache Kafka)
Но здесь появляется:
min.insync.replicas
Например:
replication.factor = 3
min.insync.replicas = 2
acks = all
Это классическая схема durability: запись считается успешной только при выполнении ISR requirement; при недостаточном количестве in-sync replicas producer получает ошибку вместо успешной записи. (Apache Kafka)
На интервью:
replication.factor=3,acks=all,min.insync.replicas=2означает, что я предпочитаю отказ записи потере durability guarantee, если число доступных ISR падает ниже требуемого уровня.
Вот это звучит архитектурно.
15. Retention
Kafka не удаляет record потому, что consumer его прочитал.
Удаление определяется retention policy.
Условно:
retention.ms
retention.bytes
То есть:
consumer read -> nothing deleted
Kafka может иметь:
consumer A at offset 100
consumer B at offset 10000
Records живут независимо от позиции consumer.
Это принципиальное отличие от традиционной queue model.
16. Log compaction
Есть topic:
user-state
Records:
key=1 value=Igor v1
key=2 value=Alex v1
key=1 value=Igor v2
key=1 value=Igor v3
При log compaction Kafka стремится сохранить последнее значение для каждого key в compacted portion log. Log compaction является отдельным механизмом Kafka и предназначена для сохранения актуального состояния по key, а не для удаления records после consumer acknowledgement. (Apache Kafka)
После compaction концептуально:
key=2 Alex v1
key=1 Igor v3
Это удобно для:
current user state
configuration
entity state
cache rebuild
CDC-derived state
Удаление key:
key=1
value=null
То есть tombstone.
Важное различение
Retention:
delete old data by time/size
Compaction:
retain latest value by key
Это разные оси.
17. Rebalance
Consumer group:
C1 -> P0 P1
C2 -> P2 P3
C3 присоединяется.
Assignment надо изменить.
C1 -> P0
C2 -> P1 P2
C3 -> P3
Это rebalance.
Причины могут включать:
consumer joins
consumer leaves
consumer dies
subscription changes
partition count changes
Классическая проблема rebalance:
processing disruption
partition reassignment
state movement
possible reprocessing around committed offsets
В Kafka 4.0 новый consumer rebalance protocol KIP-848 получил статус GA; это актуальный 4.x контекст, хотя на обычном собеседовании чаще проверят саму природу partition assignment и rebalance, а не номер KIP. (Apache Kafka)
Тебе нужно знать настройки концептуально:
session timeout
heartbeat
max.poll.interval.ms
Особенно max.poll.interval.ms.
Сценарий:
poll()
↓
processing takes 20 minutes
Kafka может решить:
consumer isn't progressing correctly
и перераспределить partitions.
А старый consumer ещё что-то обрабатывает.
Начинается веселье с ножами 🎪🔪
Поэтому long processing нужно проектировать осознанно.
18. Consumer lag
Это одна из главных operational metrics.
Условно:
log end offset = 100000
consumer offset = 90000
Lag:
10000
То есть consumer отстал на 10 000 records.
Но:
Lag в records сам по себе не говорит, насколько всё плохо.
Потому что:
10 000 tiny records
и:
10 000 expensive ML jobs
это совершенно разная реальность.
Нужно смотреть:
lag trend
processing rate
produce rate
time lag
partition skew
consumer errors
rebalance rate
Вопрос:
Что будете мониторить в Kafka?
Я бы отвечал:
Consumer lag по partition и его динамику, under-replicated/offline replica состояние, ISR degradation, produce/fetch throughput, request latency и ошибки, disk usage, broker load distribution и rebalance behaviour.
Kafka официально предоставляет broker/client metrics, включая partition и replica-related показатели. (Apache Kafka)
19. Hot partition
Допустим:
key = country
Traffic:
Russia 70%
USA 15%
Germany 5%
other 10%
Получаем:
partition 3 = Russia = 🔥🔥🔥
Другие partitions:
zzz
Это partition skew.
Kafka физически распределена, но плохой key может создать логическую централизацию нагрузки.
На system design надо спросить:
Какая cardinality key?
Как распределён traffic?
Есть heavy hitters?
Нужен ли strict per-key ordering?
Иногда делают:
key = userId
вместо:
key = countryId
Или sharded key:
tenantId + bucket
Но sharding может уничтожить strict ordering.
То есть trade-off:
ordering <-> parallelism
Вот это Kafka почти в чистом виде.
20. Kafka против RabbitMQ
Не надо отвечать:
Kafka для больших нагрузок
RabbitMQ для маленьких
Это детский ответ.
Различие модели.
Kafka
persistent distributed log
offset-based consumption
replay
multiple independent consumer groups
partition-based ordering
event streaming
Traditional queue model
message delivery
ack
message lifecycle tied to processing
per-message work distribution
И вот здесь важная поправка к нашему прошлому разговору.
В современной Kafka появился Share Consumer
Начиная с Kafka 4.2, Queues for Kafka / share groups production-ready. Share consumers могут совместно обрабатывать records без классического правила «одна partition назначена одному consumer группы»; partitions могут обслуживаться несколькими share consumers, records acknowledge’ятся индивидуально, а Kafka считает delivery attempts. Apache прямо рекомендует share groups для случаев, где records обрабатываются по одному как work items, а не как ordered stream. (Apache Kafka)
То есть наш старый вопрос:
«Я хочу, чтобы один worker взял необработанную job, а другой её не взял».
Исторически классический Kafka consumer group решал это на уровне partition ownership и offsets, поэтому это была не настоящая per-message queue semantics.
В Kafka 4.2 share groups уже специально закрывают именно этот разрыв.
Там буквально модель:
record acquired by C1
↓
temporary acquisition lock
↓
C2 cannot acquire it
C1 может:
acknowledge
release
reject
renew
Если ничего не сделал, acquisition lock истекает и record снова становится доступен для delivery attempt. (Apache Kafka)
Вот это, кстати, охуенно актуальная деталь для твоего собеса в 2026 году. Большая часть кандидатов будет рассказывать Kafka из мира 2.x и 3.x.
21. KRaft и ZooKeeper
Если тебя спросят современную Kafka:
В Kafka 4.x ZooKeeper больше нет. Kafka 4.0 и выше работает только в KRaft mode.
В KRaft control plane интегрирован в Kafka. Brokers обслуживают data requests, controllers управляют metadata, а controller quorum использует Raft protocol. ZooKeeper mode был удалён начиная с Kafka 4.0. (Apache Kafka)
Модель:
KRaft controllers
|
cluster metadata
|
Kafka brokers
|
partition data
На собеседовании не надо читать лекцию про Raft.
Достаточно:
Раньше cluster metadata и coordination были связаны с ZooKeeper. В современной Kafka 4.x ZooKeeper удалён, metadata управляется KRaft controller quorum.
22. Что такое Kafka Connect
Нужно хотя бы знать термин.
External System
↓
Source Connector
↓
Kafka
↓
Sink Connector
↓
External System
Например:
PostgreSQL -> Kafka
Kafka -> Elasticsearch
Kafka -> S3
То есть:
Kafka Connect является framework для построения и запуска connectors между Kafka и внешними системами.
Не надо руками писать:
while true:
select postgres
producer.send(...)
когда существует нормальный connector ecosystem.
23. Kafka Streams
Kafka Streams:
Kafka topic
↓
filter
map
group
aggregate
join
window
↓
Kafka topic
Например:
Payment events
↓
groupBy(accountId)
↓
window 5 min
↓
sum
↓
Fraud candidates
Нужно знать понятия:
KStream
KTable
state store
window
join
aggregation
Но для обычной backend architecture роли я бы не тратил сейчас часы на Streams DSL.
Тебе важнее понимать:
partitioning
consumer groups
offsets
delivery semantics
transactions
replication
lag
rebalance
24. Transactional Outbox
Это тебе обязательно.
Есть:
POST /orders
Сервис делает:
INSERT order
KafkaProducer.send(OrderCreated)
Проблема:
DB commit success
Kafka send failed
Получаем:
Order exists
OrderCreated doesn't
Или наоборот:
Kafka success
DB rollback
Получаем событие о несуществующем order.
Transactional Outbox:
BEGIN
INSERT orders
INSERT outbox
COMMIT
Одна DB transaction.
После этого:
Outbox Publisher
↓
Kafka
Например:
orders
----------------
id=100
outbox
----------------
event_id=abc
type=OrderCreated
aggregate_id=100
payload=...
Publisher:
SELECT outbox
↓
Kafka send
↓
mark published
Да, возможен duplicate publish.
Поэтому:
event_id
idempotent consumer
Именно здесь видно настоящее понимание distributed systems:
Мы не изображаем единую магическую transaction поверх PostgreSQL и Kafka. Мы меняем архитектуру так, чтобы атомарность нужна была только внутри одной transactional boundary.
Вот эту фразу я бы тебе вообще запомнил.
25. Что тебе надо уметь ответить без размышления
Вот прям контрольный список:
Что такое Kafka?
Distributed event log / event streaming platform. Topics разбиваются на ordered partitions, consumers читают log по offsets.
Kafka очередь?
Классическая Kafka не является обычной destructive queue. Records хранятся независимо от чтения, consumers управляют position через offsets и могут replay данные. Но начиная с Kafka 4.2 появились production-ready share groups для queue-like work processing. (Apache Kafka)
Что такое partition?
Ordered log и основная единица parallelism и distribution.
Где гарантируется ordering?
Только внутри partition.
Зачем key?
В том числе для partition routing. События одной entity можно направлять в одну partition и сохранять их порядок.
Что такое consumer group?
Группа consumers, между которыми Kafka распределяет partitions для параллельной обработки одного logical stream.
Могут два consumer одной классической group читать одну partition одновременно?
В каждый момент одна partition назначена одному consumer внутри subscribing consumer group. (Apache Kafka)
Могут разные groups прочитать одно сообщение?
Да. Groups читают topic независимо.
Что такое offset?
Position record внутри partition; committed consumer position определяет, откуда группа восстановит чтение.
At least once?
Process before commit. Возможны duplicates.
At most once?
Commit before processing. Возможна loss.
Exactly once?
Нужно определить boundary. Kafka transactions обеспечивают EOS для Kafka-to-Kafka processing при правильной работе с transactional producer, offsets и read_committed. Внешняя DB автоматически в эту guarantee не входит. (Apache Kafka)
Что такое ISR?
In-sync replicas, участвующие в replication durability model.
Что такое acks=all?
Leader ждёт acknowledgement от текущего ISR set. Это strongest producer acknowledgement guarantee Kafka. (Apache Kafka)
Зачем min.insync.replicas?
Не позволить успешно подтверждать запись при слишком деградировавшей replication state.
Что такое lag?
Отставание consumer position от конца log.
Что такое rebalance?
Redistribute partition assignments между consumers группы.
Что такое compaction?
Сохранение latest value per key в compacted log.
Как решить DB + Kafka atomicity?
Transactional outbox плюс idempotency/deduplication.
Игорь, вот это уже тот Kafka-каркас, который тебе нужен. У тебя проблема на прошлом собесе была не в том, что ты «не знаешь Kafka». У тебя знания были фрагментами: очередь, consumer, partitions, offset. А интервьюер мог ткнуть между фрагментами и получалась щель.
Здесь центральная цепочка одна:
LOG
↓
PARTITIONS
↓
KEY DETERMINES DISTRIBUTION
↓
CONSUMER PULLS FROM POSITION
↓
GROUP OWNS PARTITIONS
↓
OFFSET DEFINES RECOVERY
↓
COMMIT ORDER DEFINES DELIVERY SEMANTICS
↓
REPLICATION DEFINES DURABILITY
↓
TRANSACTION BOUNDARY DEFINES WHAT "EXACTLY ONCE" ACTUALLY MEANS
Вот её надо держать в голове. Остальная Kafka на неё навешивается.