сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Kafka и стриминг

Окна и время события в стриминге

Время и окна в потоке

Агрегаты в стриминге считают по окнам, и всё упирается в то, какое время брать и как ждать опоздавших. Собес проверяет, различаешь ли ты event time от processing time и понимаешь ли роль watermark.

Стержень: считать надо по event time, а watermark решает, когда перестать ждать поздние события и закрыть окно.

// Формулировки: «event time против processing time?», «зачем watermark?», «какие бывают окна?».

Два времени и типы окон

Event time - когда событие случилось (записано в самом событии); processing time - когда его обработали. Агрегаты по processing time искажает любая задержка доставки: ретрай или всплеск нагрузки нарисуют ложный пик не в тот момент.

Окна бывают трёх видов. Tumbling - встык без нахлёста («по минутам»). Sliding/hopping - с нахлёстом («за 5 минут каждую минуту»). Session - по паузе неактивности (пользовательские сессии).

// Поэтому корректная аналитика почти всегда по event time: время берут из события, а не из часов обработчика, иначе «всплеск в 9:00» окажется артефактом ретрая доставки.

event / processing time
время события / время его обработки
tumbling / sliding / session
встык / с нахлёстом / по неактивности

Watermark и поздние события

События приходят не по порядку и с опозданием - мобильный клиент офлайн шлёт вчерашнее сегодня, и event-time обработка обязана это терпеть. Watermark - движущаяся оценка «события старше T мы больше не ждём»: он закрывает окна и триггерит выдачу результата. Это баланс полноты против латентности: ждать дольше - полнее, но позже.

Allowed lateness - грейс после watermark: поздние события ещё обновляют уже выданный результат; совсем поздние идут в side output или отбрасываются, но подсчитанно, не молча.

// Ловушка watermark - считать его по максимальному таймстемпу без допуска: один клиент с кривыми часами из будущего закроет все окна разом. Watermark ставят с запасом на разброс времени.

watermark
порог «старше - не ждём», закрывающий окна
allowed lateness
грейс на поздние события после watermark

Результат окна не финален

Session-окна закрываются только watermark'ом - пауза ведь может продлиться. Таймаут сессии плюс задержка данных дают реальную латентность результата, поэтому ждать «финальную» сессию в реальном времени нельзя.

Результат окна - не финальная истина, а «лучшее знание на сейчас»: даунстрим должен уметь принимать обновления (upsert по ключу окна) либо ждать финализации.

// Отсюда фундамент - идемпотентность по ключу окна (окно + сущность): повторная выдача окна перезаписывает, а не задваивает. Дописывать результаты окон append'ом - прямой путь к дублям строк.

late event
событие, пришедшее после закрытия своего окна
ключ окна
окно + сущность: единица идемпотентного upsert

Как отвечать: «Event time против processing time, и зачем watermark?»

Event time это время, когда событие реально произошло, оно записано в самом событии; processing time - когда обработчик его увидел. Считать агрегаты надо по event time, потому что доставка задерживается: мобильный клиент был офлайн и прислал вчерашние события сегодня, случился ретрай или всплеск. Если считать по processing time, я получу ложный пик не в тот момент - «всплеск в 9 утра», который на самом деле догон отставшей доставки. Но у event time есть проблема: раз события опаздывают, непонятно, когда окно можно закрывать и выдавать результат - вечно ждать нельзя. Ровно это решает watermark: движущаяся оценка «событий старше такого-то времени мы больше не ждём». Он закрывает окна и триггерит выдачу, а его величина - компромисс между полнотой и латентностью. Поздние сверх watermark ловлю через allowed lateness как обновление, а не теряю молча.

Почему это сильный ответ: разведены два времени с конкретным вредом processing time (ложный пик от ретрая), объяснена роль watermark как компромисса полнота/латентность и обработка поздних через allowed lateness.

На чём валят

  • Агрегировать по processing time и получить «всплеск в 9:00» из-за ретрая доставки.
  • Watermark по максимальному таймстемпу без допуска - клиент с кривыми часами закрыл все окна.
  • Игнорировать поздние события молча - минус процент выручки в отчёте без следов.
  • Ждать «финального» результата session-окон в реальном времени - они финализируются с лагом.
  • Дописывать результаты окон append'ом - обновления окна создают дубли строк.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 14, остальные разбираются в тренажёре.

  1. #windowing_time1 / 5
    Чем tumbling, sliding и session окна отличаются?
    A)Это три синонимичных названия одного и того же типа окна
    B)Различаются единицей времени: секунды, минуты и часы соответственно
    C)Tumbling — непересекающиеся; sliding — с нахлёстом; session — по паузам активности
    D)Окно tumbling считает агрегат по числу событий, а два других не считают агрегатов по своим непрерывным окнам во времени
    показать ответ и разбор
    +C)Tumbling — непересекающиеся; sliding — с нахлёстом; session — по паузам активности

    // разбор: Tumbling — фиксированные непересекающиеся окна (каждые 5 минут). Sliding — окна фиксированной длины с шагом меньше длины, поэтому они перекрываются (последние 5 минут, обновляя каждую минуту). Session — динамические окна, закрывающиеся после паузы бездействия, — естественны для пользовательских сессий переменной длины.

  2. #windowing_time2 / 5
    Чем event time отличается от processing time в потоковой обработке?
    A)Это два синонимичных названия для одного момента времени
    B)Event time — когда система обработала событие; processing time — когда оно случилось
    C)Оба совпадают до миллисекунды в реальной распределённой системе
    D)Event time — когда событие случилось; processing time — когда система его обработала
    показать ответ и разбор
    +D)Event time — когда событие случилось; processing time — когда система его обработала

    // разбор: Event time — момент, когда событие реально произошло (метка в самом событии); processing time — когда оно дошло до обработчика. Из-за сетевых задержек, буферов и офлайн-клиентов они расходятся. Корректная аналитика опирается на event time (иначе «клик в 23:59» упадёт не в тот час), но это требует механизма ожидания опоздавших — watermark.

  3. #windowing_time3 / 5
    Что такое watermark в потоковой обработке и зачем он нужен?
    A)Оценка «событий раньше этого времени уже не придёт» — когда можно закрывать окно
    B)Водяной знак, который физически наносится на каждое сообщение для защиты авторских прав
    C)Метрика загрузки процессора на узлах потокового кластера обработки данных
    D)Запрет на приём новых событий сразу после старта обработки
    показать ответ и разбор
    +A)Оценка «событий раньше этого времени уже не придёт» — когда можно закрывать окно

    // разбор: Watermark — движущаяся оценка «все события с event time до этого момента уже, скорее всего, пришли». Он позволяет закрывать окна по event time, балансируя полноту и задержку: ждать вечно нельзя, но и закрыть слишком рано — потерять опоздавших. События позже watermark (late data) обрабатываются отдельной политикой (allowed lateness, side output, пересчёт).

  4. #windowing_time4 / 5
    Почему поздние события (late data) ломают наивную оконную агрегацию?
    A)Поздних событий в реальности встречается редко — события приходят по порядку
    B)Событие за уже закрытое окно либо потеряется, либо потребует пересчёта агрегата
    C)Late data автоматически удаляет всю ранее посчитанную статистику этого окна по всем ранее закрытым во времени окнам агрегации
    D)Поздние события приходят раньше своего времени, а не позже него
    показать ответ и разбор
    +B)Событие за уже закрытое окно либо потеряется, либо потребует пересчёта агрегата

    // разбор: Клиент был офлайн, сеть тормозила, ретрай задержался — и событие с ранним event time приходит, когда его окно уже закрыто и агрегат посчитан. Наивная обработка его либо теряет (недосчёт), либо вынуждена пересчитывать окно. Поэтому нужна явная политика: watermark задаёт момент закрытия, allowed lateness — сколько ещё ждать опоздавших.

  5. #windowing_time5 / 5
    Чем скользящее (sliding) окно отличается от кувыркающегося (tumbling)?
    A)Tumbling-окна перекрываются, а sliding-окна идут строго встык без какого-либо перекрытия
    B)Это полные синонимы: оба вида окон нарезают поток одинаковым образом
    C)Tumbling — непересекающиеся окна встык; sliding — окна перекрываются, событие попадает в несколько
    D)Sliding-окна применимы к историческим батч-данным, но не к живому потоку
    показать ответ и разбор
    +C)Tumbling — непересекающиеся окна встык; sliding — окна перекрываются, событие попадает в несколько

    // разбор: Tumbling-окна фиксированной длины идут встык и не пересекаются: каждое событие ровно в одном окне (почасовые суммы). Sliding-окна той же длины сдвигаются с шагом меньше длины и перекрываются, поэтому событие попадает в несколько окон — так считают скользящие метрики («сумма за последний час, обновляемая каждую минуту»). Session-окна — третий вид, по паузам активности. Выбор диктует нужная семантика метрики.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.