Поток событий содержит сообщения с разным временем доставки. Как механизм должен определить момент закрытия временного окна, чтобы учесть опоздавшие события?
events:
- id: A
event_time: "10:00"
arrival_time: "10:01"
- id: B
event_time: "09:58"
arrival_time: "10:05"
window: "09:55-10:00"
allowed_lateness: "5 minutes"
Для закрытия окна следует использовать время события и водяной знак (watermark), а не только время обработки. Watermark показывает, до какого момента система считает поток в основном полученным; окно закрывается после прохождения этого порога, но в пределах допустимого опоздания ещё может быть пересчитано.
В примере событие B пришло поздно, но относится к окну 09:55–10:00. При допустимом опоздании в пять минут система должна не закрывать это окно раньше watermark 10:05 либо принять корректировку уже опубликованного результата.
В распределённых потоках порядок доставки не обязан совпадать с порядком возникновения событий. Сетевые задержки, повторные попытки отправки и временная недоступность потребителей приводят к тому, что событие с ранним временем может прийти после события с более поздним временем.
Простая обработка по времени получения решает задачу низкой задержки, но искажает временные отчёты. Поэтому потоковые системы разделяют event time — время фактического события — и processing time — время его обработки, а для определения полноты данных используют watermark.
Если агрегатор закрывает окно сразу после поступления первых сообщений, позднее событие не попадёт в итог. Например, продажа, совершённая в 09:58, может быть доставлена в 10:05 и ошибочно оказаться за пределами отчёта за 09:55–10:00.
Слишком ранний watermark уменьшает задержку результата, но повышает риск неполных данных. Слишком осторожный watermark сохраняет точность, однако увеличивает задержку отчётов и объём состояния, которое нужно хранить.
Watermark — это не обязательно точная отметка последнего события. Обычно это нижняя граница времени событий, которые система ожидает увидеть с высокой вероятностью. Когда watermark превышает конец окна, окно считается готовым к закрытию.
Допустимое опоздание задаёт период, в течение которого закрытый результат ещё можно исправлять. Если позднее событие пришло в этот период, агрегатор обновляет состояние окна и публикует корректировку. После его окончания событие могут отбросить, отправить в отдельный поток ошибок или обработать через специальный пересчёт.
В распределённом источнике watermark обычно вычисляется с учётом нескольких партиций. Глобальный прогресс нельзя считать выше минимального прогресса ещё не опередившей партиции, иначе события из неё будут ошибочно признаны опоздавшими.
Watermark не гарантирует абсолютную полноту: это политика управления неопределённостью. Если источник может задерживать события дольше выбранного порога, нужны увеличенное допустимое опоздание, поздние исправления, периодический пересчёт или модель данных с поддержкой версий результата.
В примере окно можно закрыть при watermark 10:05, а событие с event_time 09:58 ещё должно изменить его результат. В реальной системе закрытие и исправление результата должны быть согласованы с потребителем: он должен уметь принимать обновления, версии или компенсирующие записи.
Сервис аналитики строит пятиминутный отчёт о платежах. Большинство событий приходит за секунды, но после восстановления мобильной сети часть событий задерживается на две-три минуты, а редкие сообщения могут опаздывать на десять минут.
Вариант с processing time дал минимальную задержку, но систематически занижал старые окна. Вариант с ожиданием десяти минут обеспечивал большую полноту, однако делал отчёты слишком медленными для оперативного контроля.
Выбрали watermark с допустимым опозданием в пять минут и отдельный канал исправлений для более поздних событий. Это дало приемлемую задержку для большинства отчётов, а редкие поздние платежи не терялись — они обновляли витрину при наличии идентификатора окна и версии результата.
Время последнего полученного события может быть высоким из-за одной быстрой партиции и ничего не говорить о задержавшихся данных в другой. Watermark — это консервативная оценка прогресса, обычно учитывающая все партиции или явно заданную стратегию неполноты. Если использовать максимум вместо безопасной нижней границы, окна будут закрываться преждевременно.
Его нельзя молча считать обычным событием: итог уже мог быть опубликован и использован downstream-системами. Возможные варианты — отбросить событие по бизнес-политике, отправить его в поток исключений, выполнить пересчёт затронутого периода или опубликовать корректирующую запись. Выбор зависит от цены ошибки и способности потребителей принимать изменения.
Увеличение периода ожидания повышает полноту, но задерживает все результаты, увеличивает объём состояния и стоимость хранения. Кроме того, оно не устраняет проблему неограниченных задержек и остановившейся партиции. На практике выбирают компромисс между свежестью и полнотой, а остаточную неопределённость покрывают исправлениями или периодическим пересчётом.