Когда State Store становится проблемой

Одной из сильных сторон Kafka Streams является возможность хранить состояние приложения локально. Благодаря State Store сервис может агрегировать события, выполнять join’ы, хранить промежуточные результаты вычислений и мгновенно находить ранее обработанные данные без обращения к внешней базе данных. По умолчанию в качестве локального хранилища используется RocksDB — встроенная key-value база данных, расположенная на диске. Она обеспечивает высокую производительность, может сохранять состояние между перезапусками приложения и автоматически восстанавливать его из changelog topic в случае потери локального состояния. Именно это делает Kafka Streams настолько удобным инструментом для построения stateful-сервисов.
Но есть одна особенность, о которой редко задумываются в начале проекта. State Store отлично умеет хранить данные, но совершенно не знает, когда их пора удалить. Если приложение однажды записало объект в Store, он останется там до тех пор, пока приложение самостоятельно его не удалит. Никакого встроенного TTL для обычного State Store в Kafka Streams нет. На небольших объёмах это практически незаметно. Однако если сервис работает месяцами и ежедневно обрабатывает миллионы сообщений, локальное состояние начинает непрерывно расти. RocksDB постепенно превращается в архив всех когда-либо обработанных данных, большая часть которых уже никогда не понадобится.
Эта статья является логическим продолжением моей предыдущей публикации «Пять дней ожидания: опыт длинных временных окон Kafka Streams». Если там основной темой была реализация агрегации с длительным ожидания данных, то здесь речь пойдёт о дополнительных сложностях в переполнении State Store и реализации механизма самоочистки.
Как эта проблема проявилась у нас
В нашем продукте существует множество State Store. Наиболее наглядно заполнение проявилось в сервисе Verification-hub. В него каждый день складывается в среднем по 18 млн. чеков. К тому же в одном сервисе существует несколько Store использующиеся сразу для нескольких задач.
Каждый новый чек увеличивал объём хранилища. Бизнес-процес обработки данных заканчивался и обращение к ним повторно уже не требуется. При этом данные продолжали оставаться в RocksDB. Первое время это никак не ощущалось. Сервисы стабильно работали, а объём State Store рос формируя бомбу замедленного действия. Спустя какое то время эксплуатации начали проявляться вполне ожидаемые последствия:
увеличился размер локального RocksDB и следствие стоимость инфры;
выросло время восстановления State Store после рестарта приложения;
увеличилась нагрузка на диск во время compaction;
менее заметно, но стали тяжелее операции чтения и записи;
увеличилось потребление памяти и дискового пространства.
Оказалось, что с этим сталкивались не только мы. Коллеги из других команд, использующих Kafka Streams, также наблюдали неконтролируемый рост State Store. В некоторых командах объём RocksDB становился настолько большим, что начинали обсуждать переход на другие технологии. На самом деле проблема была не в Kafka. Управление жизненным циклом состояния является ответственностью инженера и об этом надо задумываться сразу.
Почему обычный TTL не подходит?

Кто-то может скажет — “Нужно просто удалять старые записи по TTL. Ведь Kafka предоставляет возможность настройки параметра retention.ms”. Но к сожалению, это не всегда может помочь. Параметр retention.ms относиться к управлению временем хранения данных в changelog топике . Он никак не удаляет записи из локального RocksDB State Store. Это два независимых хранилища. Давай разберемся. На практике я выделил два сценария.
Первый — локальный State Store внутри Pod. Если объём относительно небольшой, то его можно хранить в самом Pod. В этом случае потеря локального RocksDB не является серьёзной проблемой. При обновлении приложения, перезапуске Pod или переносе его на другую ноду локальное хранилище будет создано заново, а Kafka Streams автоматически восстановит его содержимое из changelog топика. Время восстановления зависит от объёма локального state, размера changelog и скорости восстановления. Для небольших State Store оно может быть относительно быстрым. Но при больших объёмах восстановление существенно отодвигает время запуска бизнес задачи. В идеале все равно стоит реализовать механизм самоочистки. Бесконечно растущий Store хорошей практикой я бы не назвал. Но если объём состояния остаётся небольшим, а инфраструктурные работы выполняются в регламентные окна, такую задачу вполне можно оставить как технический долг.
И второй, который более реалистичный — у тебя большой State Store. Если команда и ты решили использовать Kafka Streams, то видимо есть потребность в обработки больших данных в потоке. А это значит объекты “весят” внушительно. В этом случае восстановление RocksDB из changelog может занимать часы. Всё это время сервис либо недоступен, либо работает с ограниченной производительностью.
В таком случае корректно выносить хранилище в Persistent Volume, а приложение разворачивают как StatefulSet. В отличие от Deployment, StatefulSet в сочетании с volumeClaimTemplates позволяет привязать отдельный PersistentVolumeClaim к каждому экземпляру Pod. Kafka Streams продолжает использовать RocksDB как локальное хранилище, но каталог state.dir располагается уже не на временной файловой системе Pod, а на Persistent Volume. Благодаря этому при пересоздании Pod локальные данные RocksDB могут сохраниться, и в сценарии, когда соответствующее состояние возвращается вместе с тем же экземпляром, не потребуется полностью восстанавливать его из changelog topic. При этом Persistent Volume не отменяет механизм восстановления Kafka Streams: если локальное состояние потеряно, task переместилась на другой экземпляр или состояние оказалось недоступно, Kafka Streams по-прежнему может восстановить его из changelog topic. И только после решения вопроса с сохранностью локального состояния требуется реализовать механизм его самоочистки.
В нашем случае ситуация усложняется тем, что один бизнес-объект распределен между несколькими State Store. Помимо самого агрегата существовали индексы, результаты сверки и вспомогательные структуры, связанные между собой. И удаление записи только из одного Store приводило бы к появлению “зависших” данных в остальных. Формируя неконсистентное состояние приложения.
Поэтому потребовался механизм, который бы:
определял срок жизни данных;
запускал очистку независимо от поступления новых сообщений;
удалял связанные записи сразу из нескольких State Store;
выполнял удаление небольшими пакетами, не влияя на обработку потока.
Именно такой механизм самоочистки мы реализовали в verification-hub. В следующих разделах разберём, как он устроен и почему в его основе появился отдельный timestampStore.
Подготовка RocksDB
Прежде чем реализовывать сам механизм массовых удалений нужно реализовать кастомную конфигурация RocksDB. RocksDB построена на основе LSM-деревьев (Log Structured Merge Tree). При удалении запись не удаляется физически. Вместо этого в базу записывается специальный маркер удаления — tombstone. Во время чтения RocksDB понимает, что запись удалена, и больше не возвращает её приложению. Однако сама запись вместе с tombstone продолжает находиться в SST-файлах до выполнения compaction.
Для большинства приложений это не является проблемой. Но наш сервис регулярно удаляет десятки тысяч объектов сразу. Логически данные исчезали из State Store, а физически размер RocksDB не менялся. Более того, накопление большого количества tombstones начинало усложнять обход State Store, увеличивать количество SST-файлов и приводить к дополнительной нагрузке во время compaction. По умолчанию RocksDB не слишком активно освобождает место на диске. Поэтому перед реализацией механизма очистки нужно изменить стандартную конфигурацию RocksDB.
Первым делом мы включили периодический compaction. Настраиваем periodic compaction с интервалом 1 час. Устаревшие SST-файлы периодически становятся кандидатами на compaction, что ускоряет физическое удаление данных и tombstones после логического удаления.
// Принудительный compaction каждый час (дефолт = 0, выключен!) options.setPeriodicCompactionSeconds(3600);
Вторым шагом мы сократили интервал удаления устаревших SST-файлов. По умолчанию RocksDB очищает неиспользуемые файлы примерно раз в шесть часов. Мы сократили этот интервал до одного часа ускорив освобождение дискового пространства.
// Удаление obsolete файлов каждый час (3600 секунд × 1 000 000 микросекунд) options.setDeleteObsoleteFilesPeriodMicros(3600 * 1000000L);
Эти изменения сами по себе не удаляют устаревшие записи. Они лишь гарантируют, что после выполнения нашей очистки RocksDB действительно освободит место на диске, а не будет бесконечно накапливать tombstones. И только после того как RocksDB была подготовлена можно переходить к реализации самого механизма очистки.

Реализация механизма очистки
После того как стало понятно, что State Store будет бесконечно увеличиваться, возник главный вопрос: “Как понять, какие записи уже можно удалить?”
Может хранить время создания прямо внутри агрегата? Не думаю! Агрегат — это бизнес-сущность и добавлять в неё служебную информацию только ради очистки раздует сам файл и не ускорит работу бизнес логики. Смешивать бизнес-логику и техническую не хочется. Мы пошли другим путём. Для каждого объекта сверки появилось отдельное хранилище, содержащее только время последнего изменения записи.
Получилась следующая схема:
receiptStore //Бизнес Store ReceiptKey -> VerificationObject timestampStore //Технический Store ReceiptKey -> lastUpdateTimestamp
Таким образом, вся логика управления жизненным циклом данных оказалась изолирована от бизнес-логики сервиса.
Конфигурация нового Store получилась достаточно простой.
@Bean public StoreBuilder<KeyValueStore<ReceiptKeyRecord, Long>> timestampStateStore() { var cleanUpPolicy = String.join(",", TopicConfig.CLEANUP_POLICY_DELETE, TopicConfig.CLEANUP_POLICY_COMPACT); var changelogConfig = Map.of( TopicConfig.RETENTION_MS_CONFIG, String.valueOf(appProperties.storeRetentionDuration().toMillis()), TopicConfig.CLEANUP_POLICY_CONFIG, cleanUpPolicy ); StoreBuilder<KeyValueStore<ReceiptKeyRecord, Long>> timestampStoreBuilder = Stores .keyValueStoreBuilder( Stores.persistentKeyValueStore(appProperties.timestampStoreName()), (Serde<ReceiptKeyRecord>) null, Serdes.Long() ) .withLoggingEnabled(changelogConfig); return timestampStoreBuilder; }
Store хранит всего одно значение — Long, соответствующее времени последнего обновления агрегата. Стоит обратить внимание на настройку changelog topic.
TopicConfig.CLEANUP_POLICY_DELETE TopicConfig.CLEANUP_POLICY_COMPACT
Используется одновременно compaction и delete. Compaction позволяет со временем удалять устаревшие версии записей одного ключа и оставлять актуальное значение, уменьшая объём changelog и объём данных, которые требуется прочитать при восстановлении. Стоит помнить это не означает, что в topic в каждый момент времени физически существует только последняя версия ключа. Delete позволяет автоматически удалять устаревшие записи из самого changelog topic после истечения retention.ms.
Пока мы продолжаем говорить только о подготовке механизма очистки. Мы настроили хранение changelog topic в Kafka и начали сохранять время последнего изменения каждой записи, но сами данные из RocksDB ещё не удаляем. Эти настройки управляют только жизненным циклом changelog topic в Kafka. Локальный State Store по-прежнему будет содержать запись до тех пор, пока приложение самостоятельно её не удалит. Переходим к следующему шагу — сохранение времени каждого изменения объекта.
Каждый раз, когда агрегат обновляется, сервис одновременно обновляет и его временную метку.
/ Обновляем агрегат receiptStore.put(receiptKey, updatedAgg); // Фиксируем время последнего изменения timestampStore.put(receiptKey, context.currentSystemTimeMs());
Теперь для каждой записи известен её “возраст”. Оставалось только периодически проходить по timestampStore и удалять всё, что старше заданного времени жизни.
Почему реализация отдельного Store удобнее? Использование отдельного timestampStore дало сразу несколько преимуществ. Во-первых, бизнес-модель не пришлось изменять — агрегаты остались полностью независимыми от механизма очистки. Во-вторых, один и тот же механизм можно применять сразу к нескольким State Store. Достаточно знать ключ объекта и время его последнего изменения. И наконец, срок жизни данных оказался полностью вынесен в конфигурацию приложения. Изменить TTL теперь можно без изменения кода и структуры агрегатов.
Запуск очистки
После того как для каждой записи появился собственный timestamp, оставалось решить вторую задачу — когда запускать очистку. Очистку точно нужно отделять от обработки данных и запускать по времени. Kafka Streams имеет встроенный механизм schedule(), позволяющий выполнять периодические задачи внутри Processor API. Во время инициализации процессора регистрируется планировщик:
@Override public void init(ProcessorContext<ReceiptKeyRecord, VerificationObjectRecord> context) { ... this.context.schedule( staggeredInterval, PunctuationType.WALL_CLOCK_TIME, this::performCleanup ); }
В качестве типа триггера используется PunctuationType.WALL_CLOCK_TIME. Очистка выполняется по реальному времени, а не зависит от появления новых сообщений в Kafka. Даже если поток данных временно остановился, механизм продолжит регулярно освобождать устаревшие записи.
Jitter и ограничения
При реализации механизма очистки мы сразу подумали о том, как он будет вести себя в распределённом окружении. Сегодня сервис может работать в одном экземпляре, а завтра — в двух или четырёх Pod’ах. Каждый из них будет иметь собственный State Store и выполнять очистку. Если все экземпляры приложения запускаются одновременно (что довольно часто происходит после деплоя или перезапуска Kubernetes), то без дополнительных мер они будут выполнять очистку практически синхронно.
Но даже в нашем текущем окружении проблема уже существовала. Несмотря на то что сервис был развернут в одном Pod, Kafka Streams использовал 14 потоков при 14 партиций. Каждый поток обслуживал собственную задачу обработки данных и имел доступ к локальному RocksDB. Если все потоки одновременно запускали очистку, они начинали конкурировать за один и тот же диск и внутренние ресурсы RocksDB.
Поскольку в нашей конфигурации cleanup выполнялся каждым процессором, при одновременном срабатывании расписаний потенциально могли запускаться до 14 операций очистки одновременно. Полный обход State Store, чтение большого количества данных с диска, удаление записей и обновление внутренних структур RocksDB. Это не приводило к ошибкам, но создавало пики нагрузки, увеличивало конкуренцию за диск и внутренние блокировки RocksDB.
Поэтому ещё на этапе разработки мы заложили механизм случайного смещения времени запуска (jitter).
// Генерируем случайный джиттер до 15 минут long jitterMs = ThreadLocalRandom.current().nextLong(0, Duration.ofMinutes(15).toMillis()); // Рассчитываем итоговый интервал Duration staggeredInterval = appProperties.cleanupIntervalDuration().plusMillis(jitterMs);
Теперь каждый Processor/task запускает первую очистку в случайный момент времени, после чего продолжает работать по собственному расписанию. Например, при интервале очистки в один час расписание может выглядеть следующим образом:
Processor/task-1 cleanupIntervalDuration + 2 мин
Processor/task-2 cleanupIntervalDuration + 11 мин
Processor/task-3 cleanupIntervalDuration+ 4 мин
Processor/task-4 cleanupIntervalDuration + 14 мин
В результате дорогостоящие операции по обходу RocksDB распределяются во времени, а нагрузка на кластер становится более равномерной. На текущий момент количество экземпляров сервиса невелико, и эффект от такого подхода не сильно заметен. Однако механизм был реализован заранее как часть архитектуры масштабируемого решения. При увеличении числа Pod’ов не потребуется дорабатывать сервис — очистка уже изначально рассчитана на работу в распределённой среде.
Поиск устаревших записей
Во время запуска очистки первым делом вычисляется граница, после которой запись считается устаревшей.
long threshold = timestamp - appProperties.storeRetentionDuration().toMillis();
Параметр storeRetentionDuration определяется в в файле application (8 days). Все записи, чей timestamp меньше этого значения, должны быть удалены. Поскольку вся информация о возрасте объектов хранится в timestampStore, проход выполняется именно по нему.
// Открываем итератор по всему стору меток времени try (KeyValueIterator<ReceiptKeyRecord, Long> iterator = timestampStore.all()) { long scanStartTime = System.currentTimeMillis(); while (iterator.hasNext() && deleted < batchLimit) { scanned++; KeyValue<ReceiptKeyRecord, Long> entry = iterator.next(); // Проверяем, устарела ли запись if (entry.value < threshold) { // Замеряем время поиска до момента нахождения первой записи на удаление if (deleted == 0) { firstMatchTimeMs = System.currentTimeMillis() - scanStartTime; }
timestampStore содержит минимально возможный объём данных. Полный проход по timestampStore всё равно имеет линейную стоимость относительно количества записей. Однако поскольку Store содержит только ключ и Long, такое сканирование обычно дешевле по объёму данных и стоимости десериализации, чем аналогичный проход по Store с крупными агрегатами.
Каскадное удаление связанных данных
Оставалась ещё одна важная задача. В нашем случае объект сверки существовал сразу в нескольких State Store, каждый из которых отвечал за свою часть обработки:
State Store | Назначение |
|---|---|
receiptStore | агрегаты объектов сверки |
reportStore | результаты сверки |
indexStore | индекс для быстрого поиска чеков |
timestampStore | время последнего обновления записи |
Во время очистки происходит последовательный обход timestampStore. Каждая запись представляет собой пару ключ — время последнего обновления. Например:
ReceiptKey | Timestamp |
|---|---|
FN=123, FD=456 | 1719500000000 |
FN=123, FD=457 | 1719500100000 |
FN=123, FD=458 | 1719500200000 |
Итератор Kafka Streams возвращает сразу обе части этой записи:
KeyValue<ReceiptKeyRecord, Long> entry = iterator.next();
где: entry.key — идентификатор объекта сверки (ReceiptKeyRecord), а entry.value — время его последнего обновления.
Сначала проверяется, истёк ли срок жизни записи:
if (entry.value < threshold) {
Если условие выполняется, значит объект считается устаревшим. После этого используется уже сам ReceiptKeyRecord, чтобы найти и удалить все связанные записи из остальных State Store.
ReceiptKeyRecord keyToDelete = entry.key;
Именно благодаря тому, что timestampStore хранит тот же ключ, что и остальные State Store, он становится своеобразным каталогом для очистки. По нему определяется не только что удалить, но и где именно это удалить.
Вывод

Встроенного механизма TTL для KeyValueStore нет, поэтому без собственной стратегии очистки локальное состояние будет непрерывно расти. В нашем случае проблему удалось решить за счёт отдельного timestampStore, периодического запуска очистки, пакетного удаления записей и каскадного удаления связанных данных из всех State Store. Такой подход позволил сохранить согласованность локального состояния и сделать рост RocksDB предсказуемым и управляемым.
Главный вывод, который мы сделали в процессе разработки, довольно простой: при проектировании сервисов на Kafka Streams важно заранее продумывать не только то, как данные попадут в State Store, но и когда и каким образом они будут из него удалены. Если этого не сделать на этапе проектирования, рано или поздно проблема проявится уже в промышленной эксплуатации.
