1. Главная модель Kafka: это не очередь, а распределённый log

Самая важная фраза:

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 на неё навешивается.

Прокрутить вверх