Компактация логамеханизм, который удаляет устаревшие значения по ключу и оставляет только последнее в Kafka — механизм, который удаляет устаревшие значения по ключу и оставляет только последнее. Каждый брокер выполняет её независимо над своим логом. Если реплика отстаёт или недоступна, а компактация на остальных успевает удалить запись метаданных раньше, чем отставшая реплика её получит, содержимое логов расходится навсегда. Разберём, как устроен выбор партиции для компактации, как строится карта ключей, какие записи решено сохранять, и почему независимая компактация на каждой реплике приводит к расхождению.
Как LogCleanerManager выбирает партицию для компактации
grabFilthiestCompactedLogМетод LogCleanerManager, который выбирает следующую партицию для компактации — самую «грязную» из тех, что уже можно чистить. Под блокировкой он фиксирует время последнего запуска и собирает карту последних контрольных точек очистки по всем партициям. Затем логи фильтруются: остаются только те, у которых включена компактация, которые не находятся в процессе очистки и не помечены как неочищаемые. Для каждого такого лога вычисляются очищаемые смещения, при необходимости обновляется контрольная точка, и создаётся объект LogToCleanОписание лога, подлежащего компактации: содержит границы очищаемого диапазона и вычисленную по ним «грязность» лога с firstDirtyOffsetСмещение, с которого начинается ещё не компактированная часть лога — нижняя граница очищаемого диапазона, firstUncleanableDirtyOffsetСмещение, до которого чистить нельзя — верхняя граница очищаемого диапазона и флагом needCompactionNowПризнак того, что лог нужно компактировать немедленно, не дожидаясь обычного порога (истинен, если maxCompactionDelay > 0). Пустые логи отбрасываются.
Далее вычисляется dirtiestLogCleanableRatio — максимум cleanableRatio среди всех грязных логов (или 0, если список пуст). Из грязных логов отбираются те, которые либо требуют немедленной компактации и имеют cleanableBytes() > 0, либо имеют cleanableRatio() > minCleanableRatio из конфигурации лога. Если подходящих логов не нашлось, компактация в этот раз не запускается. Иначе выбирается лог с максимальным cleanableRatio, помечается как LOG_CLEANING_IN_PROGRESS в карте inProgress и возвращается.
Формула грязности: cleanBytes — сумма размеров сегментов от -1 до firstDirtyOffset; cleanableBytes и firstUncleanableOffset берутся из calculateCleanableBytes; totalBytes = cleanBytes + cleanableBytes; cleanableRatio = (double) cleanableBytes / totalBytes.
Как строится карта ключ→offset
Cleaner строит карту ключ→offset, сначала очищая её и собирая сегменты в диапазоне [start, end), затем для каждого сегмента вызывает buildOffsetMapForSegment. Тот читает записи и для каждой записи с ключом и offset >= startOffset добавляет пару (ключ, offset) в карту, пока размер карты меньше maxDesiredMapSize = map.slots() * dupBufferLoadFactor.
SkimpyOffsetMapструктура, хранящая хеш ключа и offset в слотах фиксированного размера хранит хеш ключа и offset в слотах, где каждый слот занимает bytesPerEntry = hashSize + 8 байт, а количество слотов slots = memory / bytesPerEntry. При вставке put вычисляет хеш ключа через hashInto, затем ищет пустой слот с помощью positionOf. Для попытки attempt вычисляется probe = ByteUtils.readIntBE(hash, Math.min(attempt, hashSize - 4)) + Math.max(0, attempt - hashSize + 4), берётся Utils.abs(probe) % slots и умножается на bytesPerEntry. Если слот не пуст, сравнивается сохранённый хеш с hash1; при совпадении перезаписывается offset, иначе attempt увеличивается и поиск повторяется. При поиске get аналогично пробегает слоты, пока не найдёт совпадение хешей или пустой слот, возвращая -1 при неудаче. clear обнуляет все байты буфера и сбрасывает entries и lastOffset, делая все существующие ключи ненаходимыми.
Как Cleaner решает, сохранять ли запись
Внутри cleanInto создаётся фильтр MemoryRecords.RecordFilter, переопределяющий checkBatchRetention и shouldRetainRecord. checkBatchRetention сначала вычисляет canDiscardBatch через shouldDiscardBatch(batch, transactionMetadata). Для control-батчаслужебный батч с маркером транзакции COMMIT или ABORT discardBatchRecords истинно, когда батч можно отбросить и его deleteHorizonMs уже наступил; иначе discardBatchRecords совпадает с canDiscardBatch.
Далее вычисляется isBatchLastRecordOfProducer: берётся lastRecordsOfActiveProducers.get(batch.producerId()), и если у lastRecord есть lastDataOffset, сравнивается batch.lastOffset() == lastRecord.lastDataOffset().getAsLong(), иначе возвращается batch.isControlBatch() && batch.producerEpoch() == lastRecord.producerEpoch(); если записи нет — false. На основе этого выбирается BatchRetentionрежим судьбы батча при очистке: RETAIN_EMPTY — батч сохраняется, но без записей, чтобы не потерять информацию о последнем offset; DELETE — батч удаляется целиком; DELETE_EMPTY — батч удаляется, если после очистки в нём не осталось записей. Режим RETAIN_EMPTY выбирается, если batch.hasProducerId() && isBatchLastRecordOfProducer, а также если batch.nextOffset() == highWatermark; DELETE — если discardBatchRecords; иначе DELETE_EMPTY.
В shouldRetainRecord фильтра: если discardBatchRecords — записи батча удаляются, а сам батч сохраняется только ради sequence продюсера; иначе если batch.isControlBatch() — батч сохраняется; иначе вызывается Cleaner.this.shouldRetainRecord(...). Тот сначала проверяет pastLatestOffset = record.offset() > map.latestOffset(), и если true — запись сохраняется. Если у записи нет ключа — фиксируется stats.invalidMessage() и запись удаляется. Для записи с ключом вычисляется foundOffset = map.get(key), latestOffsetForKey = record.offset() >= foundOffset, legacyRecord = batch.magic() < RecordBatch.MAGIC_VALUE_V2. Для не-legacy записи shouldRetainDeletes = batch.deleteHorizonMs().isEmpty() || currentTime < batch.deleteHorizonMs().getAsLong(); для legacy — retainDeletesForLegacyRecords. Итог: isRetainedValue = record.hasValue() || shouldRetainDeletes, запись сохраняется, если latestOffsetForKey && isRetainedValue.
Как группируются сегменты
Сегменты группируются функцией groupSegmentsBySize: первый сегмент становится началом группы, к ней добавляются следующие, пока суммарный размер данных лога не превысит maxSize, а суммарные размеры offset-индекса и time-индекса не превысят maxIndexSize. Дополнительно проверяется, что оценка lastOffsetForFirstSegment(segments, firstUncleanableOffset) минус baseOffset последнего сегмента группы не превышает Integer.MAX_VALUE, кроме случая, когда размер первого сегмента равен 0. lastOffsetForFirstSegment оценивает последний offset первого сегмента: если есть следующий сегмент, берётся его baseOffset минус 1, иначе firstUncleanableOffset минус 1. Лимиты задаются при вызове из doClean: log.config().segmentSize() как maxSize и log.config().maxIndexSize как maxIndexSize. В cleanSegments при переполнении текущий очищенный сегмент завершается и создаётся новый.
Как обрабатываются транзакционные метаданные
Cleaner собирает метаданные транзакций в CleanedTransactionMetadata: при построении карты смещений он вызывает log.collectAbortedTransactions(start, end) и добавляет результат через addAbortedTransactions, который добавляет список в очередь abortedTransactions. При чтении каждого батча: если батч контрольный, вызывается onControlBatchRead; иначе onBatchRead. onBatchRead сначала продвигает очередь abortedTransactions до lastOffset батча через consumeAbortedTxnsUpTo, затем для транзакционного батча проверяет producerId в ongoingAbortedTxns — это набор продюсеров, чьи транзакции уже признаны прерванными по собранным метаданным. Если продюсер там найден, значит, все его батчи относятся к прерванной транзакции и подлежат удалению: Cleaner запоминает lastObservedBatchOffset и помечает батч как удаляемый. Если продюсера там нет, его producerId добавляется в ongoingCommittedTxns — набор продюсеров, чьи транзакции считаются закоммиченными, — и батч сохраняется.
onControlBatchRead для ABORT-маркера удаляет producerId из ongoingAbortedTxns и сохраняет маркер, если у abortedTxnMetadata задан lastObservedBatchOffset, иначе маркер удаляется; для COMMIT-маркера маркер удаляется, если транзакция не встречалась среди прочитанных батчей. В фильтре очистки shouldRetainRecord для контрольного батча всегда сохраняет его, а discardBatchRecords заставляет удалять записи батча, сохраняя сам батч для информации о последовательности продюсера.
Как рассчитывается диапазон очищаемых смещений
Диапазон вычисляется в cleanableOffsets: сначала берётся lastCleanOffset (последнее сохранённое в чекпоинте смещение) или logStartOffset, если его нет; если это значение меньше logStartOffset или больше log.logEndOffset(), то firstDirtyOffset сбрасывается на logStartOffset и forceUpdateCheckpoint=true, иначе firstDirtyOffset равен checkpointDirtyOffset. Верхняя граница firstUncleanableDirtyOffset — минимум из трёх значений: log.lastStableOffset(), log.activeSegment().baseOffset() и (если minCompactionLagMs > 0) результата findFirstUncleanableSegment. Тот перебирает неактивные сегменты начиная с firstDirtyOffset и возвращает baseOffset первого сегмента, у которого segment.largestTimestamp() > now - minCompactionLagMs.
calculateCleanableBytes берёт неактивные сегменты от uncleanableOffset, определяет firstUncleanableOffset как baseOffset первого из них (или активного сегмента, если список пуст) и суммирует размеры сегментов в диапазоне от min(firstDirtyOffset, firstUncleanableOffset) до firstUncleanableOffset.
Как отслеживается состояние очистки
Состояние очистки хранится в карте inProgress, где ключ — TopicPartition, а значение — LogCleaningStateсостояние очистки партиции: очистка идёт, очистка прервана или очистка приостановлена. Возможны три состояния: очистка идёт (LOG_CLEANING_IN_PROGRESS) — партиция сейчас обрабатывается cleaner'ом; очистка прервана (LOG_CLEANING_ABORTED) — работа по партиции остановлена; очистка приостановлена (LogCleaningPaused) — партиция временно исключена из очистки и хранит счётчик пауз, показывающий, сколько раз её приостанавливали. Состояние партиции читается из этой карты, а установка нового состояния записывает его туда же. Проверка конкретного состояния сравнивает текущее состояние с ожидаемым, а проверка паузы срабатывает только тогда, когда партиция находится в состоянии приостановленной очистки.
Пауза для некомпактируемых партиций: выбираются логи без компактации, пропускаются уже находящиеся в inProgress, и каждому ставится состояние приостановленной очистки с счётчиком пауз, равным единице. Снятие паузы уменьшает счётчик: при единственной паузе запись о партиции убирается, при нескольких — счётчик уменьшается на единицу, а если состояние не является приостановленной очисткой, это считается ошибкой. Прерывание и пауза очистки: если партиция ещё не в работе, ей ставится приостановленная очистка с единичным счётчиком; если очистка идёт, состояние меняется на прерванную; если партиция уже приостановлена, счётчик пауз увеличивается на единицу; после этого система дожидается, пока партиция не окажется в состоянии приостановленной очистки. Для выбранного самого «грязного» лога выставляется состояние идущей очистки, и тем же состоянием помечаются все партиции, подлежащие удалению.
В чём именно гонка
Гонка состоит в том, что каждый брокер компактит свой лог независимо, и пока одна реплика отстаёт или офлайн, компактация на остальных может удалить критичную запись метаданных раньше, чем отставшая реплика успеет её получить. Тогда отставшая реплика о записи так и не узнает, и содержимое логов навсегда расходится. Удаляться могут: маркеры удаленияtombstone — запись, помечающая ключ как удалённый, управляющие батчиcontrol-батч с маркером COMMIT или ABORT, а также оставшийся после маркера пустой батч.
Маркер удаления может быть удалён через delete.retention.ms после его записи. Управляющий батч сначала заменяется пустым батчем с ID продюсера и флагом COMMIT/ABORT через delete.retention.ms, а затем этот пустой батч может быть удалён через producer.id.expiration.ms. Если реплика отстаёт дольше этих таймеров, она пропускает и маркер, и пустой батч, и восстановить информацию уже невозможно.
Четыре проявления проблемы
Различаются они тем, какая именно запись метаданных теряется.
Проблема 1 — потеря маркера удаления для ключа K. Маркер удаления для ключа K записывается, пока Broker 2 недоступен. Brokers 1 и 3 в ходе компактации удаляют и исходное значение, и сам маркер удаления. Для Brokers 1 и 3 ключ K удалён, а Broker 2 по-прежнему отдаёт K=V. Что именно увидит консьюмер, зависит от того, какой брокер в этот момент будет лидером.
Проблема 2 — потеря маркера ABORT для TX1. Если Broker 2 не получил маркер ABORT для TX1, а на других брокерах тот уже был удалён компактацией, данные poison останутся в его логе. На Brokers 1 и 3 консьюмер с read_committed видит good=data, но не видит никаких записей poison; на Broker 2 тот же консьюмер видит poison SHOULD_NOT_SEE_THIS.
Проблема 3 — потеря маркера COMMIT для TX1. Если Broker 2 не получил маркер COMMIT для TX1, а на других брокерах он уже был удалён компактацией, то при получении ABORT из TX2 Broker 2 применит его и к данным TX1. Закоммиченные данные будут ошибочно помечены как отменённые и исчезнут для консьюмеров с read_committed.
Проблема 4 — потеря маркера COMMIT вместе с пустым батчем. Если Broker 2 не получил маркер COMMIT, а затем тот был удалён компактацией вместе с оставшимся после него пустым батчем, на Broker 2 по-прежнему остаются транзакционные данные, но сам брокер не знает, что транзакция завершена. Он считает данные незакоммиченными и фиксирует последний стабильный offset (LSO) на этой позиции. Когда Broker 2 становится лидером, консьюмеры с read_committed перестают видеть всё, что записано после этого места: для них партиция фактически замирает, хотя новые данные продолжают в неё поступать. Такое состояние сохраняется, пока не истечёт producer.id.expiration.ms, отсчитываемый от последней записи этого продюсера в логе, — по умолчанию 24 часа.
Воспроизводим баг пошагово
Сначала создаётся топик с компактацией. Единственный переопределённый параметр — delete.retention.ms: по умолчанию он равен 24 часам, но setup.sh уменьшает его в соответствии с длительностью теста.
kafka_topics --create --topic foo --partitions 1 --replication-factor 3 \
--config cleanup.policy=compact \
--config delete.retention.ms=$DELETE_RETENTION_MSЗатем в фоне запускается продюсер aborted-to-committed.py, после чего ожидается сигнал tx1_produced. Далее подаётся сигнал do_abort через touch /tmp/signals/do_abort и ожидается tx1_aborted с таймаутом 180. После этого подаётся do_tx2 и ожидается tx2_committed.
Сначала всех лидеров __transaction_state переносят с Broker 2, чтобы избежать зависания commit при смене координатора: иначе, если именно Broker 2 окажется координатором нужного transactional.id, предстоящий commit может зависнуть при переключении координатора. Это делается командой move_tx_coord_off 2. Затем Broker 2 останавливают через docker compose kill kafka2 и ждут, пока лидерство партиции foo перейдёт на другой брокер, циклом while [ "$(get_leader)" = "2" ]; do sleep 1; done. Далее Broker 2 возвращают через docker compose start kafka2, ждут его вхождения в ISR циклом while ! kafka_topics --describe --topic foo | grep -qP 'Isr:\s*[123],[123],[123]'; do sleep 1; done и принудительно делают лидером командой force_leader 2.
Ротация сегмента форсируется записью около 1 ГБ данных-заполнителя. После этого скрипт ждёт delete.retention.ms, чтобы маркер ABORT стал доступен для удаления, и снова пишет 1 ГБ — в коде это два вызова pump_1gb с паузой sleep "$SLEEP_S" между ними. Затем в цикле while проверяется, остались ли в выводе kafka_consume с режимом read_uncommitted строки, начинающиеся с poison; пока такие строки есть, снова вызывается pump_1gb и пауза 15 секунд. После исчезновения poison выполняется пять дополнительных циклов pump_1gb; sleep 15, чтобы cleaner успел удалить отменённые данные, запись с маркером ABORT и оставшийся пустой батч.
Потребитель запускается дважды: сначала kafka_consume "kafka1:9092,kafka3:9092" read_committed | grep "^poison" — вывода нет, транзакция корректно отменена; затем kafka_consume "kafka2:9092" read_committed | grep "^poison" — вывод содержит poison SHOULD_NOT_SEE_THIS.
Как сравнивали и что получилось
Стенд — Docker Compose с тремя брокерами Kafka (kafka1, kafka2, kafka3) и контейнером-продюсером txproducer. Топик foo создаётся с одной партицией и фактором репликации 3, cleanup.policy=compact, а delete.retention.ms уменьшается setup.sh пропорционально длительности теста. Окружение пересоздаётся командами docker compose down --volumes 2>/dev/null, docker compose up -d и sleep 10.
Ручной сценарий: подключается setup.sh, загружающий вспомогательные функции и уменьшенные значения параметров; в контейнере txproducer запускается python3 /scripts/aborted-to-committed.py в фоне; ожидается сигнал tx1_produced через wait_for_signal; создаётся файл-сигнал /tmp/signals/do_abort командой touch; ожидается tx1_aborted с таймаутом 180 секунд; создаётся /tmp/signals/do_tx2; ожидается tx2_committed. Автоматический запуск: git clone репозитория redpanda-data-blog/kafka-log-compaction-bug-fix, переход в kafka-compaction-divergence/aborted-to-committed и запуск ./diverge.sh 10m — агрессивные настройки компактации, расчёт примерно на 10 минут.
Результат по сценарию с потерей маркера ABORT: на Brokers 1 и 3 консьюмер с read_committed видит good=data, но не видит записей poison; на Broker 2 тот же консьюмер видит poison SHOULD_NOT_SEE_THIS. Потребитель с bootstrap-серверами kafka1:9092,kafka3:9092 в режиме read_committed не выдаёт строк poison — транзакция корректно отменена; потребитель с bootstrap-сервером kafka2:9092 выдаёт poison SHOULD_NOT_SEE_THIS.
Оговорка: измерения ограничены одним сценарием — отменённая транзакция, ставшая закоммиченной, на партиции с одной партицией и фактором репликации 3.
Решение Redpanda Streaming: координированная компактация
В Redpanda Streaming каждая реплика отслеживает MCCOмаксимальный offset, до которого компактация её собственного лога полностью завершена — ниже этой точки нет повторяющихся ключей, а для каждого ключа остаётся не более одного значения или маркера удаления. Лидер периодически спрашивает у каждой follower-реплики её MCCO, затем вычисляет MTROминимальный MCCO среди всех реплик, ниже которого маркеры удаления можно безопасно удалять как минимальный MCCO среди всех реплик, включая временно недоступные (для них используется последнее известное значение MCCO), и рассылает MTRO всем репликам. Маркер удаления ниже MTRO можно безопасно удалить, потому что каждая реплика уже выполнила компактацию старых значений, которые этот маркер заменяет, и ни одна из них не останется с удалённым ключом как с актуальным.
Аналогично для транзакционных маркеров: каждая реплика отслеживает MXFOoffset, до которого состояние транзакций на реплике полностью определено — все транзакции продюсеров ниже этой точки либо закоммичены, либо отменены, и ни одна не остаётся активной, — а лидер вычисляет MXROминимальный MXFO среди всех реплик, ниже которого маркеры транзакций можно безопасно удалять как минимальный MXFO среди всех реплик; маркер COMMIT/ABORT можно безопасно удалить, когда он находится ниже MXRO. Если реплика долго недоступна, MTRO/MXRO не продвигаются вперёд, поэтому удаление маркеров приостанавливается во всём кластере, что предотвращает расхождение реплик. Когда реплика возвращается и выполняет компактацию, её MCCO/MXFO продвигается вперёд, лидер пересчитывает MTRO/MXRO, и очистка возобновляется.
Причина, по которой координация необходима, сформулирована так: ничто не гарантирует, что реплика не будет недоступна или не отстанет дольше, чем на delete.retention.ms, поэтому, чтобы система работала корректно даже при длительной недоступности или сильном отставании брокеров, Redpanda дополняет компактацию небольшим протоколом координации. Компромисс по задержке компактации назван прямо: при недоступной реплике удаление маркеров приостанавливается во всём кластере — это осознанное архитектурное решение, корректность данных гарантируется, а компактация выполняется по возможности. По нагрузке сказано лишь, что лидер периодически опрашивает follower-реплики об их MCCO и рассылает MTRO, а при смене лидера повторно рассылает MTRO, даже если значение не изменилось, — некоторые реплики могли пропустить последнее обновление во время смены лидера. Итоговая заявленная выгода — брокеры совместно определяют, какие записи уже можно удалить, и освобождают максимум дискового пространства без ущерба для безопасности данных.
Что из этого следует на практике
Компактация в Kafka — локальная операция каждого брокера над собственным логом, и её безопасность относительно репликации ничем не координируется. Отсюда прямое следствие: любая реплика, отставшая дольше delete.retention.ms (для маркеров удаления) или дольше producer.id.expiration.ms (для пустых батчей после управляющих маркеров), может пропустить критичную запись метаданных, и восстановить её будет невозможно.
Что именно теряется, определяет характер поломки. Потеря маркера удаления возвращает удалённые данные, и наблюдаемый результат зависит от того, какой брокер станет лидером. Потеря маркера ABORT делает отменённые данные видимыми для read_committed. Потеря маркера COMMIT либо превращает закоммиченные данные в отменённые, либо фиксирует LSO на позиции и замораживает партицию для read_committed до истечения producer.id.expiration.ms (по умолчанию 24 часа).
Пороги, от которых зависит, как быстро компактация добирается до записей, задаются конфигурацией: minCleanableRatio определяет, при какой доле очищаемых байт лог считается достойным компактации, maxCompactionDelay — когда компактация запускается немедленно, minCompactionLagMs — насколько свежие сегменты остаются неочищаемыми, delete.retention.ms — сколько живёт маркер удаления. Уменьшение delete.retention.ms (как в сценарии воспроизведения) сокращает окно, в течение которого отставшая реплика может получить маркер, и делает расхождение воспроизводимым за минуты вместо суток.
Единственный описанный способ закрыть гонку — координация: реплики сообщают лидеру, до какого offset их компактация завершена, лидер берёт минимум по всем репликам (включая недоступные, по последнему известному значению) и разрешает удалять маркеры только ниже этого минимума. Плата за это — приостановка удаления маркеров во всём кластере, пока хотя бы одна реплика отстаёт.
Где смотреть в коде
- LogCleanerManager.java: grabFilthiestCompactedLog
- Cleaner.java: cleanInto
- LogToClean.java: LogToClean
- Cleaner.java: shouldRetainRecord
- LogCleanerManager.java: cleanableOffsets
- Cleaner.java: groupSegmentsBySize
- LogCleanerManager.java: deletableLogs
- LogCleanerManager.java: resumeCleaning