Две политики очистки, OFFSET_OUT_OF_RANGE и почему после компакции в оффсетах появляются дырки.
Kafka удаляет данные по расписанию, а не по факту прочтения. Это принципиальное отличие от очередей: сообщение не исчезает после обработки, но и не хранится вечно.
Политика cleanup.policy=delete удаляет старые сегменты по времени (retention.ms) или по размеру (retention.bytes). Опасность в том, что удаление не спрашивает, прочитал ли кто-нибудь эти записи. Если консьюмер отстал сильнее, чем retention, он запросит оффсет, которого больше нет, получит OFFSET_OUT_OF_RANGE и по умолчанию (auto.offset.reset=earliest) продолжит с начала живого лога — молча пропустив всё, что не успел прочитать. Это, пожалуй, самый неочевидный способ потерять данные в Kafka: кластер здоров, ошибок в логах нет, данные просто исчезли.
Отсюда правило: retention считают не от объёма, а от максимально допустимого простоя потребителя плюс запас на разбор инцидента. И ставят алерт на приближение лага к retention.
Политика cleanup.policy=compact решает другую задачу: хранить не историю, а последнее состояние каждого ключа. Компактор проходит по «чистой» части лога и оставляет по одной последней записи на ключ. Так делают топики состояния, справочники и changelog для Kafka Streams.
У компакции есть следствие, которое ломает наивный код: она удаляет записи, но не перенумеровывает оставшиеся. Оффсет записи неизменен на всю её жизнь, поэтому после компакции лог идёт с разрывами. Консьюмер обязан это переживать, просто читая следующую существующую запись. Никогда не считайте, что offset+1 существует.
Разобрать это в симуляторе Интерактивный сценарий: настраиваешь кластер, ломаешь его и смотришь, что происходитДа, cleanup.policy=compact,delete. Тогда старые записи удаляются по времени, а внутри оставшегося окна работает компакция.
Запись с ключом и пустым значением. Для компактируемого топика она означает удаление ключа и сама удаляется через delete.retention.ms.
По росту logStartOffset и по ошибкам OFFSET_OUT_OF_RANGE в логах консьюмера. Лучше не доводить: алертить на лаг заранее.