Watermarks
Watermark — монотонно возрастающая нижняя оценка времён событий в потоке (подробнее о концепции: Watermarks). В данном разделе описана настройка watermarks в потоковых запросах YDB.
Время события
В потоковой обработке каждое событие имеет временную метку, по которой система отслеживает прогресс времени в потоке. В текущей реализации источником времени события может быть только время записи события в топик, доступное через системную колонку __ydb_write_time.
Примечание
Поддержка произвольных выражений для извлечения времени из данных события (например, поля event.created_at) планируется в следующих версиях.
Использование
Watermark используется операциями, зависящими от прогресса времени события в потоке. В YDB к таким операциям относится оконная агрегация HoppingWindow — она определяет скользящие временные окна, по которым группируются события. При получении watermark HoppingWindow закрывает все окна, которые полностью покрыты этим значением.
Вычисление watermark
Когда система получает событие, она обновляет watermark — продвигает его вперёд по временной оси. Watermark вычисляется как максимальное наблюдаемое время события − delay, где delay — величина отставания, заданная в выражении WATERMARK (например, Interval("PT5S") в WATERMARK = __ydb_write_time - Interval("PT5S")).
События в потоке могут приходить не в хронологическом порядке: событие с временем 10:00:03 может быть обработано после события с временем 10:00:05. Причины: расхождение часов в распределённой системе, сетевые задержки, неравномерная нагрузка на партиции топика.
Параметр delay задаёт допустимый «запас» времени для событий, поступающих с задержкой. Например, при delay в 5 секунд событие с временем 00:00:48 будет принято, даже если уже пришли события с временем 00:00:50: watermark ещё не дошёл до 00:00:48. Если то же событие придёт позже, когда watermark уже продвинулся за 00:00:48, оно будет признано опоздавшим и отброшено.
О компромиссе между точностью и задержкой выдачи результатов: Компромисс точности и задержки.
Простаивающие партиции
Если входной топик содержит несколько партиций, каждая из них продвигает watermark независимо. Общий watermark запроса не обгоняет самую медленную партицию: окна не закрываются, пока хотя бы одна партиция не достигла соответствующего момента времени.
Если одна из партиций перестаёт получать данные, её watermark перестаёт продвигаться вперёд. Такая партиция называется простаивающей (idle). Пока простаивающая партиция учитывается при вычислении общего watermark, он тоже перестаёт продвигаться, и результаты не выдаются, несмотря на поступление данных от других партиций.
Чтобы избежать этой блокировки, простаивающая партиция исключается из вычисления общего watermark по истечении настраиваемого периода ожидания (параметр WATERMARK_IDLE_TIMEOUT, подробнее в разделе Настройка).
Настройка
Watermarks включаются и настраиваются в секции WITH при чтении из топика.
Параметры настройки:
WATERMARK— выражение для вычисления watermark. Сейчас поддерживается только время записи в топик с константной задержкой. Формат:__ydb_write_time - Interval("<delay>"), где<delay>задаётся в формате ISO 8601.WATERMARK_GRANULARITY— периодичность генерации watermarks. Чем она меньше, тем больше потребление CPU запросом и тем меньше задержка ответа. Имеет смысл только для потоковых запросов. Задаётся в формате ISO 8601. Значение по умолчанию — 1 секунда.WATERMARK_IDLE_TIMEOUT— период, после которого простаивающая партиция будет исключена из вычисления объединённого watermark. Имеет смысл только для потоковых запросов. Задаётся в формате ISO 8601. Значение по умолчанию — 5 секунд.
Важно
При использовании HoppingWindow первый параметр (time extractor) и источник времени в выражении WATERMARK должны совпадать. В текущей реализации оба должны использовать __ydb_write_time.
Пример
Ниже приведён пример потокового запроса с watermark и оконной агрегацией. Запрос читает события из топика, фильтрует их по полю pass и агрегирует значения payload в окнах по 10 секунд с шагом 5 секунд. Watermark настроен с отставанием в 5 секунд.
Входные данные
{"pass": 1, "payload": "a"} // время записи: 1970-01-01T00:00:40Z
{"pass": 1, "payload": "b"} // время записи: 1970-01-01T00:00:42Z
{"pass": 0, "payload": "c"} // время записи: 1970-01-01T00:00:50Z
{"pass": 1, "payload": "d"} // время записи: 1970-01-01T00:00:40Z
Запрос
CREATE STREAMING QUERY example AS
DO BEGIN
$input = (
SELECT
t.*,
__ydb_write_time AS ts
FROM
Input
WITH (
FORMAT = json_each_row,
SCHEMA = (
pass Int64,
payload String
),
WATERMARK = __ydb_write_time - Interval("PT5S")
) AS t
);
$output = (
SELECT
AGGREGATE_LIST(payload) AS result,
CAST(HOP_END() AS String) AS ts
FROM
$input
WHERE pass > 0
GROUP BY
HoppingWindow(ts, "PT5S", "PT10S")
);
INSERT INTO Output
SELECT
ToBytes(Unwrap(Yson2::SerializeJson(Yson::From(TableRow()))))
FROM $output;
END DO;
Где:
CREATE STREAMING QUERY— создаёт именованный потоковый запрос.__ydb_write_time— системная колонка, содержащая время записи события в топик.FORMAT = json_each_row— формат данных в топике, каждая строка содержит отдельный JSON-объект.WATERMARK = __ydb_write_time - Interval("PT5S")— watermark с отставанием 5 секунд.Interval("PT5S")задаёт интервал в формате ISO 8601.AGGREGATE_LIST— агрегатная функция, собирающая значения в список.HOP_END()— возвращает временную метку конца текущего окна.HoppingWindow(ts, "PT5S", "PT10S")— оконная функция с шагом 5 секунд и размером окна 10 секунд.
Результат
{"result": ["a", "b"], "ts": "1970-01-01T00:00:45.000000Z"}
Пояснение
- Первое событие (
"a", время записи 40с) проходит фильтр (pass > 0) и попадает в окна[35; 45)и[40; 50). Watermark продвигается до 35с и не закрывает ни одного окна. - Второе событие (
"b", время записи 42с) аналогично попадает в окна[35; 45)и[40; 50). Watermark продвигается до 37с. - Третье событие (
"c", время записи 50с) отбрасывается на фильтре (pass = 0). Несмотря на это, событие всё равно продвигает watermark до 45с. Watermark закрывает окно[35; 45)— результат["a", "b"]выдаётся. - Четвёртое событие (
"d", время записи 40с) не продвигает watermark: он уже находится на отметке 45с. Событие проходит фильтр, но отбрасывается как опоздавшее — его время записи (40с) меньше текущего watermark (45с). Хотя окно[40; 50)ещё открыто, watermark уже обещал, что событий с временем < 45с больше не будет, поэтомуdне учитывается ни в одном из своих окон.
См. также
- GROUP BY ... HoppingWindow — оконная функция, использующая watermarks.
- WITH — секция WITH для настройки параметров чтения из топика.
- Гарантии доставки данных — гарантии доставки данных.
- Чекпоинты — механизм чекпоинтов.