Обрабатывать поток кажется просто — читай сообщения и считай, — пока не появляются два вопроса, на которых ломаются самодельные решения: по какому времени считать (когда событие произошло или когда до нас дошло?) и когда закрывать окно, если данные приходят с задержкой и не по порядку. Apache Flink построен вокруг честных ответов на них: event-time как основа и watermarks как способ рассуждать о полноте данных во времени. Обзорная статья Kafka Streams против Flink сравнивает движки и помогает выбрать; эта серия — про то, как Flink устроен внутри, и начинается она с модели времени и окон.
Первая статья серии «Apache Flink: глубокое погружение».
В статье
- Архитектура: JobManager, TaskManager, слоты
- Dataflow и ключи: как поток течёт по графу
- Модель времени: event / processing / ingestion
- Watermarks: «всё до момента T уже видели»
- Окна: tumbling, sliding, session
- Опоздавшие данные: допуск и отброс
- Эксплуатационные ошибки
- Checklist: готовность к event-time
- Демо и версии
- Документация и первоисточники
Архитектура: JobManager, TaskManager, слоты
Flink — не библиотека, которую вызывают в цикле, а распределённый рантайм. Задание (job) описывается как dataflow-граф операторов: источники, трансформации, стоки. Этот граф Flink планирует и разносит по кластеру, а не исполняет как последовательную программу.
Кластер держится на двух ролях:
| Роль | Отвечает за |
|---|---|
| JobManager | координация: планирует граф, раздаёт задачи, управляет чекпоинтами и восстановлением после сбоя |
| TaskManager | выполнение: гоняет операторы на данных, хранит состояние ключей локально |
TaskManager делит ресурсы на слоты (slots) — единицы параллелизма. Слот — это доля памяти и CPU воркера, в которую помещается подзадача (или сцепленная цепочка операторов). Параллелизм оператора — это число его параллельных экземпляров, и каждый занимает слот. Отсюда простое следствие: суммарного числа слотов в кластере должно хватать на суммарный параллелизм графа, иначе задание не запланируется.
В демо-стенде это видно буквально: JobManager поднимает Flink Dashboard на localhost:8081, где отображаются слоты и метрики, а TaskManager исполняет операторы и пишет результат в лог. Одно INSERT ... SELECT в SQL-клиенте — это уже полноценный dataflow-граф: источник datagen, оператор оконной агрегации, сток print.
Dataflow и ключи: как поток течёт по графу
Логически задание — цепочка: источник → трансформации → сток. Физически каждый оператор разбивается на параллельные экземпляры, а поток между операторами перераспределяется. Ключевой вопрос — по какому признаку событие попадает в тот или иной параллельный экземпляр, потому что от этого зависит, где живёт состояние.
Здесь работает keyed stream: поток разбивается по ключу (в SQL — по столбцам GROUP BY, в DataStream API — по keyBy). Все события с одним ключом гарантированно попадают в один и тот же экземпляр оператора. Это даёт два свойства сразу:
- распределение — разные ключи расходятся по разным экземплярам, нагрузка масштабируется;
- локальность состояния — состояние (счётчик, агрегат, содержимое окна) для ключа хранится там же, где обрабатывается ключ; оператору не нужно ходить за ним по сети.
Именно поэтому в оконной агрегации ниже группировка идёт по window_start, window_end: окно — это, по сути, ключ, и его состояние (частичный счётчик) копится локально до момента закрытия. Как это состояние переживает сбои и почему оно не задваивается — тема отдельной статьи серии про состояние и exactly-once; здесь достаточно, что состояние привязано к ключу и живёт на TaskManager.
Модель времени: event / processing / ingestion
У каждого события в потоке есть несколько «времён», и путать их — источник самых коварных багов аналитики.
| Время | Что означает | Проблема |
|---|---|---|
| Event-time | когда событие произошло в реальности (штамп внутри самого события) | приходит не по порядку и с задержкой |
| Ingestion-time | когда событие попало во Flink | зависит от лага доставки |
| Processing-time | когда конкретный оператор его обработал | зависит от нагрузки и перезапусков |
Разница не академическая. Считать «клики за 10:00–10:01» по processing-time значит считать не по времени клика, а по времени, когда сообщение доехало и Flink до него добрался. Всплеск лага, перезапуск оператора, изменение параллелизма — и одно и то же событие в разных прогонах попадёт в разные минуты. Результат перестаёт быть свойством данных и становится свойством инфраструктуры.
Event-time честнее ровно потому, что он воспроизводим: событие с меткой 10:00:37 попадёт в минуту 10:00 всегда — сегодня в реальном времени, завтра при переигрывании архива, независимо от лага и перезапусков. Это и есть каноническое требование Тайлера Акидау из Streaming 101: для аналитики потока считать надо по времени события, а не по времени обработки.
Плата за честность — событий по event-time нельзя дождаться «по часам»: они опаздывают и приходят вперемешку. Нужен механизм, который скажет оператору «по event-time мы уже дошли до момента T, окна раньше T можно закрывать». Это и есть watermark.
Watermarks: «всё до момента T уже видели»
Watermark(T) — это метка в потоке, означающая «событий с event-time меньше T дальше почти не будет». Оператор, увидев watermark, вправе считать, что окна, чей конец не позже T, собрали все свои данные, и закрыть их.
Watermark порождается из данных: обычно по максимальному наблюдённому event-time минус допуск на неупорядоченность (bounded out-of-orderness). Далее он распространяется по графу вместе с потоком: оператор с несколькими входами берёт минимум входящих watermark (медленный вход тормозит продвижение — иначе можно закрыть окно, не дождавшись отстающей партиции).
В стенде watermark задаётся прямо в определении источника:
CREATE TEMPORARY TABLE clicks (
user_id INT,
url STRING,
event_time AS CAST(CURRENT_TIMESTAMP - INTERVAL '0.001' SECOND * CAST(RAND() * 10000 AS INT) AS TIMESTAMP(3)),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'datagen', 'rows-per-second' = '20',
'fields.user_id.min' = '1', 'fields.user_id.max' = '5', 'fields.url.length' = '8'
);
Разберём две ключевые строки:
event_time AS CAST(CURRENT_TIMESTAMP - INTERVAL '0.001' SECOND * CAST(RAND() * 10000 AS INT) AS TIMESTAMP(3))— генератор искусственно раскидывает event_time на 0..10 секунд в прошлое (разброс задаётся в миллисекундах). Это создаёт out-of-order: событие, «случившееся» 8 секунд назад, приходит позже события, случившегося 2 секунды назад. Множитель интервала обязан быть целым числом: соблазнительноеINTERVAL '10' SECOND * RAND()проходит валидацию и даже создаёт таблицу, а падает только наINSERT, когда доходит до кодогенерации —Unsupported casting from DOUBLE to INTERVAL.WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND— watermark отстаёт от наблюдённого времени на 5 секунд. Это и есть допуск на неупорядоченность: Flink готов ждать опоздавшие события до 5 секунд, прежде чем закрыть окно.
Соотношение этих двух чисел — разброс 0..10 секунд против допуска 5 секунд — и есть суть демонстрации. Часть событий укладывается в допуск и попадает в свои окна; часть опаздывает сильнее допуска — о них ниже.
Окна: tumbling, sliding, session
Окно назначает событию интервал по его event-time и агрегирует всё, что в интервал попало. Flink SQL выражает окна через табличные функции (TVF). Основной стенд считает tumbling — непересекающиеся окна фиксированной длины:
CREATE TEMPORARY TABLE out_counts (window_start TIMESTAMP(3), window_end TIMESTAMP(3), cnt BIGINT)
WITH ('connector' = 'print');
INSERT INTO out_counts
SELECT window_start, window_end, COUNT(*) AS cnt
FROM TABLE(TUMBLE(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '10' SECOND))
GROUP BY window_start, window_end;
TUMBLE(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '10' SECOND) нарезает поток на 10-секундные окна по столбцу event_time. Результат каждого закрытого окна печатается стоком print в лог TaskManager строкой вида:
+I[2026-08-19T15:00:10, 2026-08-19T15:00:20, 187]
Здесь +I — вставка (insert) новой строки результата, дальше границы окна и число событий. Главное — окно печатается по watermark, а не по стенным часам. Flink не смотрит на реальное «сейчас»; он ждёт, пока watermark пройдёт правую границу окна (то есть event-time потока дойдёт до конца окна плюс допуск), и только тогда фиксирует счётчик. Поэтому результат зависит от данных, а не от того, когда оператор до них добрался.
Кроме tumbling стенд позволяет попробовать ещё два вида окон — достаточно заменить TVF в INSERT:
Sliding (HOP) — окна перекрываются:
HOP(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '5' SECOND, INTERVAL '10' SECOND)
Окна длиной 10 секунд с шагом 5 секунд. Соседние окна перекрываются наполовину, поэтому каждое событие попадает сразу в два окна. Sliding берут, когда нужна скользящая метрика («за последние 10 секунд, обновляемая каждые 5»).
Session — окна по активности:
SESSION(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '...' SECOND)
Session-окно не имеет фиксированной длины: оно растёт, пока события идут, и закрывается после паузы без событий заданной длины (gap). Такое окно естественно ложится на пользовательские сессии: активность склеивается в один интервал, пауза разрывает его.
Выбор окна — это выбор вопроса, на который вы отвечаете: «сколько за фиксированный интервал» (tumbling), «сколько за скользящий интервал» (sliding), «сколько за сеанс активности» (session).
Опоздавшие данные: допуск и отброс
Watermark с допуском в 5 секунд — это обещание закрыть окно, дождавшись опоздавших не более чем на 5 секунд. Но генератор раскидывает event_time на 0..10 секунд. Значит, часть событий опаздывает сильнее допуска — они «слишком поздние» и в окно не попадают: к моменту их прихода watermark уже прошёл конец их окна, окно закрыто и посчитано.
Это ровно тот компромисс, вокруг которого крутится вся модель:
- больше допуск → полнее результат, но выше задержка. Ждём дольше, ловим больше опоздавших, но и окно закрывается позже.
- меньше допуск → ниже задержка, но теряем опоздавших. Окна закрываются быстрее, зато «хвост» опоздавших событий отбрасывается.
Проверяется это прямо на стенде: если уменьшить допуск с 5 до 1 секунды (event_time - INTERVAL '1' SECOND), счётчики окон падают — часть событий, укладывавшихся в 5-секундный допуск, теперь опаздывает сильнее секунды и отбрасывается. Данные буквально исчезают из агрегата, хотя источник даёт тот же поток.
Числа со стенда: генератор даёт 20 событий в секунду, то есть в среднем на 10-секундное окно приходится 200 событий. С допуском 5 секунд окна набирают в среднем около 180 (по отдельным окнам разброс примерно 150–205), с допуском 1 секунда — около 130. Смотреть здесь нужно именно на среднее: разброс event_time перемешивает события между соседними окнами, поэтому отдельно взятое окно может набрать и больше 200. А вот систематическая недостача при допуске 5 секунд — уже не шум: это хвост самых опоздавших событий, не успевших к закрытию своего окна. Потеря идёт тихо, без единой ошибки в логе.
Это наглядный урок: допуск на неупорядоченность — не косметический параметр, он напрямую управляет полнотой результата.
Контрольный эксперимент — перевести то же окно на processing-time (TUMBLE по PROCTIME() вместо event_time). Тогда результат перестаёт зависеть от разброса задержек: ничего «не опаздывает», потому что окно считает по времени обработки. Но честность теряется — одно и то же событие в разных прогонах попадёт в разные окна, потому что зависит от того, когда оператор до него добрался. Это та же дилемма полнота ↔ задержка, что и в хвостовых задержкахСкоро: нельзя одновременно и мгновенно, и полно — приходится осознанно выбирать точку компромисса.
Эксплуатационные ошибки
- Считать по processing-time там, где нужна аналитика. Результат становится свойством инфраструктуры, а не данных: перезапуск, лаг, изменение параллелизма — и цифры «поплыли». Processing-time уместен для грубого мониторинга «здесь и сейчас», но не для отчётов, которые должны воспроизводиться.
- Допуск watermark выставлен наугад. Слишком маленький — тихо отбрасывает опоздавших (счётчики занижены, и никто не замечает); слишком большой — окна закрываются с большой задержкой. Допуск подбирают под реальный разброс event-time в потоке, а не «на глаз».
- Watermark не продвигается на простое источника. Если партиция замолчала, её watermark стоит, а оператор берёт минимум по входам — общий watermark замирает, и окна не закрываются вовсе. Idle-источники нужно учитывать явно, иначе результат «зависает».
- Забыли, что окно закрывается по watermark, а не по часам. На тихом потоке окна не печатаются не потому, что «сломалось», а потому что event-time не дошёл до их границы. Диагностика начинается с вопроса «докуда доехал watermark», а не «сколько сейчас времени».
- Медленный вход тормозит весь граф. Оператор с несколькими входами продвигает watermark по минимуму — одна отстающая партиция держит окна всего задания закрытыми. Следят за перекосом (skew) входов.
- Sliding-окно недооценили по стоимости.
HOPс мелким шагом кладёт каждое событие сразу в несколько окон — состояния и вычислений кратно больше, чем у tumbling. Шаг выбирают исходя из реальной потребности, а не «почаще на всякий случай».
Checklist: готовность к event-time
- Выбрано ли время осознанно — event-time для аналитики (воспроизводимость) против processing-time для грубого мониторинга?
- Есть ли в событиях надёжный event-time штамп и объявлен ли
WATERMARK FORна источнике? - Подобран ли допуск на неупорядоченность под реальный разброс event-time (а не наугад)?
- Понятен ли компромисс: насколько увеличение допуска бьёт по задержке, а уменьшение — по полноте?
- Учтены ли idle-источники, чтобы watermark не замирал на молчащей партиции?
- Выбран ли тип окна под вопрос (tumbling / sliding / session) и оценена ли стоимость sliding с мелким шагом?
- Есть ли наблюдаемость продвижения watermark и отбрасываемых «слишком поздних» событий (иначе потери данных незаметны)?
- Хватает ли слотов в кластере на суммарный параллелизм графа?
Если на большинство — «да», вы считаете поток по времени данных, а не по капризам инфраструктуры.
Демо и версии
- Demo:
digital-cookbook/flink/00-model/— самодостаточный стенд без Kafka: источник — встроенный коннекторdatagen, сток —print. Tumbling-окно по event-time с out-of-order-данными и watermark; в комментариях — как переключить на sliding (HOP), session (SESSION) и processing-time (PROCTIME()), и как уменьшение допуска роняет счётчики окон. Запуск черезdocker compose up -d, работа — в Flink SQL client, результат окон — в логе TaskManager (+I[start, end, cnt]), Dashboard наlocalhost:8081. - Версии: Apache Flink 1.20.1 (Java 17), Flink SQL. Поведение TVF-окон, watermark и допуска зафиксировано на этой версии — при обновлении пиновать и перепроверять, семантика оконных функций чувствительна к версии.
Документация и первоисточники
- Канон стриминга: Streaming 101 и Streaming 102 Тайлера Акидау — event-time, watermarks, окна и компромисс полнота↔задержка «от первоисточника».
- Документация: Apache Flink 1.20 — docs (концепции времени, генерация watermark, оконные TVF в Flink SQL).
- Смежное на сайте: Kafka Streams против Flink — что выбрать, состояние и exactly-once во Flink, Flink SQL: коннекторы и эксплуатация, проектирование событийных пайплайнов, хвостовые задержки и перцентилиСкоро.
Комментарии