
Когда речь заходит о потоковой обработке, временные окна обычно рассматриваются как основной механизм агрегирования событий. Но просто «подождать пять дней» не получится. В потоковой архитектуре ожидание — это не обычный таймер. Необходимо учитывать задержки доставки, порядок поступления сообщений и другие пограничные состояния.
В моём случае агрегат формировался из данных нескольких независимых источников, каждый из которых работал по своим правилам. Одни системы публиковали события практически сразу после их возникновения. Другие могли прислать информацию спустя несколько часов или даже дней. Где-то происходила повторная отправка сообщений, а где-то — исправление уже ранее переданных данных.
Каждый источник передавал собственную часть информации в своём формате. На этапе обработки все сообщения приводились к единой структуре — JSON-объекту, описывающему кассовый чек. В рамках проекта рассматривалось два варианта агрегации: на уровне отдельных чеков и на уровне Z-отчёта. В этой статье речь пойдёт об агрегате на уровне Z-отчёта.
При этом не требовалось сохранять полное содержимое каждого чека. Для решения поставленной задачи было достаточно знать количество чеков и их общую сумму. Поэтому в агрегате хранились только счётчик чеков и суммарная стоимость, без детальной информации по каждому документу.
По результату решения проблемы мы сформировали более сложную задачу в потоковой обработки — Как корректно собрать итоговое состояние, если мы не знаем, когда именно придёт последнее сообщение?
Чтобы понять, почему этот вопрос оказался ключевым, сначала разберём, как Kafka Streams работает со временем и временными окнами.
Немного теории: что такое время и окна в понимании Kafka Streams?
Прежде чем перейти к проблеме длинных временных окон, разберём базовые концепции потоковой обработки. Постараюсь не сильно подробно углубляться в теорию. Когда мы говорим об обработке событий во времени, существует несколько понятий времени:
Event Time — время, когда событие фактически произошло. Например, чек пробит в 12:00, но в систему попал в 14:00 из-за задержки. Для анализа важен именно момент 12:00, так как он отражает бизнес-событие.
Ingestion Time — время, когда сообщение попало в систему (например, в Kafka поступил в 14:00). Используется для анализа задержек доставки и работы интеграций.
Processing Time — время обработки события системой. Зависит от нагрузки и скорости сервиса и чаще применяется для оценки производительности, а не бизнес-логики.
Kafka Streams поддерживает несколько основных типов окон:
Tumbling Window (кувыркающиеся окна)

Фиксированные непересекающиеся интервалы.
Например:
окно в 5 минут: 12:00–12:05
следующее окно: 12:05–12:10
Каждое окно имеет протяженность 5 минут и перемещается вперед на 5 минут. Каждое событие попадает только в одно окно. Подходит для периодических отчётов и простых агрегатов.
Hopping Window (прыгающие окна)

Окна фиксированного размера, которые могут пересекаться. Например, каждое окно имеет протяженность 5 минут, но каждую 1 минуту происходит сдвиг. Есть перекрытие в 4 минут у соседних окон. Получаем:
12:00–12:05
12:01–12:06
12:02–12:07
Используется, когда нужны агрегаты с перекрытием периодов.
Sliding Window (скользящие окна)

Окно непрерывно скользит по таймлайну и включает только события, происходящие в интервале, заданном фиксированным размером окна. Окно формируется относительно каждого события и может пересчитываться при появлении новых данных. В отличие от фиксированных окон, оно не привязано к определённым временным границам. Подходит для сценариев, где важна реакция на изменения состояния в реальном времени. При размере окна в 3 минуты получаем:
12.00 - 12.03
12.01 - 12.04
12.05 - 12.08
Session Window (сессионные окна)

Окно существует только пока есть активность. Когда поток событий прекращается на заданный период, сессия считается завершённой. Сеансы будут продолжать увеличиваться в размер, пока не наступит период бездействия, превышающий заданный интервал. Например: Пользователи совершали действия, далее одну минуту ничего не происходило — сессия закрылась.
При ознакомлении технологией временных окон авторы убеждали, что переживать за недосборку агрегата не стоит. В Kafka Streams есть механизм учета неупорядоченных данных. Все риски нивелируются за счет Grace периода. Поэтому виделось, что временные окна это практически обычная группировка, только поверх потока. Но как всегда есть особенности, с которыми нужно уметь работать.
Первая реализация: «давайте просто сделаем окно на 5 дней»

После анализа требований и изучения технологий временных окон первый вариант решения выглядел достаточно очевидным — использовать временное окно в период ожидания в пять дней. Логика была простой: бизнес определил, что финальные данные могут появляться с задержкой до пяти дней после фактического события. Значит, потоковая обработка должна удерживать состояние этого периода и финализировать агрегат только после завершения допустимого времени ожидания.
Выбор пал на использование Tumbling Window с размером окна в один день и Grace периодом в 4 дня. Такой подход соответствовал бизнес-логике: каждый календарный день должен формировать собственный агрегат, а дополнительные четыре дня давались источникам для доставки запоздалых данных. Это позволяло учитывать изменения в данных и формировать агрегаты относительно даты события (Event Time). В нашем случае каждое окно было связано с датой начала периода:
01.01 – 05.01 → агрегат по Event Time 01.01
02.01 – 06.01 → агрегат по Event Time 02.01
03.01 – 07.01 → агрегат по Event Time 03.01
Событие с Event Time равный 01.01 попадало в соответствующие окна независимо от того, когда оно физически пришло в систему. Пока событие укладывалось в допустимый период ожидания, Kafka Streams могла корректно учесть его при расчёте агрегата. Grace Period позволял принимать задержки доставки после окончания основного временного интервала.
На тестовых данных решение показывало отличные результаты. Симуляция отражала ожидаемый сценарий. События приходили с корректными временными метками, задержки были предсказуемыми, источники передавали согласованные данные, а агрегаты формировались стабильно. На этом этапе решение выглядело достаточно простым и полностью соответствующим требованиям.
Основное предположение было следующим: если данные не пришли в момент возникновения события, то они обязательно появятся в течение пяти дней и успеют попасть в соответствующий агрегат. Оставалось проверить, насколько хорошо этот подход будет работать на реальных данных от всех источников.
Реальность production-данных

На тестовых данных пятидневное окно выглядело идеальным решением. Настоящие проблемы появились только после запуска в прод, где поведение источников оказалось значительно менее предсказуемым. В реальной системе регулярно возникали ситуации, которые невозможно было воспроизвести на тестовых данных:
события приходили с задержкой в несколько дней;
ранее отправленные данные переотправлялись с изменениями;
одно бизнес-событие существовало в нескольких версиях;
порядок сообщений нарушался.
Главная проблема заключалась не в объёме данных, а в отсутствии финальности. После закрытия окна продолжали поступать новые или исправленные события, поэтому агрегат нельзя было считать окончательно сформированным. Вместо одного результата система постепенно начинала создавать несколько версий одного и того же периода.
Попытка выбрать «правильную» модель времени также не помогла. Event Time зависит от качества данных источников, а Processing Time и Ingestion Time — от скорости доставки. Ни одна из моделей не отвечает на главный вопрос: когда можно перестать ждать новые данные?
Из-за этого усложнилось и тестирование. На тестовых данных результат всегда был определен, а в проде любой агрегат мог измениться после получения позднего события. Проверить корректность становилось значительно сложнее.
Стоит отдельно отметить, что Exactly Once не решает эту проблему. Он гарантирует отсутствие технических дублей при обработке, но не защищает от некорректных или исправленных бизнес-данных, поступающих из внешних систем.
Кроме того, длинные окна требуют длительного хранения состояния в State Store. При сбоях восстановление занимает больше времени, а цена потери накопленного состояния существенно возрастает. Именно эта проблема впоследствии привела нас к отдельной задаче — управлению жизненным циклом состояния Kafka Streams и настройке самоочистки RocksDB. В следующей статье я отдельно разберу, как устроено хранение состояния, какую роль играет changelog-топик, как он связан с локальным State Store и почему одного retention недостаточно для управления размером хранилища.
Итог
Пятидневное окно оказалось не просто параметром Kafka Streams, а архитектурным решением, затрагивающим модель времени, качество данных, хранение состояния, восстановление после сбоев и тестирование.
Главное ограничение оказалось не в Kafka Streams, а в природе самих данных: если невозможно определить момент, после которого события больше не изменяются, то временное окно перестаёт быть надёжным механизмом финализации агрегата.
Именно поэтому мы отказались от длинных временных окон и перешли к модели, в которой жизненным циклом агрегата управляет бизнес-логика, а не время.
Эволюция подхода к агрегации данных

После анализа ограничений временных окон стало понятно, что проблема заключалась не в выборе конкретного размера окна или настройки Grace Period. Мы пытались решить временным механизмом задачу, которая по своей природе не являлась временной. В оконной модели система должна была ответить на вопрос: «Когда можно считать, что данные за период больше не изменятся?». Для этого требовалось определить момент финальности данных, выбрать подходящую модель времени и ожидать закрытия окна. Однако реальные источники не могли гарантировать такой момент. Данные могли приходить с задержкой, изменяться задним числом или повторно отправляться после корректировок.
Поэтому мы изменили сам принцип определения готовности агрегата. Вместо вопроса «Сколько времени прошло с момента события?» появился другой вопрос: «Получены ли все необходимые части объекта для формирования результата?». Вместо привязки к дате был введён составной бизнес-ключ: “смена + номер фискального накопителя”. Фискальный накопитель являлся уникальным номером, а номер смены определяет закрытость промежутка времени без самого времени.
Мы перешли к новой модели: Business Key → State Store → Aggregate State → Complete Result. Теперь агрегат существовал не фиксированный период времени, а до тех пор, пока не выполнялось условие его завершения.
В процессе участвовали три источника данных:
ОФД (оператор фискальных данных) — данные Z-отчётов. Документ, который порождает касса;
DWH (корпоративное хранилище) — искусственный аггрегат из данных чеков. Подобие Z-отчета;
RMS (оперативное хранилище) — искусственный аггрегат из данных операций. Подобие Z-отчета.
Важно отметить, что только Z-отчёт ОФД являлся первичным бизнес-документом. Объекты DWH и RMS представляли собой производные агрегаты, которые строились внутри нашей системы. Именно поэтому триггером завершения процесса был выбран документ ОФД. Его получение означало, что смена действительно закрыта, поэтому на шестой день система инициировала его получение и использовала как финальный сигнал для проверки полноты агрегата.
Каждый источник передавал только свою часть итогового объекта. Kafka Streams должна была собрать эти части в единый агрегат независимо от порядка поступления и времени доставки. Для этого вместо временной агрегации была реализована stateful-агрегация по ключу объекта:
var result = input .groupByKey() .aggregate( () -> ReportObjectRecord.newBuilder().build(), (reportKey, value, aggregateGroup) -> reportService.collectAggregate(value, aggregateGroup), MaterializedwithRetention .<ReportKeyRecord, ReportObjectRecord, KeyValueStore<Bytes, byte[]>>as(appProperties.aggregateStore()) .withRetention(appProperties.retention()) .withCachingDisabled() .withLoggingEnabled( Map.of( TopicConfig.RETENTION_MS_CONFIG, String.valueOf(appProperties.retention().toMillis()), TopicConfig.CLEANUP_POLICY_CONFIG, TopicConfig.CLEANUP_POLICY_COMPACT) ) ) .filterNot((k, v) -> reportService.isSkipped(v)) .toStream() .process(() -> new StoreProcessor(storageService, reportService));
При поступлении первого сообщения для ключа Kafka Streams создавала новый пустой агрегат. Последующие сообщения с тем же ключом обновляли существующее состояние, добавляя недостающие части объекта. В отличие от оконной агрегации, состояние агрегата больше не зависело от временного интервала. Жизненный цикл объекта определялся не датой, а фактом получения необходимых данных. Агрегат мог существовать в State Store столько времени, сколько требовалось для получения всех частей объекта. Ограничением становился только срок хранения состояния в State Store.
Логика объединения данных
Каждое сообщение содержало часть итогового объекта. В зависимости от типа источника соответствующий блок добавлялся в существующий агрегат. Основная логика объединения:
public ReportObjectRecord collectAggregate( ReportObjectRecord reportObject, ReportObjectRecord aggregateGroup ) { if (isComplete(aggregateGroup)) { aggregateGroup.setSkipped(true); } // Данные ОФД if (reportObject.getOfdReport() != null) { int ofdReportCounter = aggregateGroup.getOfdReportCounter() == null ? 0 : aggregateGroup.getOfdReportCounter(); aggregateGroup.setOfdReportCounter(ofdReportCounter + 1); aggregateGroup.setOfdReport( reportObject.getOfdReport() ); if (aggregateGroup.getOfdReportCounter() == 1) { storageService.updateOfdReportCounter(reportObject); } } // Данные DWH if (reportObject.getDwhReport() != null) { int dwhReportCounter = aggregateGroup.getDwhReportCounter() == null ? 0 : aggregateGroup.getDwhReportCounter(); aggregateGroup.setDwhReportCounter(dwhReportCounter + 1); aggregateGroup.setDwhReport( reportObject.getDwhReport() ); } // Данные RMS if (reportObject.getOperationReport() != null) { int operationReportCounter = aggregateGroup.getOperationReportCounter() == null ? 0 : aggregateGroup.getOperationReportCounter(); aggregateGroup.setOperationReportCounter(operationReportCounter+ 1); aggregateGroup.setOperationReport( reportObject.getOperationReport() ); } return aggregateGroup; }
Определение полноты агрегата
Проверка готовности перестала определяться временем и стала формироваться логикой бизнес-условия. Агрегат считался полным при наличии всех трех источников. Логика работы проста: каждый новый источник дополняет существующий агрегат. Например, после получения данных только от RMS объект остаётся неполным. После прихода любого нового объекта выполняется повторная проверка, и если все необходимые части присутствуют, агрегат считается готовым.
public boolean isComplete(ReportObjectRecord reportObject) { boolean hasOfd = reportObject.getOfdReport() != null; boolean hasDwh = reportObject.getDwhReport() != null; boolean hasRms = reportObject.getOperationReport() != null; return hasOfd && hasDwh && hasRms;
После сборки поток разделяется на полные и неполные объекты.
return result.branch( (key, value) -> reportService.isComplete(value), (key, value) -> !reportService.isComplete(value) );
Оба агрегата складываются в два разных топики Kafka. Неполные агрегаты продолжали находиться в состоянии ожидания отображая журнал изменений. Полные складываются в топик только при наличии всех трех источников. Далее они передаются дальше по процессу.
Ограничения и результаты
Переход на stateful-агрегацию решил основную проблему — определение полноты данных перестало зависеть от времени. Однако при проектировании пришлось отдельно решить ещё два вопроса.
Первый — обработка дубликатов. Возможность многократного поступления одного и того же объекта существует практически в любой распределённой системе. Вместо усложнения логики агрегации мы решили перенести эту задачу на более ранний этап обработки. Перед формированием чека был реализован Quality Gate, который проверяет корректность и полноту входящих данных. Кроме того, между участниками интеграции был согласован контракт: консистентность системы достигается за счёт исключения ручных изменений уже сформированных объектов. Если чек содержит все обязательные поля, он считается неизменяемым и передаётся дальше в процесс агрегации.
Второе ограничение связано со сроком хранения состояния в State Store. Если одна из частей агрегата поступит позже периода хранения (в нашем случае — спустя более пяти дней), ранее накопленное состояние уже будет удалено. Это осознанное ограничение. Для защиты от повторной обработки была реализована дополнительная проверка на стороне источника. Важно отметить, что наш сервис не создаёт новый чек самостоятельно, а лишь передаёт информацию о потенциальной необходимости его формирования. Непосредственно перед созданием чека источник проверяет, существует ли он уже в системе. Если чек найден, повторное создание не выполняется, что исключает появление дубликатов даже в случае позднего поступления данных.
Основное отличие новой модели от временных окон заключается в том, что она изменила сам принцип определения готовности данных. В оконной модели система отвечала на вопрос: Достаточно ли прошло времени, чтобы считать данные завершёнными? В новой реализации вопрос стал другим: Получены ли все необходимые части бизнес-объекта?
Новая же модель перевела вопрос из временной плоскости в бизнесовую: Получены ли все необходимые части бизнес-объекта? В нашем случае нам удалось определить триггер. Это наличие всех 3 источников. Вместо ожидания окончания периода система стала отслеживать фактическое состояние агрегата. Пока отсутствует хотя бы один источник, объект остаётся незавершённым и хранится в State Store.
Такой подход имеет несколько преимуществ.
исчезла зависимость от качества временных меток. Ошибка Event Time больше не могла привести к попаданию объекта в неправильный период;
исчезла необходимость угадывать размер периода ожидания. Пять дней в первоначальной реализации были не бизнес событием, а техническим предположением: мы ожидали, что за этот период придут все необходимые данные;
агрегат стал более прозрачным с точки зрения состояния. В любой момент можно определить, чего именно не хватает. Каждое изменение записывался в топик неполных агрегатов формируя журнал изменений.
Важно отметить, что нет проблемы в самих временных окнах Kafka Streams. Они являются полезным инструментом, но решают другую задачу. Их основное назначение — разделение непрерывного потока событий на временные интервалы для последующего анализа, расчётов и обработки данных по периодам. Временное окно нельзя использовать как механизм определения полноты бизнес-агрегата. Завершение должно определяться не тем, сколько прошло времени, а наличием всех необходимых данных.

