The Jerusalem PostJapanese PM Takaichi protests to US after Marine accused of murder in OkinawaESPN DeportesFalcons y Saints ponen fin a la Semana 4 con emotivo reencuentro en "Monday Night Football"ESPNWhat are the top 10 upsets in Super Bowl history?PunchMotor park killing: Why we handed over operatives to police – Osun AmotekunZDF heuteAktuelle Pressemitteilungen des ZDFVilaWebEls EUA reclamen a Rússia transparència per un suposat cas de pesta pulmonar a SibèriaBBC NewsMan detained for stabbing dog walker to deathWirtualna PolskaTragedia na A2. Motocyklista wbił się w tira i zginąłCNN بالعربيةاستمرار مارك تومسون رئيسًا لشبكة CNN مع تغير الشركة الأمSouth China Morning PostBid to turn 14 subsidised sale flats into subdivided homes raises abuse concernsالشرقالسيسي: لن نسمح بأي إجراء يمس أمن مصر المائي.. و"حصة النيل حق تاريخي وقانوني"VanguardNigeria may crash chaotically under Atiku – Babachir Lawal says he’d rather back Tinubu
The Daily Newsstand · Free, Always
Monday, October 5, 2026

[Перевод] Компактация лога в Apache Kafka повреждала данные. Вот как мы это исправили

Translate

В топиках с компактацией Apache Kafka сохраняет только последнее значение для каждого ключа. Для удаления ключа используются маркеры удаления (tombstone) – записи со значением null. После того как компактация удалит записи со значениями для такого ключа, Kafka ждёт как минимум delete.retention.ms, а затем удаляет и маркер удаления. Благодаря этому топик не разрастается из-за маркеров удаления для давно исчезнувших ключей.

Но здесь есть проблема. Точнее, сразу четыре. В этой статье мы разберём найденный баг и покажем, как координированная компактация решает его в Redpanda Streaming.

Как работает компактация лога в Kafka

Чтобы разобраться в проблеме, сначала вкратце посмотрим, как сейчас устроена компактация лога в Kafka.

Компактация затрагивает и управляющие батчи транзакций. При транзакционной записи продюсер сначала записывает данные – возможно, сразу в несколько партиций, – а затем добавляет в каждую партицию управляющий батч COMMIT или ABORT. Консьюмер с isolation.level=read_committed использует эти маркеры, чтобы определить, нужно ли отдавать записи транзакции или скрыть их.

Управляющие батчи хранятся в логе наряду с обычными записями, и в топиках с компактацией к ним применяются те же правила удаления по сроку хранения. Когда связанные с ними данные уже удалены компактацией и прошло достаточно времени, сам маркер тоже можно удалить. Это позволяет эффективно очищать лог от старых данных и метаданных. Компактация удаляет не только записи с данными и маркеры удаления для удалённых ключей, но и связанные с ними управляющие батчи транзакций.

Маркеры удаления и управляющие батчи COMMIT/ABORT – единственные сигналы о том, что связанные с ними записи соответственно удалены, закоммичены или отменены. Как только маркер удаления или управляющий батч исчезает в результате компактации, эта информация теряется.

Последствия могут быть катастрофическими: компактация может удалить маркер удаления или управляющий батч на одной реплике, пока другой он ещё нужен. Каждый брокер выполняет компактацию своего лога независимо. Реплика, на которую маркер удаления или маркер транзакции не попал, продолжает хранить связанные с ним записи. Когда она снова подключается, у лидера уже нет этого маркера, поэтому реплицировать его невозможно. В результате содержимое логов на репликах навсегда расходится, а то, какую версию данных увидит консьюмер, зависит от того, какой брокер окажется лидером в момент чтения.

Баг стабильно воспроизводится в Kafka версий 3.9–4.2. Мы обнаружили четыре его варианта: от «удалённые данные снова появляются» до «отменённые данные выдаются как закоммиченные». Дальше разберём все четыре случая, пошагово воспроизведём один из них и покажем, как мы устранили эту проблему.

Корень проблемы: гонка между компактацией и репликацией

Когда брокер начинает отставать или уходит в офлайн, он выпадает из набора синхронизированных реплик (ISR, in-sync replica set). Остальные брокеры тем временем продолжают принимать записи и выполнять компактацию как обычно. Если критически важная запись – маркер удаления, маркер COMMIT или ABORT – появляется, пока одна из реплик недоступна, а компактация удаляет её раньше, чем реплика успевает синхронизироваться, эта реплика о записи так и не узнает. С её точки зрения, такой записи никогда не существовало.

Защитные механизмы Kafka основаны на времени. Маркер удаления можно удалить через delete.retention.ms после его записи – по умолчанию это 24 часа. Для управляющих батчей транзакций очистка проходит в два этапа:

  1. Через delete.retention.ms сам батч с маркером заменяется пустым батчем, в заголовке которого по-прежнему хранятся ID продюсера и флаг COMMIT/ABORT.

  2. Через producer.id.expiration.ms – по умолчанию тоже 24 часа с момента последней активности продюсера – этот пустой батч также можно удалить.

Если брокер остаётся в офлайне или отстаёт дольше этих таймеров – например, из-за аппаратного сбоя, длительного обслуживания или медленного восстановления, – он пропустит и сам маркер, и оставшийся после него пустой батч. Восстановить эту информацию уже будет невозможно.

Мы обнаружили четыре проявления этой проблемы – в зависимости от того, какая именно запись с метаданными теряется. Во всех сценариях ниже используется кластер из трёх брокеров, в котором Broker 2 надолго уходит в офлайн.

Проблема 1: расхождение из-за маркера удаления – удалённые данные возвращаются

Маркер удаления для ключа K записывается, пока Broker 2 недоступен. Brokers 1 и 3 в ходе компактации удаляют и исходное значение, и сам маркер удаления. Когда Broker 2 возвращается в кластер, реплицировать на него уже нечего, поэтому исходная запись остаётся в его логе. Для Brokers 1 и 3 ключ K удалён, а Broker 2 по-прежнему отдаёт K=V. Что именно увидит консьюмер, зависит от того, какой брокер в этот момент будет лидером.

Проблема 2: отменённая транзакция становится закоммиченной

Продюсер выполняет две транзакции с одним и тем же transactional.id:

  1. TX1: записывает poison=SHOULD_NOT_SEE_THIS, затем выполняет ABORT.

  2. TX2: записывает good=data, затем выполняет COMMIT.

Если Broker 2 не получил маркер ABORT для TX1, а на других брокерах тот уже был удалён компактацией, данные poison останутся в его логе. При восстановлении состояния транзакций следующим управляющим батчем от того же продюсера окажется COMMIT из TX2, и Broker 2 применит его в том числе к данным TX1. В результате запись poison будет отдаваться консьюмерам с read_committed как обычные закоммиченные данные. То есть данные, которые приложение явно откатило, попадут в последующие системы обработки как настоящие, а Kafka никак не просигнализирует о проблеме.

Проблема 3: закоммиченная транзакция становится отменённой

Сценарий похож на предыдущий. В TX1 продюсер коммитит корректные данные, а затем в TX2 записывает мусорные данные и выполняет ABORT. Если Broker 2 не получил маркер COMMIT для TX1, а на других брокерах он уже был удалён компактацией, то при получении ABORT из TX2 Broker 2 применит его и к данным TX1. Закоммиченные данные будут ошибочно помечены как отменённые и исчезнут для консьюмеров с read_committed.

Проблема 4: партиция зависает, READ_COMMITTED перестаёт продвигаться

Продюсер начинает транзакцию, записывает K=V и выполняет COMMIT. Маркер COMMIT сообщает консьюмерам с read_committed, что теперь эти данные можно читать. Если Broker 2 не получил маркер COMMIT, а затем тот был удалён компактацией вместе с оставшимся после него пустым батчем, на Broker 2 по-прежнему остаются транзакционные данные, но сам брокер не знает, что транзакция завершена. Он считает данные незакоммиченными и фиксирует последний стабильный offset (LSO) на этой позиции.

Когда Broker 2 становится лидером, консьюмеры с read_committed перестают видеть всё, что записано после этого места: для них партиция фактически замирает, хотя новые данные продолжают в неё поступать. Такое состояние сохраняется, пока не истечёт producer.id.expiration.ms, отсчитываемый от последней записи этого продюсера в логе, – по умолчанию 24 часа. А если тот же продюсер продолжает создавать новые транзакции и тем самым обновляет время последней активности PID, зависание может длиться практически бесконечно.

Воспроизводим баг пошагово

Скрипты для воспроизведения лежат в сопутствующем GitHub-репозитории. Из зависимостей нужен только Docker Compose. Каждый скрипт diverge.sh принимает аргумент командной строки, с помощью которого можно настроить Kafka так, чтобы баг воспроизводился быстрее. С настройками по умолчанию это займёт около двух дней.

Вот как запустить сценарий, в котором отменённая транзакция становится закоммиченной, в автоматическом режиме и с уменьшенным delete.retention.ms:

git clone https://github.com/redpanda-data-blog/kafka-log-compaction-bug-fix.git
cd kafka-log-compaction-bug-fix/kafka-compaction-divergence/aborted-to-committed
./diverge.sh 10m   # агрессивные настройки компактации, расчёт примерно на 10 минут

Можно пройти сценарий и вручную, шаг за шагом. Сначала подключим setup.sh: он загружает вспомогательные функции, которые понадобятся дальше, и задаёт значения параметров, уменьшенные пропорционально длительности теста.

git clone https://github.com/redpanda-data-blog/kafka-log-compaction-bug-fix.git
cd kafka-log-compaction-bug-fix/kafka-compaction-divergence/aborted-to-committed
source ./setup.sh 10m

Запустим свежий кластер из трёх брокеров:

docker compose down --volumes 2>/dev/null
docker compose up -d
sleep 10

Создадим топик с компактацией. Единственный переопределённый параметр – 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

Запустим Python-продюсер. Он использует confluent-kafka, уже включённый в образ txproducer через txproducer.Dockerfile. Продюсер начинает TX1, записывает key=poison и value=SHOULD_NOT_SEE_THIS, а затем ждёт:

docker compose exec -d txproducer python3 /scripts/aborted-to-committed.py
wait_for_signal tx1_produced

Убедимся, что транзакционные данные попали в ISR на всех трёх брокерах. Теперь в логе Broker 2 есть запись poison, но состояние её транзакции ещё не разрешено:

while ! kafka_topics --describe --topic foo | grep -qP 'Isr:\s*[123],[123],[123]'; do sleep 1; done

Сначала перенесём всех лидеров __transaction_state с Broker 2. Иначе, если именно Broker 2 окажется координатором нашего transactional.id, предстоящий commit может зависнуть при переключении координатора. Затем остановим Broker 2 и дождёмся, пока лидерство партиции foo перейдёт на другой брокер. Всё, что произойдёт дальше, Broker 2 уже пропустит:

move_tx_coord_off 2
docker compose kill kafka2
while [ "$(get_leader)" = "2" ]; do sleep 1; done

Теперь дадим продюсеру сигнал выполнить ABORT для TX1. Маркер ABORT получат только Brokers 1 и 3:

docker compose exec -T txproducer touch /tmp/signals/do_abort
wait_for_signal tx1_aborted 180

Запишем около 1 ГБ данных-заполнителя, чтобы форсировать ротацию сегмента, подождём delete.retention.ms, пока маркер ABORT не станет доступен для удаления, затем запишем ещё 1 ГБ. После этого выполним ещё несколько циклов, чтобы cleaner успел удалить отменённые данные, запись с маркером ABORT и оставшийся после неё пустой батч:

pump_1gb; sleep "$SLEEP_S"; pump_1gb
while [ "$(kafka_consume "$BOOTSTRAP" read_uncommitted | grep -c '^poison')" -gt 0 ]; do
    pump_1gb; sleep 15
done
for i in 1 2 3 4 5; do pump_1gb; sleep 15; done

Теперь дадим продюсеру сигнал выполнить TX2 (key=good value=data) и сделать COMMIT:

docker compose exec -T txproducer touch /tmp/signals/do_tx2
wait_for_signal tx2_committed

На этом этапе свежий маркер COMMIT для TX2 уже находится в логах Brokers 1 и 3. На этих брокерах консьюмер с read_committed видит good=data, но не видит никаких записей poison:

kafka_consume "kafka1:9092,kafka3:9092" read_committed | grep "^poison"
# вывода нет; транзакция корректно отменена

Теперь вернём Broker 2, дождёмся, пока он снова войдёт в ISR, и принудительно сделаем его лидером:

docker compose start kafka2
while ! kafka_topics --describe --topic foo | grep -qP 'Isr:\s*[123],[123],[123]'; do sleep 1; done
force_leader 2

Прочитаем данные с Broker 2 с уровнем read_committed:

kafka_consume "kafka2:9092" read_committed | grep "^poison"
# poison	SHOULD_NOT_SEE_THIS

В логе Broker 2 по-прежнему лежат данные poison. Когда брокер восстанавливал состояние транзакций, первым управляющим батчем от этого продюсера, который он увидел, оказался COMMIT из TX2. Поэтому Broker 2 считает данные TX1 закоммиченными. Приложение отменило TX1, но консьюмеры, читающие через Broker 2, получают запись poison как валидные данные.

Тот же топик, та же партиция, те же данные лога на диске – но консьюмеры с read_committed получают разные результаты в зависимости от того, с какого брокера читают.

Решение Redpanda Streaming: координированная компактация

Redpanda Streaming соблюдает delete.retention.ms. Без этого ограничения медленный консьюмер мог бы увидеть старое значение ключа, но пропустить его маркер удаления или маркер транзакции, если тот успели удалить. При этом ничто не гарантирует, что реплика не будет недоступна или не отстанет дольше, чем на delete.retention.ms.

Чтобы система работала корректно даже при длительной недоступности или сильном отставании брокеров, Redpanda дополняет компактацию небольшим протоколом координации. Он не позволяет удалить маркер удаления или управляющий маркер, пока каждая реплика не выполнит компактацию связанных с ним записей данных.

Протокол удаления маркеров удаления

Координированная компактация использует два значения для каждой партиции:

  • MCCO (максимальный offset завершённой компактации), отдельно для каждой реплики. Каждая реплика отслеживает offset, до которого компактация её собственного лога полностью завершена: ниже этой точки нет повторяющихся ключей, а для каждого ключа остаётся не более одного значения или маркера удаления. MCCO движется только вперёд: если данные уже прошли такую компактацию, они остаются скомпактированными. Компактация выполняется только ниже high watermark, поэтому MCCO не может быть усечён при репликации.

  • MTRO (максимальный offset удаления tombstone-маркеров), один на набор реплик. Лидер вычисляет MTRO как минимальный MCCO среди всех реплик, включая временно недоступные: для них используется последнее известное значение MCCO. Поэтому MTRO не может продвинуться дальше прогресса компактации офлайн-реплики. Маркер удаления ниже MTRO можно безопасно удалить: каждая реплика уже выполнила компактацию старых значений, которые этот маркер заменяет, а значит, ни одна из них не останется с удалённым ключом как с актуальным.

Протокол работает в два этапа:

  1. Сбор. Лидер периодически спрашивает у каждой follower-реплики: «Какой у тебя MCCO?» Реплики сообщают текущий прогресс локальной компактации.

  2. Распространение. Лидер вычисляет MTRO = min(всех MCCO) и рассылает полученное значение всем репликам. После этого каждая реплика знает верхнюю безопасную границу для удаления маркеров удаления.

Схема работы протокола координированной компактации

Схема работы протокола координированной компактации

Теперь каждая реплика знает, что записи с offset ниже 80 можно безопасно удалять, а записи с offset 80 и выше нужно сохранять, пока самая медленная реплика не догонит остальные по компактации.

Обработка граничных случаев

Смена лидера. Когда выбирается новый лидер, в качестве отправной точки он использует ранее распространённое значение MTRO. Затем он начинает собирать MCCO со всех follower-реплик и вычисляет новый MTRO. Новый лидер повторно рассылает MTRO, даже если значение не изменилось: некоторые реплики могли пропустить последнее обновление во время смены лидера.

Изменение состава реплик. Когда в группу добавляется новая реплика, её MCCO инициализируется текущим MTRO группы. MCCO при этом может оказаться выше конца локального лога реплики. На первый взгляд это выглядит странно, но всё корректно: новая реплика получит лог от другой реплики, на которой компактация уже завершена вплоть до MTRO. Если реплику удалить из группы, MTRO может продвинуться вперёд, поскольку именно у выбывшей реплики MCCO мог быть самым низким.

MTRO никогда не движется назад. Если решение об очистке уже принято, оно окончательное. Попытки уменьшить MTRO, например из-за запоздавшего RPC от предыдущего лидера, игнорируются.

Протокол удаления маркеров транзакций

Для безопасного удаления маркеров транзакций используется похожая пара offset:

  • MXFO (максимальный offset с полностью разрешёнными транзакциями), отдельно для каждой реплики. Это offset, до которого состояние транзакций на реплике полностью определено: все транзакции продюсеров ниже этой точки либо закоммичены, либо отменены, и ни одна из них не остаётся активной. Как и MCCO, MXFO движется только вперёд.

  • MXRO (максимальный offset удаления маркеров транзакций), один на набор реплик. Это минимальный MXFO среди всех реплик, включая недоступные, для которых используется последнее известное значение. Маркер COMMIT/ABORT можно безопасно удалить, когда он находится ниже MXRO: к этому моменту каждая реплика уже обработала все маркеры и определила состояние всех транзакций ниже MXRO.

Безопасность данных важнее всего. Очистка выполняется по возможности

Даже если реплика долго остаётся недоступной, MTRO/MXRO не продвигаются вперёд, поэтому удаление маркеров удаления и транзакционных маркеров выше этих offset приостанавливается во всём кластере. Это осознанное архитектурное решение: корректность данных гарантируется, а компактация выполняется по возможности.

Когда реплика возвращается в кластер и выполняет компактацию, её MCCO/MXFO продвигается вперёд, лидер пересчитывает MTRO/MXRO, и очистка возобновляется.

Компактация без компромиссов

Алгоритм координированной компактации позволяет Redpanda Streaming принимать оптимальные решения об очистке даже в экстремальных условиях: например, при высокой нагрузке или длительной недоступности узла. Брокеры совместно определяют, какие записи уже можно удалить, и освобождают максимум дискового пространства без ущерба для безопасности данных.

Работа с Kafka быстро выходит за рамки настройки продюсеров и консьюмеров: приходится учитывать репликацию, транзакции, обмен между сервисами и поведение системы при сбоях. Разобраться в этих механизмах на практике можно на бесплатных вебинарах — от базовой работы с Kafka до её использования в микросервисной архитектуре.

  • 6 октября в 20:00. «Apache Kafka в микросервисной архитектуре — лучшие практики асинхронного обмена». Записаться

  • 21 октября в 20:00. «Apache Kafka: быстрый старт для разработчиков и инженеров». Записаться

Полный список бесплатных уроков октября по всем ИТ-направлениям можно посмотреть в дайджесте.

View the original on Хабр →

KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.