Представьте: вы разрабатываете антифрод-систему для банка. Пользователь совершает пять покупок подряд за пару минут. Все транзакции успешно попадают в Kafka, никаких ошибок, consumer lag почти нулевой, Flink работает стабильно.
Но модель выдачи решений получает на вход признак: transactions_in_10_minute_window = 3 Хотя по факту их было пять, а вместо вероятности мошенничества 0.91 она выдаёт 0.42.
Естественно, мы начинаем копать модель: может, переобучилась? Но модель честно отработала - получила три транзакции и посчитала вероятность исходя из них.
Проблема в другом: пайплайн решил, что к окну относятся только три события.
Раньше, в батче, такой проблемы обычно не возникает: данные уже лежат в хранилище, их можно отсортировать по времени и спокойно пересчитать любой признак.
Однако, в стриминге всё сложнее: событие может произойти в 12:00:03, а прийти в обработку в 12:00:08, следующее — случиться в 12:00:04, но оказаться в системе раньше предыдущего.
Out-of-order — не баг, а фича
Давайте рассмотрим пять транзакций одного пользователя. В таблице ниже — их время возникновения (Event Time) и реальное время поступления в Flink (Processing Time).
№ | Event Time | Processing Time |
|---|---|---|
1 | 12:00:01 | 12:00:01 |
2 | 12:00:02 | 12:00:02 |
3 | 12:00:03 | 12:00:08 |
4 | 12:00:04 | 12:00:09 |
5 | 12:00:05 | 12:00:05 |
По бизнесу порядок: 1->2->3->4->5, а в потоке они прилетели: 1->2->5->3->4 и это не сбой — причины разные: партиции, сеть, продюсеры перегружены, узлы тормозят. У каждого события минимум два времени: когда случилось (Event Time) и когда Flink его увидел (Processing Time). Разница между ними может составлять секунды, минуты, а иногда и часы.
Теперь представьте: признак «число транзакций за последние 10 минут», если считать по Event Time, транзакция в 12:09:59 попадает в окно 12:00–12:10, а если по Processing Time, она приходит в 12:10:07, и уже относится к следующему окну. На мелких задержках это незаметно, а на границе окна — разница в разы.
Чтобы считать по Event Time, timestamp должен приезжать вместе с событием. Например, в объекте Transaction мы храним поле eventTime:
public record Transaction( String userId, long amount, Instant eventTime ) {}
Дальше сообщаем Flink, какой timestamp использовать, и задаём стратегию формирования водяных знаков:
WatermarkStrategy<Transaction> strategy = WatermarkStrategy .<Transaction>forBoundedOutOfOrderness( Duration.ofSeconds(5) ) .withTimestampAssigner( (event, timestamp) -> event.eventTime().toEpochMilli() ); DataStream<Transaction> timedStream = transactions.assignTimestampsAndWatermarks(strategy);
Здесь две разные сущности:
TimestampAssigner — говорит Flink: «это timestamp события, используй его как Event Time»;
WatermarkStrategy — определяет, как Flink должен учитывать out-of-order события и продвигать Event Time.
Watermark — это метка прогресса Event Time, Flink получил события с временами 12:00:01, 12:00:02, 12:00:05. Максимальное время — 12:00:05, но это не значит, что событие на 12:00:03 не может прийти следующим. Поэтому мы делаем запас:
forBoundedOutOfOrderness(Duration.ofSeconds(5)).
Упрощённо можно думать, что водяной знак примерно равен максимальному EventTime – задержка, но на практике он генерируется по таймеру (по умолчанию каждые 200 мс) и не продвигается, если поток событий затихает — даже если максимальный EventTime большой. А если одна из партиций Kafka простаивает, то водяной знак вообще может остановиться, если не настроена стратегия withIdleness(). Поэтому реальное поведение зависит не только от задержки, но и от активности источников, а также от того, как обрабатываются несколько входных потоков.
Сколько секунд задержки ставить?
В примере выше мы поставили 5 секунд практически «с потолка», однако в реальном проекте я бы сначала посмотрел на распределение задержек в источнике.
Например:
Перцентиль | Задержка |
|---|---|
p50 | 200 ms |
p95 | 2.5 s |
p99 | 8 s |
Если выбрать задержку в 5 секунд, часть событий с большим нарушением порядка будет приходить уже после продвижения водяных знаков и станут опаздавшими. Насколько большой будет эта доля, нужно измерять на реальных данных.
Получаем:
Меньше задержка | Больше задержка |
|---|---|
Меньше latency | Больше полнота |
Быстрее закрываются окна | Больше времени на приход событий |
Выше риск late events | Ниже риск late events |
Потенциально меньше state | Потенциально больше state |
Нужно выводить значение из наблюдаемого распределения задержек, а не копировать число из примера в документации и обязательно учитывать бизнес-требования, так для антифрода дополнительные 5 секунд ожидания могут стоить дороже, чем несколько поздних событий. Для аналитического pipeline ситуация может быть совершенно обратной.
Как окна работают с водяными знаками
Допустим, нам нужны десятиминутные окна по Event Time:
12:00 ───────── 12:10 ← окно 1 12:10 ───────── 12:20 ← окно 2
Окно не закрывается просто потому, что часы сервера показали 12:10. Для Event Time оно закрывается, когда Watermark пересекает границу окна и здесь появляется вся цепочка:
Event Time → Watermark → Window → Aggregation → Feature
Поздние события: allowedLateness и side outputs
Допустим, что окно 12:00–12:10 уже закрылось, водяной знак пересёк 12:10, и вдруг приходит событие: eventTime = 12:09:59. По умолчанию, после закрытия окна, такое событие уже не изменит его результат, но иногда терять такие данные слишком дорого — особенно в антифроде. Можно разрешить окну принимать поздние события ещё некоторое время с помощью allowedLateness:
.window(TumblingEventTimeWindows.of(Time.minutes(10))) .allowedLateness(Time.seconds(10))
Теперь окно может принимать опоздания ещё 10 секунд после закрытия. Результат окна может измениться постфактум. Например, сначала transactions = 3, а через 7 секунд — уже 4.
И тут есть подводный камень: Flink может сгенерировать несколько выходных записей для одного окна — первую при закрытии, потом ещё при каждом позднем событии, если агрегат поменялся. Если наш сервис не умеет обновлять результат, то повторная отправка может всё сломать.
Для совсем уж опоздавших событий, которые не влезают даже в allowedLateness, можно сделать side output:
OutputTag<Transaction> lateTag = new OutputTag<>("late-transactions") {}; // в оконной операции: .sideOutputLateData(lateTag) // затем: DataStream<Transaction> lateStream = features.getSideOutput(lateTag);
Такой поток можно использовать для мониторинга, аудита, отдельного пересчёта или просто чтобы ругаться на источник.
Важно: sideOutputLateData относится именно к событиям, которые пришли уже после того, как для окна истёк допустимый период позднего прихода.
Stateful processing: почему память — не просто кэш
Для многих признаков недостаточно посмотреть только на текущее событие.
Например:
средняя сумма транзакций за день;
количество транзакций за последние 24 часа;
максимальная сумма покупки;
количество уникальных устройств.
Это называется stateful processing. Во Flink можно использовать разные state backend, например:
HashMapStateBackend;EmbeddedRocksDBStateBackend.
Выбор зависит от объёма состояния и требований к latency и производительности. При большом объёме state RocksDB помогает, но требует серьёзной настройки: размер блоков, кеши, политика компакции, количество параллельных операций. Для маленького состояния RocksDB может оказаться медленнее, чем HashMap, из‑за сериализации/десериализации и дисковых операций. Выбор бэкенда всегда должен сопровождаться нагрузочным тестированием на вашем профиле данных.
Самая опасная ловушка: Training-Serving Skew
Допустим, мы обучили модель на исторических данных, все события уже лежат в хранилище, мы можем упорядочить их по Event Time и посчитать: transactions_in_10_minute_window = 5
Модель показывает отличные метрики. В продакшен streaming pipeline использует другую временную семантику, другую политику обработки поздних событий или вообще Processing Time.
Окружение | Значение признака |
|---|---|
Обучение (batch) | 5 |
Продакшен (streaming) | 3 |
Модель одна, пользователь один, транзакции те же — а вход разный. Это и есть training-serving skew. Именно поэтому недостаточно просто договориться об имени поля в Feature Store, нужно зафиксировать семантику признака, например, количество транзакций пользователя с Event Time в заданном десятиминутном интервале, с определённой политикой обработки поздних событий.
А тренировочный и продакшен-пайплайны должны реализовывать эту семантику согласованно, иначе можно месяцами улучшать модель, а проблема окажется вообще не в модели.
Однако бывает, что инфраструктурные метрики идеальны: CPU 30%, Memory 40%, Kafka lag ≈ 0. А водяной знак не двигается.
Для streaming-системы это серьёзный сигнал. Время застыло, окна могут не закрываются, а признаки — не обновляются.
Причина может быть в:
отстающей partition;
остановившемся продюсере;
проблемах upstream;
особенностях Watermark Strategy;
задержке на одном из участков pipeline.
Поэтому нужно смотреть не только на инфраструктурные метрики, но и на прогресс Event Time:
currentInputWatermark и currentOutputWatermark позволяют видеть прогресс watermark на входе и выходе оператора
watermarkLag — производная метрика, позволяющая оценивать, насколько Event Time отстаёт от текущего времени или другого выбранного эталона..
Система при этом может выглядеть абсолютно здоровой с точки зрения CPU и памяти. Мониторинг Event Time —продакшен-инструмент, как например мониторинг CPU, memory или Kafka lag.
Что происходит после падения Flink
Stateful пайплайн сохраняет состояние через чекпоинты. После сбоя он восстанавливается из последнего чекпоинта и продолжает. Но для ML есть нюанс: вы уже отправили предсказание во внешний сервис, а потом Flink упал. При восстановлении он может обработать те же события и снова отправить результат.
И если операция неидемпотентна, то можем получить:
повторную обработку;
дублирование результата;
повторную отправку уведомления;
рассинхронизацию с внешним Feature Store;
в худшем случае — повторное бизнес-действие.
Чекпоинт восстанавливает состояние Flink, но не гарантирует, что внешний сервис не получит один и тот же результат повторно. Для некоторых источников вывода можно использовать exactly-once, например транзакционный Kafka Sink. В этом случае Flink помогает гарантировать обработку результата ровно один раз в рамках поддерживаемой семантики sink. Но такая схема сложнее и может увеличить latency. Поэтому в реальном проекте приходится выбирать: делать внешнюю операцию идемпотентной или использовать механизм exactly-once там, где он действительно нуже
Иногда говорят: «А почему бы не добавить сюда LLM или векторную БД?» Но это не лучшее решение, так как подсчёт транзакций за 10 минут — это классическая агрегация, которую Flink делает лучше всех. LLM и векторные базы хороши для текста и поиска, но они не заменяют оконные функции. Не надо усложнять там, где достаточно счётчика.
Ниже привожу базовый пример, со всеми элементами в одном месте:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); KafkaSource<Transaction> source = KafkaSource.<Transaction>builder() .setBootstrapServers("localhost:9092") .setTopics("transactions") .setGroupId("fraud-detector") .setValueOnlyDeserializer(new TransactionDeserializer()) .build(); WatermarkStrategy<Transaction> wmStrategy = WatermarkStrategy .<Transaction>forBoundedOutOfOrderness( Duration.ofSeconds(5) ) .withTimestampAssigner( (tx, ts) -> tx.eventTime().toEpochMilli() ) .withIdleness(Duration.ofSeconds(10)); // считаем партицию простаивающей после 10 секунд без событий DataStream<Transaction> transactions = env.fromSource( source, wmStrategy, "источник-транзакций" ); OutputTag<Transaction> lateTag = new OutputTag<>("опоздавшие-транзакции") {}; SingleOutputStreamOperator<Feature> features = transactions .keyBy(Transaction::getUserId) .window( TumblingEventTimeWindows.of( Time.minutes(10) ) ) .allowedLateness(Time.seconds(10)) .sideOutputLateData(lateTag) .aggregate(new FeatureAggregator()); DataStream<Transaction> lateStream = features.getSideOutput(lateTag); DataStream<Prediction> predictions = features.map(new MLScorer()); predictions.addSink(new MySink()); env.execute("Антифрод: обработка транзакций");
Здесь можно увидеть всю цепочку: Event Time → Watermark → Window → Allowed Lateness → Aggregation → Feature → ML
Flink часто хвалят за throughput, масштабируемость и fault tolerance, но для Real-Time ML есть ещё одна вещь, которая порой важнее: Flink позволяет явно определить, что мы считаем «произошедшим сейчас». Событие может случиться в 12:00:03, а прийти в систему в 12:00:08, если pipeline этого не учитывает, модель получит не тот признак ещё до вызова predict() и тогда проблема выглядит как ML-проблема, хотя корень находится в streaming-слое.

