Когда речь заходит о потоковой обработке, временные окна обычно рассматриваются как основной механизм агрегирования событий. Но просто «подождать пять дней» не получится. В потоковой архитектуре ожидание — это не обычный таймер. Необходимо учитывать задержки доставки, порядок поступления сообщений и другие пограничные состояния.
В моём случае агрегат формировался из данных нескольких независимых источников, каждый из которых работал по своим правилам. Одни системы публиковали события практически сразу после их возникновения. Другие могли прислать информацию спустя несколько часов или даже дней. Где-то происходила повторная отправка сообщений, а где-то — исправление уже ранее переданных данных.
Каждый источник передавал собственную часть информации в своём формате. На этапе обработки все сообщения приводились к единой структуре — JSON-объекту, описывающему кассовый чек. В рамках проекта рассматривалось два варианта агрегации: на уровне отдельных чеков и на уровне Z-отчёта. В этой статье речь пойдёт об агрегате на уровне Z-отчёта.
При этом не требовалось сохранять полное содержимое каждого чека. Для решения поставленной задачи было достаточно знать количество чеков и их общую сумму. Поэтому в агрегате хранились только счётчик чеков и суммарная стоимость, без детальной информации по каждому документу.