Kafka хранит сообщения на диске дольше, чем кажется — и именно это делает его логом, а не очередью: можно перечитать поток заново, отмотать consumer’а назад, восстановить состояние. Но диск не бесконечен, поэтому есть два способа чистки с разной семантикой: retention (удалить старое по времени/размеру) и log compaction (оставить последнее значение по каждому ключу). Путать их — значит либо терять нужное, либо копить лишнее.
Это четвёртая статья серии «Kafka: глубокое погружение». Опирается на модель партиции из первой статьи — здесь та же партиция впервые предстаёт не абстрактной последовательностью offset’ов, а конкретными файлами на диске, которые можно перечислить ls. Все числа ниже — из живого трёхброкерного KRaft-кластера: реальные размеры сегментов, реальный ход компакции по проходам, реальные коэффициенты сжатия. Одну вещь — tiered storage — воспроизвести на этом кластере честно не удалось, и почему именно, разобрано отдельно, без выдачи неудачи за успех. Живой стенд целиком открыт в репозитории digital-cookbook — kafka/go/storage (franz-go) и kafka/java/storage (kafka-clients), сценарий компакции — отдельным скриптом ops/inspect-segments.sh.
В статье
- Сегменты лога: как партиция лежит на диске
- Retention: три сегмента исчезают, earliest перескакивает на 600
- Log compaction: снимок последнего значения по ключу
- Два прохода log cleaner’а: почему tombstone не пропадает сразу
- Retention или compaction: как выбирать
- Компрессия: none, gzip, lz4, zstd на одних и тех же данных
- Dynamic и static конфиги: не всё можно поменять на лету
- Tiered storage: архитектура и честный итог
Сегменты лога: как партиция лежит на диске
Партиция физически — не один файл, а последовательность сегментов: каждый сегмент — пара файлов <base-offset>.log (сами записи) и <base-offset>.index/<base-offset>.timeindex (разреженные индексы offset→позиция в файле и время→offset, которые дают брокеру искать нужную запись без полного скана). Имя файла — 20-значный, дополненный нулями base offset первой записи в этом сегменте: 00000000000000000206.log начинается с offset 206.
В любой момент времени у партиции ровно один активный сегмент — тот, в который сейчас идёт запись. Все остальные сегменты закрыты («rolled») и неизменяемы. Roll закрывает текущий активный сегмент и открывает новый, когда срабатывает один из двух порогов: segment.bytes (сегмент достиг предельного размера) или segment.ms (сегмент существует дольше заданного времени, даже если ещё маленький). Это разграничение важно не только эксплуатационно: retention удаляет целые закрытые сегменты, а log compaction (дальше в статье) вообще не видит активный сегмент — компакции физически подвергаются только уже закрытые файлы.
segment.bytes — не просто конфиг «сколько байт до roll»: у брокера есть жёсткий нижний предел, обнаруженный прямо на живом кластере. Попытка создать топик с segment.bytes=2048 заканчивается отказом брокера: INVALID_CONFIG ("Value must be at least 1048576") — 1 МиБ это не рекомендация, а минимум, зашитый в валидацию LogConfig. Меньше выставить нельзя, и это стоит учитывать, если кажется удобным сделать сегменты совсем маленькими ради быстрой демонстрации retention.
Retention: три сегмента исчезают, earliest перескакивает на 600
retention.ms и retention.bytes — два независимых порога удаления: время жизни записи и суммарный размер лога партиции. Ключевое свойство обоих — retention удаляет целые закрытые сегменты, а не отдельные записи внутри них: как только весь сегмент старше retention.ms (или лог в сумме превышает retention.bytes), файл сегмента целиком помечается на удаление. Позиция consumer’ов при этом брокеру безразлична — retention ничего не знает о том, кто и что уже прочитал, и может удалить сегмент, даже если какой-то consumer ещё не добрался до его записей.
Живой прогон: топик demo-storage-retention, segment.bytes=1048576 (минимум брокера), retention.ms=8000, segment.ms намеренно выставлен большим (3600000 — час), чтобы roll управлялся исключительно размером, а не временем — иначе результат зависел бы от того, насколько быстро выполнился сценарий. 600 записей без ключа, ~5000 байт каждая, отправлены синхронно и без компрессии — важная оговорка, к которой статья ещё вернётся в разделе про сжатие: франц-go по умолчанию сжимает батчи Snappy, если явно не указать NoCompression(), и без этой явной настройки реальный размер на диске оказался бы примерно в 15 раз меньше расчётного, а roll срабатывал бы в разы позже.
// newSyncProducer — клиент-продюсер для синхронной (ProduceSync) отправки.
//
// ⚠️ Если compression не передан явно, ЯВНО ставим NoCompression() — franz-go
// по умолчанию использует Snappy (проверено живьём: без этой строки
// retention-produce с pad('x', 5000) на диске давал ~341 байт/запись вместо
// ~5000 — однобайтовый filler чудовищно сжимается снаппи, и segment.bytes
// roll срабатывал в разы позже, чем рассчитано по логическому размеру
// данных).
func newSyncProducer(seeds []string, compression ...kgo.CompressionCodec) *kgo.Client {
if len(compression) == 0 {
compression = []kgo.CompressionCodec{kgo.NoCompression()}
}
opts := []kgo.Opt{
kgo.SeedBrokers(seeds...),
kgo.RequiredAcks(kgo.AllISRAcks()),
kgo.ProducerBatchCompression(compression...),
}
// ...
}600 записей без компрессии при segment.bytes=1МБ дают около трёх закрытых сегментов и один пустой активный — есть что реально удалять:
ДО (600 записей отправлено, earliest=0 latest=600):
00000000000000000000.log 1044420 байт (offset 0..205)
00000000000000000206.log 1044420 байт (offset 206..411)
00000000000000000412.log 953160 байт (offset 412..599)
00000000000000000600.log 0 байт (активный, ещё пуст)
--- жду retention.ms(8с) + log.retention.check.interval.ms(2с) + запас (~14с) ---
ПОСЛЕ:
00000000000000000000.log.deleted
00000000000000000206.log.deleted
00000000000000000412.log.deleted
00000000000000000600.log 0 байт (активный, не тронут)
[assert] OK: earliest=600 > 0Все три закрытых сегмента, физически содержавшие данные, помечены .log.deleted и удалены; активный сегмент, в который после отправки 600-й записи уже ничего не писалось, остался как есть — просто пустым. earliest offset партиции сдвинулся ровно на длину удалённых сегментов, с 0 на 600: с точки зрения любого нового consumer’а лог теперь начинается там же, где заканчивается. На Java-клиенте тот же сценарий дал побайтово идентичные размеры сегментов (1044420 / 1044420 / 953160) — совпадение размеров на двух независимых клиентах подтверждает, что это предсказуемые данные без скрытой компрессии, а не совпадение по счастливой случайности.
demo-storage-retention-0"] --> S0["00000000000000000000.log
1044420 байт"] P --> S1["00000000000000000206.log
1044420 байт"] P --> S2["00000000000000000412.log
953160 байт"] P --> S3["00000000000000000600.log
активный, 0 байт"] S0 -->|"retention.ms истёк"| D["*.log.deleted
физически удалены"] S1 -->|"retention.ms истёк"| D S2 -->|"retention.ms истёк"| D D -.-> E["earliest offset: 0 -> 600"] style P fill:#f9f3e3,stroke:#8b7355 style S0 fill:#eee8d8,stroke:#8b7355 style S1 fill:#eee8d8,stroke:#8b7355 style S2 fill:#eee8d8,stroke:#8b7355 style S3 fill:#c9e4c5,stroke:#5b8a5e style D fill:#e8c9c0,stroke:#b4552f style E fill:#f9f3e3,stroke:#8b7355
flowchart LR
P["Партиция
demo-storage-retention-0"] --> S0["00000000000000000000.log
1044420 байт"]
P --> S1["00000000000000000206.log
1044420 байт"]
P --> S2["00000000000000000412.log
953160 байт"]
P --> S3["00000000000000000600.log
активный, 0 байт"]
S0 -->|"retention.ms истёк"| D["*.log.deleted
физически удалены"]
S1 -->|"retention.ms истёк"| D
S2 -->|"retention.ms истёк"| D
D -.-> E["earliest offset: 0 -> 600"]
style P fill:#f9f3e3,stroke:#8b7355
style S0 fill:#eee8d8,stroke:#8b7355
style S1 fill:#eee8d8,stroke:#8b7355
style S2 fill:#eee8d8,stroke:#8b7355
style S3 fill:#c9e4c5,stroke:#5b8a5e
style D fill:#e8c9c0,stroke:#b4552f
style E fill:#f9f3e3,stroke:#8b7355
Log compaction: снимок последнего значения по ключу
Retention чистит по времени или объёму, не глядя на содержимое. Log compaction — противоположный принцип: cleanup.policy=compact держит гарантированно не более одной живой записи на каждый ключ — самую свежую, а всё, что было записано под тем же ключом раньше, log cleaner физически удаляет. Удаление ключа целиком делается tombstone-записью — записью с тем же ключом и Value=nil (именно nil, а не пустой массив байт: null-значение в Kafka — это и есть маркер удаления, отличный от «пустого, но существующего» значения).
Живой прогон: топик demo-storage-compact, cleanup.policy=compact, min.cleanable.dirty.ratio=0.01 (почти любой объём «грязных» данных запускает компакцию), delete.retention.ms=100 (недолгая жизнь tombstone-маркера после компакции), max.compaction.lag.ms=8000 (страховка — компакция гарантированно случится в пределах 8 секунд даже без срабатывания dirty-ratio).
case "compact-setup":
t := defTopic(*topic, "demo-storage-compact")
configs := map[string]*string{
"cleanup.policy": strPtr("compact"),
"segment.bytes": strPtr(itoa(*segmentBytes)),
"segment.ms": strPtr("3600000"), // roll только по размеру
"min.cleanable.dirty.ratio": strPtr("0.01"),
"delete.retention.ms": strPtr("100"),
"min.compaction.lag.ms": strPtr("0"),
"max.compaction.lag.ms": strPtr("8000"),
}
recreateTopic(seeds, t, int32(*partitions), int16(*rf), configs)8 ключей (biz-key-1…biz-key-8), по 4 обновления на каждый (32 записи), плюс 2 tombstone на ключи biz-key-7 и biz-key-8 — итого 34 бизнес-записи. Пока сегмент с этими записями активен — компакция физически невозможна, что бы ни было настроено: log cleaner видит только закрытые сегменты. Consume «до» честно возвращает все 34 записи — ни одна версия ещё не схлопнута.
Два прохода log cleaner’а: почему tombstone не пропадает сразу
Живой прогон вскрыл деталь, которую легко упустить, читая только документацию: физическое удаление tombstone-маркера требует второго прохода log cleaner’а по одному и тому же сегменту (это задокументированное поведение, тонкость реализации из KAFKA-3137, не баг стенда). Первый проход схлопывает версии живых ключей до последней и оставляет tombstone нетронутым — это предохранитель против гонки с отстающими consumer’ами: если удалить tombstone сразу, consumer, который ещё не дочитал старые версии ключа, никогда не увидит сигнал «этот ключ был удалён», и в его локальной копии состояния ключ останется висеть вечно. Только когда delete.retention.ms истёк и у лога появились новые «грязные» данные для повторного запуска cleaner’а, второй проход удаляет tombstone уже физически. Если новых данных нет, второй проход попросту не запускается сам — на тихом топике tombstone способен провисеть без изменений минутами.
Сценарий стенда форсирует оба прохода явно, отправляя два раунда «балластных» записей (filler) — каждый раунд заставляет активный сегмент закрыться (roll), делая его видимым cleaner’у:
client_run -scenario=compact-produce-business -topic="$topic" -pad-bytes=5000 -rounds=4
echo "--- consume ДО: сегмент активен, cleaner не может его тронуть ---"
client_run -scenario=compact-consume -topic="$topic" -label=до-компакции -idle=3s
echo "--- filler раунд 1: форсирует roll -> проход 1 ---"
client_run -scenario=compact-produce-filler -topic="$topic" -pad-bytes=5000 -filler-n=250 -filler-start=0
sleep 6
client_run -scenario=compact-consume -topic="$topic" -label=после-прохода-1 -idle=3s
echo "--- filler раунд 2 (нумерация продолжена с 250): форсирует roll -> проход 2 ---"
client_run -scenario=compact-produce-filler -topic="$topic" -pad-bytes=5000 -filler-n=250 -filler-start=250
sleep 8
client_run -scenario=compact-consume -topic="$topic" -label=после-компакции -idle=3s -assertfiller-start продолжает нумерацию ключей между раундами (250, а не с нуля) намеренно: если второй раунд заново пронумеровать с 0, он тихо перезапишет те же ключи первого раунда, и cleaner схлопнёт filler-записи как обычные обновления — балласт не выполнит свою роль форсировать новый roll. Реальные числа по проходам:
до-компакции: всего читаемо 34 записей (32 обновления + 2 tombstone, сегмент ещё активен)
после-прохода-1: всего читаемо 258 записей
(250 filler + 6 живых ключей по 1 записи + 2 tombstone ЕЩЁ ЖИВЫ, offset 32/33)
после-компакции: всего читаемо 506 записей
(250 + 250 filler + 6 живых ключей по 1 записи)
[assert] OK: 6 живых ключей — ровно по 1 записи (последняя версия, маркер "-round3-");
2 tombstone-ключа — отсутствуют полностью32 обновления по 8 ключам схлопнулись в 6 итоговых записей (по одной, последней, на живой ключ), а biz-key-7 и biz-key-8 после второго прохода не оставили после себя ни значения, ни самого tombstone-маркера — как будто этих ключей никогда не было. Java-клиент на идентичном сценарии дал тот же результат.
8 ключей x 4 версии + 2 tombstone
34 бизнес-записи, компакция невозможна"] A -->|"roll + проход 1"| B["6 живых ключей x 1 запись
2 tombstone ЕЩЁ присутствуют
(предохранитель от гонки с consumer'ом)"] B -->|"roll + delete.retention.ms истёк + проход 2"| C["6 живых ключей x 1 запись
2 tombstone-ключа удалены физически"] style A fill:#eee8d8,stroke:#8b7355 style B fill:#f9f3e3,stroke:#8b7355 style C fill:#c9e4c5,stroke:#5b8a5e
flowchart LR
A["Активный сегмент
8 ключей x 4 версии + 2 tombstone
34 бизнес-записи, компакция невозможна"]
A -->|"roll + проход 1"| B["6 живых ключей x 1 запись
2 tombstone ЕЩЁ присутствуют
(предохранитель от гонки с consumer'ом)"]
B -->|"roll + delete.retention.ms истёк + проход 2"| C["6 живых ключей x 1 запись
2 tombstone-ключа удалены физически"]
style A fill:#eee8d8,stroke:#8b7355
style B fill:#f9f3e3,stroke:#8b7355
style C fill:#c9e4c5,stroke:#5b8a5e
Retention или compaction: как выбирать
Выбор между политиками — это в первую очередь вопрос про смысл данных в топике, а не про технические настройки:
- Поток событий (аудит-лог, клики, заказы как факты «что произошло») — каждая запись самоценна сама по себе, старая версия не заменяет новую, она просто предшествует ей во времени. Здесь уместен
cleanup.policy=delete: записи стареют и удаляются целиком по времени или объёму, потому что важна именно история, пока она нужна, а не «последнее состояние». - Снимок текущего состояния по ключу (профиль пользователя, последняя цена товара, конфигурация сервиса) — топик используется как changelog: важно только последнее значение по каждому ключу, а вся промежуточная история — шум, который можно безопасно выбросить. Здесь работает
cleanup.policy=compact: независимо от того, сколько времени прошло, читатель, подключившийся впервые, восстановит актуальное состояние по каждому ключу за один проход с начала топика.
Обе политики можно скомбинировать (cleanup.policy=compact,delete) — типичный случай, когда снимок состояния всё равно не должен жить вечно: compaction схлопывает версии по ключу, а retention сверху ограничивает, как долго вообще хранится история, включая последнюю версию каждого ключа.
Компрессия: none, gzip, lz4, zstd на одних и тех же данных
Компрессия в Kafka работает по батчу, а не по отдельной записи: producer собирает пачку сообщений и сжимает её целиком одним блоком, который так и лежит на диске брокера сжатым (брокер не распаковывает и не перепаковывает батчи при обычной записи). Отсюда прямое следствие для методологии замера: чтобы увидеть реальный эффект кодека, отправка должна быть асинхронной и батчевой — последовательный ProduceSync, ожидающий ответа перед каждой следующей записью, почти всегда создаёт батчи размера 1, на которых gzip/lz4/zstd почти ничего не сжимают, а иногда даже «раздувают» запись служебными заголовками кодека.
Живой прогон: 3000 записей на кодек, ~300 байт JSON-подобного значения каждая, асинхронная батчевая отправка (Produce + один Flush в конце). Размер партиции на диске измерен через kafka-log-dirs.sh — официальный способ узнать физический размер без прямого ls:
zstd дал лучший результат (16843–17039 байт в двух прогонах, разброс около 1%), gzip — между zstd и lz4 (28638 байт), lz4 — заметно хуже gzip на этих данных (41557 байт), none — худший вариант по определению (930016 байт, база сравнения). На Java (kafka-clients, тот же сценарий) порядок кодеков совпал один в один — zstd лучший, none худший, gzip между zstd и lz4 — но абсолютные числа заметно другие: none 930538, gzip 46409 (~20.1×), lz4 63790 (~14.6×), zstd 29635 (~31.4×). Разница в коэффициентах между клиентами — не разница в эффективности кодеков, а эффект размера батча: у franz-go и kafka-clients разные дефолты linger.ms/batch.size, и компрессия работает именно по батчу — крупнее батч, больше избыточности кодек успевает найти. Порядок кодеков — воспроизводимое свойство алгоритмов, абсолютные коэффициенты — свойство конкретной настройки батчинга, и эти два факта не стоит путать.
Асимметрия дефолтов между клиентами здесь работает в обе стороны: если franz-go по умолчанию сжимает Snappy (см. раздел про retention выше), то kafka-clients по умолчанию не сжимает вовсе — compression.type не задан значит none, явно выключать ничего не нужно.
static Properties producerProps(String compressionType) {
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Kafka.BOOTSTRAP);
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());
p.put(ProducerConfig.ACKS_CONFIG, "all");
if (compressionType != null && !compressionType.equals("none")) {
p.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, compressionType);
}
return p;
}Dynamic и static конфиги: не всё можно поменять на лету
Не каждый broker-конфиг можно изменить командой kafka-configs.sh --alter без пересоздания кластера — и разница между «dynamic» и «static» конфигом обнаруживается только тогда, когда пробуешь поменять его вживую, а не читая список параметров в документации. log.cleaner.backoff.ms (как часто log cleaner проверяет топики на «грязные» данные) — dynamic, применяется мгновенно как cluster-wide default:
set_cleaner_backoff() {
local ms="$1"
docker exec kafka-cookbook-1 /opt/kafka/bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type brokers --entity-default --alter --add-config log.cleaner.backoff.ms="$ms" >/dev/null
echo "[ops] log.cleaner.backoff.ms=$ms (dynamic, cluster-wide default)"
}Но у этого стенда, как и у всей его compose-топологии, нет персистентного volume под /tmp/kafka-logs: при пересоздании контейнеров (не рестарте — именно пересоздании) весь storage форматируется заново, и ранее выставленный dynamic-конфиг теряется вместе с ним. Поэтому сценарий компакции переустанавливает log.cleaner.backoff.ms=1000 (вместо дефолтных 15000мс — иначе демонстрация растянется на минуты) в начале каждого прогона, а не полагается на состояние, оставшееся от предыдущего запуска.
log.retention.check.interval.ms (как часто брокер вообще проверяет топики на предмет истёкшего retention) устроен жёстче: это не dynamic-alterable конфиг вовсе. Попытка изменить его через kafka-configs.sh --alter на живом кластере заканчивается прямым отказом брокера — Cannot update these configs dynamically. Единственный способ поменять его — статическая правка compose.yml (KAFKA_LOG_RETENTION_CHECK_INTERVAL_MS=2000) и пересоздание кластера; дефолт брокера — 300000мс (5 минут), и без явного занижения ждать эффекта retention на живой демонстрации пришлось бы соответствующе дольше. Именно за счёт этой статической правки retention выше сработал за секунды, а не за пять минут.
Тот же паттерн — часть конфигов dynamic и меняется мгновенно, часть жёстко static и требует пересборки кластера — повторяется и в операционном слое Kafka за пределами хранения: leader.imbalance.check.interval.seconds, отвечающий за авто-ребаланс лидеров, тоже не dynamic-alterable, и это разобрано в статье про эксплуатацию и тюнинг KRaft.
Tiered storage: архитектура и честный итог
Tiered storage (KIP-405) — идея разделить лог партиции на два яруса: «горячие» недавние сегменты остаются на локальном диске брокера, а «холодные» уходят в удалённое объектное хранилище (S3-совместимое или аналогичное) через плагин RemoteStorageManager и подсистему RemoteLogManager. Мотивация — разорвать связь между тем, сколько истории нужно хранить, и тем, сколько локального SSD физически стоит на брокере: retention может измеряться неделями или месяцами, а объём локального диска — оставаться скромным, потому что старые сегменты не занимают его вовсе. Чтение старых данных при этом продолжает работать прозрачно для consumer’а — просто медленнее, потому что запрос уходит в удалённое хранилище, а не читается с локального SSD.
Демонстрация этого механизма на живом кластере — честно не удалась, и причина не в ошибке настройки, а в самом образе. remote.log.storage.system.enable — static broker-конфиг (тоже не dynamic-alterable), и на использованном кластере он выключен (false); remote.log.storage.manager.class.name не сконфигурирован (null) — плагина, который умел бы физически перекладывать сегменты в объектное хранилище, попросту нет. Поиск по официальному образу apache/kafka:4.3.1 подтверждает: подходящего jar-файла в поставке нет вообще.
scenario_tiered() {
echo "--- remote.log.storage.system.enable (текущее значение брокера) ---"
docker exec kafka-cookbook-1 /opt/kafka/bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type brokers --entity-name 1 --describe --all 2>&1 \
| grep -i "remote.log.storage.system.enable" || echo "(конфиг не найден в --describe --all)"
echo "--- поиск классов RemoteStorageManager/RemoteLogMetadataManager в образе ---"
docker exec kafka-cookbook-1 sh -c \
'find /opt/kafka/libs -iname "*tiered*" -o -iname "*remote-storage*" 2>/dev/null' || true
}Поиск по /opt/kafka/libs не нашёл ни одного jar-файла с подходящим именем — только базовый фреймворк хранения (kafka-storage-4.3.1.jar, kafka-storage-api-4.3.1.jar), который определяет интерфейсы RemoteStorageManager/RemoteLogMetadataManager, но не содержит ни одной их реализации. Референсный плагин (LocalTieredStorage) существует в исходниках проекта — как тестовый инструмент для разработчиков самой Kafka — но не входит в релизный образ. Решение было честным: не пытаться имитировать перекладывание сегментов вручную поверх работающего общего кластера (это создало бы риск для остальных живых стендов серии) и зафиксировать статус как есть, без подмены реального поведения инсценировкой. Архитектура и мотивация KIP-405 от этого не становятся менее реальными — просто конкретно на этом образе продемонстрировать их физическую работу нельзя, и об этом стоит знать до того, как проектировать retention-политику, рассчитывая на прозрачную выгрузку в объектное хранилище «из коробки».
Партиция на диске — это не абстракция: конкретные файлы сегментов, конкретный минимум segment.bytes, конкретный порядок из двух проходов cleaner’а для tombstone. Retention отвечает на вопрос «сколько истории хранить», compaction — «какое состояние держать по каждому ключу», и оба механизма можно комбинировать, если топик одновременно и снимок состояния, и не бесконечен. Компрессия работает по батчу, а не по записи, и путать эффект кодека с эффектом настроек батчинга — методологическая ошибка, которую легко допустить, сравнивая абсолютные числа между разными клиентами вместо порядка величин.
Дальше в серии: как эти сегменты и их размеры влияют на время восстановления брокера после падения и на дисковый бюджет кластера в проде — в статье про эксплуатацию и тюнинг KRaft. Идемпотентность и транзакции, разобранные в контексте durability в статье про репликацию, опираются на ту же модель партиции, что и здесь, — развёрнуты дальше в статье про exactly-once semantics. Лог как источник для внешних систем — в статьях про CDC с Debezium и про потоковую обработку Kafka Streams/Flink, где ровно те же сегменты, retention и compaction определяют, что физически доступно для повторного чтения. Для навигации по всей теме messaging — карта messaging-landscape-map.
Комментарии