Сообщения в Kafka не удаляются сразу после того, как consumer их прочитал: брокер хранит их по правилам ретенции, а consumer лишь продвигает свой offset. Удаление сообщений происходит независимо от чтения — по времени, размеру или через лог-компакцию.
В Kafka чтение сообщений и их удаление жёстко разделены. Брокер просто хранит лог (последовательность сообщений) в партициях, а каждый consumer запоминает, до какого смещения (offset) он дочитал.
Когда consumer читает сообщения, они не «вытаскиваются» из очереди, как в классической очереди, а копируются из лога. После чтения consumer обычно коммитит offset (автоматически или вручную), тем самым помечая: «всё до этого offset я уже обработал». Это состояние хранится отдельно (чаще всего в служебном топике __consumer_offsets) и не влияет на сами сообщения в топиках.
Физическое удаление сообщений из брокера управляется настройками топика/кластера:
retention.ms — время жизни сообщений: старые записи удаляются, даже если их никто не прочитал, или наоборот — даже если их читали много раз.retention.bytes — максимальный размер лога: при переполнении самые старые сегменты удаляются, чтобы освободить место.log compaction) — для топиков с ключами Kafka может оставлять только «последнее состояние» по каждому ключу, удаляя устаревшие версии записей, но опять же это не связано с тем, кто и когда читал сообщения.Из этого следует важное следствие для собеседования: один и тот же топик могут читать разные consumer-группы независимо друг от друга, переигрывать сообщения, делать reprocess старых данных и поднимать новые сервисы, которые начнут читать «с начала» (если ретенция ещё позволяет), потому что сами сообщения не исчезают после чтения.
Отметьте свой прогресс