Watermarks
Watermark в потоковой обработке данных (stream processing) — это монотонно возрастающая нижняя оценка времён событий, которые ещё могут поступить в поток. Когда watermark достигает значения X, система объявляет, что все события с временем меньше X с высокой вероятностью уже получены.
YDB реализует механизм watermarks в потоковых запросах: они используются для корректного закрытия временных окон агрегации (HoppingWindow) и гарантируют, что результат окна выдаётся только тогда, когда система убедилась в полноте входных данных за этот период.
В потоковой обработке каждое событие имеет две временны́е метки: время события (event time) — момент, когда событие произошло в реальном мире, и время обработки (processing time) — момент, когда система получила событие. Из-за сетевых задержек, сбоев и неравномерной нагрузки эти два значения могут существенно расходиться. Именно поэтому системе нужен механизм watermarks: без него она не знает, когда можно считать прошедший временной диапазон достаточно полным, чтобы выдать результат.
Компромисс Точности и Задержки
В YDB этот компромисс регулируется явно: параметр WATERMARK в предложении GROUP BY задаёт выражение для вычисления watermark, в том числе величину отставания от времени последнего события. Это позволяет подобрать баланс между актуальностью результатов и полнотой учёта запоздавших событий.
Watermark не может одновременно учитывать сколь угодно большие задержки событий и продвигаться быстро: чем дольше система ждёт опоздавших событий, тем позже она выдаёт результаты. Это фундаментальный компромисс потоковой обработки.
События, поступившие после того, как watermark прошёл соответствующий временной диапазон, считаются опоздавшими и отбрасываются. Подробнее об обработке опоздавших событий и о настройке watermarks: Watermarks.