Три предыдущие статьи серии разбирали Flink через SQL — и это честно: SQL/Table API и DataStream — это два API над общим runtime (те же event-time, watermark’и, управляемое состояние, чекпоинты). Для оконных агрегатов и джойнов декларативного API хватает, и лезть ниже незачем.
Но «глубокое погружение» на нём и заканчивать нельзя. Table API — не «синтаксис поверх DataStream»: его планировщик строит собственный граф выполнения, а между представлениями есть мосты (toDataStream/fromDataStream). И часть задач в декларативном API честно не выражается: реакция на каждое событие произвольным кодом, ручное управление состоянием, срабатывание по таймеру, корреляция двух потоков нестандартной логикой, обогащение неблокирующими внешними вызовами, маршрутизация в отдельные потоки — это Java-первый DataStream API. Эта статья — про то, где кончается SQL и что даёт спуск под него. Весь код — в стенде digital-cookbook/messaging/flink/03-datastream: две джобы, одна команда, воспроизводимые числа.
В статье
- Где кончается SQL
- KeyedProcessFunction: реакция на событие и на его отсутствие
- Состояние руками
- Обогащение: broadcast state и async I/O
- Слияние: connect и CoProcessFunction
- Side outputs
- Watermark-стратегия руками
- Мост к SQL и честная граница
- Что дальше
Где кончается SQL
Сначала честная граница — чтобы не спускаться в DataStream там, где он не нужен. Оконные агрегаты, группировки и простые джойны — это SQL, и Table API делает их лучше: короче, с оптимизатором и без ручного управления состоянием. Если задача формулируется как SELECT ... GROUP BY TUMBLE(...) или интервальный JOIN — берите SQL, это модули 00-model/02-sql-ops серии.
DataStream нужен, когда логика перестаёт быть декларативной:
- произвольная реакция на каждое событие — не агрегат, а код с ветвлениями и доступом к контексту;
- состояние под полным контролем — что храним, когда чистим, какой TTL;
- таймеры — сделать что-то через N времени или при отсутствии событий (в SQL нет «сработай, если данных не было»);
- корреляция двух потоков нестандартной логикой — «сматчи A с последним B, но с таким-то таймаутом и такой-то маршрутизацией промахов»;
- обогащение внешним источником — неблокирующие async-вызовы или реплицируемый справочник;
- маршрутизация — развести один поток на несколько по произвольному признаку.
MATCH_RECOGNIZE/CEP покрывает часть паттернов, оставаясь в SQL, — но не всё из списка.
KeyedProcessFunction: реакция на событие и на таймер
Базовый примитив DataStream — KeyedProcessFunction: processElement вызывается на каждое событие ключа с доступом к контексту (ключ, время, сервис таймеров), а onTimer — при срабатывании таймера. На стенде (джоба SilenceDetector) это ловит разрыв в event-time: между событиями ключа прошло больше порога по времени событий.
Здесь критично не перепутать два домена времени — на них построена вся реакция «по таймеру», и подмена одного другим ведёт к неверной модели:
- Event-time таймер срабатывает, когда watermark пройдёт его дедлайн — то есть при приходе более позднего события или на финальном watermark в конце потока. Остановившийся источник watermark не двигает, поэтому по «умершему» ключу event-time таймер сам не сработает.
- Processing-time таймер срабатывает по стенным часам, независимо от данных. Именно он нужен для настоящего liveness — «источник физически замолчал».
Стенд детектит event-time разрыв (детерминированно, воспроизводимо), поэтому таймер тут event-time. И ещё тонкость: при событиях не по порядку дедлайн нельзя двигать назад — хранить надо максимум event-time, а не последний пришедший:
public void processElement(Event e, Context ctx, Collector<String> out) throws Exception {
long wm = ctx.timerService().currentWatermark();
if (wm != Long.MIN_VALUE && e.eventTime <= wm) { // LATE: watermark уже прошёл — политика явная
ctx.output(LATE, e); // (иначе таймер зарегистрируется в прошлом -> ложный GAP)
return;
}
Long max = maxTs.value(); // ValueState<Long> — МАКСИМУМ event-time
if (max != null && e.eventTime <= max) { // out-of-order (не late): дедлайн не откатываем
ctx.output(OUT_OF_ORDER, e);
return;
}
if (max != null) ctx.timerService().deleteEventTimeTimer(max + GAP);
maxTs.update(e.eventTime);
ctx.timerService().registerEventTimeTimer(e.eventTime + GAP);
}
public void onTimer(long ts, OnTimerContext ctx, Collector<String> out) throws Exception {
Long max = maxTs.value();
if (max != null && ts == max + GAP) { // null-safe: только актуальный дедлайн
out.collect("GAP key=" + ctx.getCurrentKey());
maxTs.clear();
}
}Три тонкости, которые здесь легко упустить (и на которых ловится «наивная» версия): late-guard до работы со state (иначе событие ниже watermark зарегистрирует таймер в уже пройденном времени — и породит ложный мгновенный GAP), хранение максимума event-time (а не «последнего»), и null-safe onTimer (после clear() state пуст).
На стенде это даёт 4 разрыва (GAP) — и читать их надо честно: один реальный (когда sensor-3 замолкает в середине потока) и три от финального MAX_WATERMARK ограниченного источника, который в конце «закрывает» все ключи (в бесконечном потоке этих трёх не было бы). Событий не по порядку в прогоне ровно 19 — все уходят в side output out-of-order и ложных алертов не создают; опоздавших (ниже watermark) в детерминированном in-order-генераторе 0, но политика в коде явная.
Для wall-clock liveness («источник замолчал») — тот же каркас, но таймер processing-time:
// deadline: ValueState<Long> — храним текущий дедлайн, чтобы отменить прежний таймер
Long prev = deadline.value();
if (prev != null) ctx.timerService().deleteProcessingTimeTimer(prev); // иначе старые таймеры сработают
long next = ctx.timerService().currentProcessingTime() + GAP;
ctx.timerService().registerProcessingTimeTimer(next);
deadline.update(next);
// onTimer: Long cur = deadline.value(); if (cur != null && ts == cur) { emit liveness; deadline.clear(); }Здесь важна та же дисциплина, что и с event-time: отменять прежний таймер и хранить дедлайн в state, иначе каждое событие оставит свой таймер, и они будут срабатывать один за другим, порождая ложные алерты. (В bounded-стенде processing-time liveness не показываем: источник завершается раньше дедлайна таймера, и число несводимо к воспроизводимому.)
Состояние руками
В processElement состояние — не «магия под SQL», а явные примитивы, привязанные к ключу: ValueState<T> (одно значение, как выше), ListState<T> (накопить список), MapState<K,V> (словарь). Живут они в state backend (heap или RocksDB — разобрано в статье про состояние и exactly-once), участвуют в чекпоинтах и переживают перезапуск.
Ключевое, чего SQL не даёт, — контроль над жизненным циклом: когда состояние чистить (state.clear()) и какой у него TTL (StateTtlConfig). Но TTL — не «само собой освободится»: он отсчитывается по processing-time, а физическая очистка зависит от backend и стратегии (при доступе к записи, фоновым проходом, при компакции RocksDB). Без явной очистки логически истёкшая запись может ещё занимать место — просто перестаёт быть видимой. В декларативном джойне жизненный цикл решает планировщик; здесь он на вас — и это ответственность, а не только свобода.
Обогащение: broadcast state и async I/O
Обогащение — присоединить к событию данные извне. SQL умеет lookup join, но без нужного контроля. DataStream даёт два инструмента:
Broadcast state — небольшой справочник или конфиг (правила, пороги, курсы), который реплицируется на все параллельные инстансы оператора и обновляется отдельным потоком. Событие джойнится со справочником локально, без шардирования по ключу справочника:
MapStateDescriptor<String, Rule> rules = new MapStateDescriptor<>("rules", String.class, Rule.class);
events.connect(ruleStream.broadcast(rules))
.process(new BroadcastProcessFunction<Event, Rule, Enriched>() {
public void processElement(Event e, ReadOnlyContext ctx, Collector<Enriched> out) {
Rule r = ctx.getBroadcastState(rules).get(e.type); // локальный lookup
out.collect(enrich(e, r));
}
public void processBroadcastElement(Rule r, Context ctx, Collector<Enriched> out) {
ctx.getBroadcastState(rules).put(r.type, r); // обновление справочника
}
});Async I/O (AsyncDataStream.unorderedWait) — когда обогащение требует внешнего вызова (БД, HTTP): синхронный вызов застопорил бы весь оператор, async отправляет запросы, не блокируя пайплайн, и собирает ответы по мере готовности. Цена не бесплатна: unorderedWait меняет порядок записей (есть orderedWait, но он дороже по латентности), а обязательными становятся timeout на запрос, ограничение числа одновременных вызовов (capacity), стратегия retry и идемпотентность внешнего эффекта — async усиливает и нагрузку, и повторы. (В стенде не поднимаем ради «одной команды» — нужен внешний сервис-заглушка.)
Слияние: connect и CoProcessFunction
Слить два потока SELECT ... JOIN-ом — это SQL. Но когда нужна кастомная логика корреляции с таймаутом и маршрутизацией промахов, идут в connect() + KeyedCoProcessFunction: два потока делят состояние по ключу, а два метода — processElement1/processElement2 — обрабатывают каждый свой вход.
На стенде (джоба EnrichAndCorrelate) — заказы и платежи по orderId: кто пришёл первым, ждёт в состоянии; пара — эмитим с вычисленным latencyMs (обогащение); нет пары за таймаут — таймер уводит в side output.
public void processElement2(Payment p, Context ctx, Collector<String> out) throws Exception {
Order o = pendingOrder.value();
if (o != null) {
long delta = p.eventTime - o.eventTime;
if (delta >= 0 && delta <= MATCH_TIMEOUT_MS) { // платёж В ОКНЕ [order, order + timeout]
out.collect("MATCH orderId=" + o.orderId + " latencyMs=" + delta);
} else { // вне окна: таймер бы не спас — матч уже случился бы
ctx.output(UNMATCHED, "out-of-window orderId=" + o.orderId + " deltaMs=" + delta);
}
clear(ctx);
} else { // пары ещё нет — ждём
pendingPayment.update(p);
arm(ctx, p.eventTime + MATCH_TIMEOUT_MS); // таймер ЛИШЬ ограничивает удержание state
}
}Здесь легко ошибиться: таймер ограничивает удержание state, но не сам матч. Без проверки окна платёж со временем order + 100с, пришедший до срабатывания таймаут-таймера (watermark ещё не дошёл до дедлайна), спокойно соединился бы с заказом — вопреки заявленным 8 c. Поэтому матч действителен, только если delta ∈ [0, timeout] по event-time; вне окна — в side output. На стенде это видно: демонстративный платёж на 100 c позже заказа даёт ровно 1 запись out-of-window.
Живые числа стенда: 265 скоррелированных пар (в окне), 68 несостыкованных по таймауту и 1 отвергнутая как вне окна — все в side output. Такую логику (произвольный порядок, окно, разная реакция на промах заказа vs платежа) интервальный SQL-JOIN так не описывает.
И контракт ключа: ValueState держит по одной записи каждого типа — модель предполагает ровно один заказ и один платёж на orderId; повторный элемент перезаписал бы прежний. Для потоков с дубликатами нужна дедупликация или ListState/MapState с явной политикой матчинга. Таймаут здесь — event-time (срабатывает по watermark, в стенде — на финальном watermark), не по стенным часам.
Side outputs
Обычный оператор отдаёт один поток. Часто же нужно развести его: основной результат — в один сток, аномалии/промахи/опоздавшие — в другой. Это side output через OutputTag: в processElement/onTimer пишут в тег, а ниже забирают getSideOutput(tag).
public static final OutputTag<String> UNMATCHED = new OutputTag<>("unmatched", Types.STRING);
// ...в onTimer, когда пары так и нет:
ctx.output(UNMATCHED, "order orderId=" + o.orderId + " без платежа за " + MATCH_TIMEOUT_MS + "ms");
// ...в main:
matched.getSideOutput(UNMATCHED).print("UNMATCHED");Несматченные на стенде (68 по таймауту + 1 вне окна) идут именно так. Классический второй кейс — опоздавшие события: то, что пришло ниже watermark, увести в отдельный поток (late) вместо тихой потери — ровно та политика из детектора разрывов выше. (В детерминированном in-order-генераторе опоздавших 0, а их число вообще зависит от тайминга watermark’ов и невоспроизводимо, поэтому в живых числах side output показан на несматченных.)
Watermark-стратегия руками
В SQL watermark настраивается в DDL; в DataStream им управляют явно — WatermarkStrategy на источнике:
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) // допуск неупорядоченности
.withTimestampAssigner((e, ts) -> e.eventTime);Здесь же — контроль над извлечением времени события (withTimestampAssigner), обработкой простаивающих источников (withIdleness) и — на уровне оконных операторов — allowed lateness. Механика event-time и watermark’ов подробно разобрана в статье про модель Flink.
Мост к SQL и честная граница
DataStream и SQL — не «или-или», а слои одной джобы: StreamTableEnvironment конвертирует DataStream в Table и обратно, так что тяжёлую декларативную часть (агрегаты, джойны) пишут на SQL, а нестандартную (таймеры, корреляция, обогащение) — на DataStream, в одном пайплайне.
Граница простая: что выражается декларативно — оставляйте в SQL (короче и с оптимизатором); спускайтесь в DataStream ровно за тем, что SQL не выражает — произвольным состоянием, таймерами, кастомной корреляцией, broadcast/async-обогащением и маршрутизацией. «Глубокое погружение» — это уметь и то, и это, и знать, где проходит шов.
Что дальше
Остальные срезы серии: модель Flink: event-time, watermarks, окна, состояние и exactly-once (те же ValueState/backends, что здесь, но со стороны отказоустойчивости) и Flink SQL, коннекторы и эксплуатация (декларативный фасад и прод). Как Flink встаёт в общую картину стриминга — на карте messaging. Обе джобы этой статьи — в стенде (одна команда, живые числа). Если тема зашла — напишите в Telegram-группе, что из DataStream разобрать подробнее.
Комментарии