Flink SQL, коннекторы и эксплуатация: от Table API до прода

Писать Flink-джобы на низкоуровневом API нужно не всегда: Table API и Flink SQL закрывают большую часть задач декларативно, включая стриминговые джойны и агрегации. Плюс то, без чего не выйти в прод: коннекторы (Kafka, CDC, lakehouse), управление параллелизмом, диагностика backpressure, тюнинг чекпоинтов и апгрейд джобы через savepoint без потери состояния. Завершает погружение практикой эксплуатации

Низкоуровневый DataStream API мощен, но избыточен для большинства задач: посчитать оконные агрегаты или соединить два потока проще декларативно. Flink SQL и Table API дают именно это — стриминговые запросы на SQL с тем же движком event-time и состояния под капотом. А дальше начинается то, что отличает демо от прода: подключить реальные источники и стоки, понять, почему джоба тормозит (backpressure), настроить чекпоинты и — самое болезненное для стрим-джоб — обновить код, не потеряв состояние. Эта статья завершает погружение практикой эксплуатации.

Финальная статья серии «Apache Flink: глубокое погружение». Опирается на модель event-time и watermark’ов из первой статьи и на состояние с exactly-once из второй. Всё ниже прогнано на живом стенде digital-cookbook/messaging/flink/02-sql-ops — только Flink SQL, без единой строки Java.

Белка-оператор Apache Flink сводит два потока-трубы orders и payments в ±1-минутном time-window join; коннекторы Kafka/CDC/datagen слева, рычаг SAVEPOINT, белки-помощники добавляют слоты (RESCALE 6→8), gauge BACKPRESSURE, а распухший бак REGULAR JOIN (INFINITE STATE) предупреждает о вечном состоянии

В статье

Flink SQL и Table API: декларативный стриминг

Table API и Flink SQL — это не отдельный движок, а декларативный фасад над тем же ядром, что и DataStream: тот же event-time, те же watermark’и, то же управляемое состояние и те же чекпоинты. Разница в том, что вы описываете что посчитать, а планировщик сам решает как — какие операторы построить, где хранить состояние, как его разбить по ключу.

Основа модели — динамическая таблица: непрерывный запрос над бесконечным потоком, результат которого сам меняется во времени. Поток и таблица здесь двойственны — таблица описывается через CREATE TABLE ... WITH ('connector'=...), где WITH подключает физический источник или сток, а сама таблица логически представляет либо поток append-only-строк, либо changelog (поток вставок, обновлений и удалений — +I, -U, +U, -D). Какой из режимов, зависит от запроса: простой SELECT ... WHERE над источником даёт append, а агрегат или regular join — changelog с ретракциями. Это различие всплывёт в разделе про джойны — оно определяет, что уедет в Kafka-сток.

Таблица orders со стенда — обычная Kafka-таблица с объявленным event-time и watermark’ом (ровно тот механизм из первой статьи):

CREATE TABLE orders (
    order_id  INT,
    user_id   INT,
    amount    INT,
    ts        TIMESTAMP(3),
    WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka', 'topic' = 'orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'properties.group.id' = 'flink-orders',
    'scan.startup.mode' = 'earliest-offset',
    'format' = 'json'
);

WATERMARK FOR ts AS ts - INTERVAL '5' SECOND объявляет столбец ts меткой события и допускает опоздание до 5 секунд — дальше вся оконная логика и джойны работают в event-time, а не по времени обработки. Запрос запускается через INSERT INTO, живёт как долгоиграющая джоба и пишет результат непрерывно. Несколько таких INSERT можно объединить в один джоб через EXECUTE STATEMENT SET ... BEGIN ... END — именно так на стенде оба фидера стартуют одной командой.

Помимо джойнов и агрегаций Flink SQL умеет распознавание паттернов последовательностей (CEP) через стандартный оператор MATCH_RECOGNIZE — «найди в потоке A, за которым в течение минуты B, но не C». Это отдельная большая тема; здесь достаточно знать, что она есть и живёт в том же SQL, не требуя спуска в DataStream.

Коннекторы: kafka, datagen, print, CDC

Коннектор в WITH (...) — это мост между логической таблицей и физическим миром. На стенде задействованы четыре, и каждый закрывает свою роль:

Коннектор Роль Где на стенде
datagen синтетический источник для нагрузки и демо фидеры orders/payments (rows-per-second, sequence, min/max)
kafka source и sink; читает/пишет топики, держит offset’ы в чекпоинтах таблицы orders, payments, enriched
print сток в stdout TaskManager — для быстрой отладки запроса — (включается точечно)
CDC (Debezium) источник изменений из БД как поток прод-замена datagen для одной из сторон джойна

datagen удобен именно для стенда: он генерирует строки без внешних зависимостей. Заказы на стенде — это последовательность order_id и случайные поля:

CREATE TABLE orders_src (
    order_id  INT,
    user_id   INT,
    amount    INT
) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '5',
    'fields.order_id.kind' = 'sequence', 'fields.order_id.start' = '1', 'fields.order_id.end' = '1000000',
    'fields.user_id.min' = '1', 'fields.user_id.max' = '100',
    'fields.amount.min' = '1', 'fields.amount.max' = '500'
);

Kafka-коннектор используется на стенде в версии 3.3.0-1.20 (внешний артефакт, версионируется отдельно от ядра Flink — суффикс -1.20 означает совместимость с Flink 1.20). Как sink он умеет транзакционную запись: enriched-таблица пишется с 'sink.delivery-guarantee' = 'exactly-once' и префиксом транзакционного ID — это тот самый two-phase commit из второй статьи, где коммит в Kafka привязан к завершению чекпоинта.

В проде источником одной из сторон джойна часто оказывается не datagen, а CDC — поток изменений реальной таблицы БД через Debezium. С точки зрения Flink SQL это просто другой connector: таблица объявляется так же, а изменения строк приходят как changelog (INSERT/UPDATE/DELETE). Как устроен сам захват изменений — в статье про CDC из Postgres в Kafka; как рядом живут Schema Registry и Kafka Connect — в обзоре экосистемы. Стоком же в аналитику всё чаще служит не Kafka, а lakehouse — Flink пишет в открытые табличные форматы; устройство этого слоя — Iceberg/DeltaСкоро.

Стриминговые джойны: interval vs regular

Соединить два бесконечных потока — операция принципиально иная, чем джойн двух таблиц: у потоков нет конца, и «дождаться второй стороны» можно только за счёт состояния. Ключевой вопрос — сколько состояния джойн обязан держать. Ответ определяет, будет джоба жить годами или упрётся в диск.

Interval join ограничивает состояние временным окном. Со стенда — обогащение заказа методом платежа в пределах ±1 минуты по event-time:

CREATE TABLE enriched (order_id INT, user_id INT, amount INT, `method` STRING)
WITH ('connector'='kafka','topic'='orders-enriched','properties.bootstrap.servers'='kafka:9092',
    'format'='json','sink.delivery-guarantee'='exactly-once','sink.transactional-id-prefix'='enriched-tx');

INSERT INTO enriched
SELECT o.order_id, o.user_id, o.amount, p.`method`
FROM orders o
JOIN payments p
  ON o.order_id = p.order_id
 AND p.ts BETWEEN o.ts - INTERVAL '1' MINUTE AND o.ts + INTERVAL '1' MINUTE;

Условие p.ts BETWEEN o.ts - INTERVAL '1' MINUTE AND o.ts + INTERVAL '1' MINUTE — не фильтр, а граница состояния. Как только watermark’и с обеих сторон проходят конец окна, Flink знает, что более ранняя строка уже не найдёт пару, и выбрасывает её из состояния. Состояние ограничено окном и не растёт бесконечно, а результат — append-only: каждая совпавшая пара уезжает в Kafka один раз, поэтому обычного JSON-стока хватает.

Regular join — тот же запрос без условия по ts (ON o.order_id = p.order_id). Синтаксически проще, семантически опаснее: раз временной границы нет, любая будущая строка любой из сторон теоретически может найти пару, и Flink обязан хранить всё состояние обеих сторон вечно. Состояние растёт монотонно с числом уникальных ключей, RocksDB пухнет, чекпоинты тяжелеют, рано или поздно джоба падает по диску. Отдельно: результат regular join — не обязательно append. Будет ли он changelog с ретракциями, зависит от характера входов и запроса (обновляющиеся входы или меняющиеся во времени соответствия заставляют отзывать и переиздавать прежние результаты); тогда в append-only Kafka-сток без upsert-режима его уже не положишь. Но главный операционный риск regular join — не формат вывода, а именно неограниченно растущее состояние (Joins в Flink SQL, docs 1.20).

Interval join Regular join
Условие p.ts BETWEEN o.ts - X AND o.ts + X только o.key = p.key
Границы состояния окно по event-time нет — всё, вечно
Рост состояния ограничен окном монотонный, неограниченный
Результат append-only append или changelog — зависит от входов/запроса
Сток обычный Kafka-sink нужен upsert/changelog-сток
Когда брать события связаны по времени (заказ↔платёж) справочник × поток, редкие ключи, или обязательный TTL

Regular join не «плохой» — он нужен, когда связь не временная (обогащение потока небольшим справочником). Но тогда неограниченный рост гасят TTL состояния (table.exec.state.ttl): состояние без обращений через заданное время выбрасывается — ценой того, что очень поздняя пара не сойдётся. Отдельная разновидность — temporal join (FOR SYSTEM_TIME AS OF): обогащение потока версией справочника, актуальной на момент события. Он держит версионированную таблицу, а не полное декартово состояние, и предназначен ровно для «подтянуть курс/тариф, действовавший тогда».

Практический вывод один: всегда осознанно отвечайте, чем ограничено состояние джойна — окном, TTL или версионностью. Незамеченный regular join — самая частая причина, по которой стрим-джоба, работавшая неделями, внезапно умирает.

Эксплуатация: savepoint, rescale, апгрейд джобы

Стрим-джоба живёт долго, и рано или поздно её нужно перезапустить: обновить логику, добавить слотов, переехать. Наивная остановка теряет состояние окон — а с ним точность агрегатов и джойнов. Инструмент, который это решает, — savepoint.

Checkpoint против savepoint — два снимка состояния с разным назначением:

Checkpoint Savepoint
Кто инициирует Flink автоматически (на стенде execution.checkpointing.interval: 10s) оператор вручную
Назначение восстановление после падения апгрейд, rescale, миграция
Формат оптимизирован под движок, может быть инкрементальным (RocksDB) портируемый, самодостаточный
Время жизни ротируется движком живёт, пока не удалите
Куда state.checkpoints.dir state.savepoints.dir

На стенде оба каталога заданы (file:///tmp/flink-checkpoints и file:///tmp/flink-savepoints), backend — RocksDB, чекпоинт раз в 10 секунд. Чекпоинт — про «пережить сбой сам по себе»; savepoint — про «я сознательно останавливаю джобу и хочу поднять её иначе».

Rescale — смена параллелизма без потери состояния. Джоба джойна держит состояние обеих сторон в окне; нужно раздать его на больше слотов. Порядок со стенда (проверено живьём):

docker compose exec jobmanager ./bin/flink list
docker compose exec jobmanager ./bin/flink stop --savepoint file:///tmp/flink-savepoints <JOB_ID>
# → выведет путь, например file:///tmp/flink-savepoints/savepoint-xxxx

flink stop --savepoint останавливает джобу, сняв согласованный снимок, и печатает его путь. Дальше — перезапуск из этого снимка с новым параллелизмом, прямо в SQL client:

SET 'execution.savepoint.path' = 'file:///tmp/flink-savepoints/savepoint-xxxx';
SET 'parallelism.default' = '4';
-- заново ТОЛЬКО INSERT INTO enriched (таблица в сессии уже создана)

Джоба поднимается с восстановленным состоянием окна, но уже на 4 слотах. Flink перераспределяет keyed-состояние по key-group’ам между новыми слотами — руками ничего перекладывать не нужно. В orders-enriched на стыке нет ни потерь, ни дублей: exactly-once-сток плюс savepoint дают консистентную границу. На живом прогоне после rescale в топике оказалось 755 сообщений — ровно 755 уникальных order_id подряд, с первого по 755-й, без пропусков и повторов.

Два подвоха, на которых этот сценарий спотыкается на практике. Первый — CREATE TABLE повторять не надо: в той же сессии таблица уже есть, и команда упадёт с TableAlreadyExistException; а если сессия SQL client закрывалась, вместе с ней пропал весь каталог (он in-memory) и таблицы придётся объявить заново. Второй — слоты. Фидеры можно не трогать, они независимы от джобы джойна, но слот они занимают: фидер-джоба держит один, а джойн после rescale просит четыре. На кластере с четырьмя слотами перезапуск встанет в CREATED и будет молча ждать пятый — без ошибки в интерфейсе, джоба при этом числится RUNNING. Поэтому в стенде taskmanager.numberOfTaskSlots: 8.

Апгрейд джобы — тот же механизм: снять savepoint, поменять запрос, подняться из снимка. Но здесь важная оговорка, специфичная для SQL. Flink генерирует операторы и раскладку состояния из текста запроса; заметное изменение запроса меняет топологию, и восстановление из старого savepoint может отвалиться по несовместимости состояния. Безопасны прежде всего изменения, не трогающие stateful-часть (например, вычисляемые проекции в SELECT); изменение ключей джойна, окон или агрегатов — повод планировать миграцию отдельно, а не рассчитывать на бесшовный restore. Практическое правило: rescale и переезд savepoint переживает почти всегда, а вот «переписать SQL и поднять из старого снимка» — проверяйте на стейджинге, не на проде.

Backpressure и наблюдаемость

Backpressure — это сигнал, что оператор ниже по потоку не успевает за тем, что в него льют: очереди заполняются, давление распространяется вверх до источника, и тот притормаживает чтение. Для Flink это штатный механизм саморегуляции (общая теория — в backpressure и load shedding и теории очередейСкоро), но устойчивый backpressure означает, что джоба недообеспечена ресурсами.

Главный инструмент диагностики — Flink Dashboard (localhost:8081): джобы, их граф операторов, чекпоинты и backpressure. Каждый оператор в графе окрашивается по метрикам busy/backpressured; красный/розовый узел показывает, где именно копится — то есть какой оператор является узким местом, а не просто «джоба тормозит».

Воспроизвести это на стенде — прямой эксперимент: поднять rows-per-second у фидеров до сотен и уменьшить число слотов TaskManager. Источник начнёт производить быстрее, чем джойн успевает обрабатывать в оставшихся слотах, и в Dashboard оператор джойна окрасится в backpressure, а давление уползёт вверх к источнику. Отсюда и лечение по порядку убывания частоты причин:

  • не хватает параллелизма — раздать оператор на больше слотов (тот самый rescale через savepoint выше);
  • перекос по ключу — один «горячий» ключ грузит один слот, остальные простаивают; лечится ключом партиционирования, а не числом слотов;
  • медленный сток — узкое место снаружи Flink (внешняя БД, S3); тогда упирается sink-оператор, и добавлять слоты джойну бесполезно;
  • раздутое состояние (см. regular join выше) — RocksDB упирается в диск/IO, чекпоинты растягиваются.

Второй экран наблюдаемости — вкладка checkpoints: их длительность и размер, растёт ли состояние, укладывается ли чекпоинт в 10-секундный интервал. Растущая длительность чекпоинта на фоне ровной нагрузки — почти всегда симптом неограниченного состояния. Для продовых пайплайнов те же сигналы (лаг, свежесть, рост состояния) выносят в общий мониторинг данных — см. data observability и lineageСкоро.

Эксплуатационные ошибки

  1. Regular join там, где нужен interval. Убрали (или забыли) условие по ts — состояние копится вечно, джоба неделями работает, потом умирает по диску. Всегда осознанно ограничивайте состояние джойна окном, TTL или версионностью.
  2. Append-сток под changelog-результат. Regular join и агрегаты выдают ретракции; обычный Kafka JSON-sink их не выражает. Под changelog нужен upsert/changelog-сток — иначе результат некорректен или джоба не стартует.
  3. Забытый или заниженный watermark. Event-time-джойны и окна не продвинутся без watermark’а, а слишком строгий (маленький допуск опоздания) молча выбрасывает опоздавшие события из окна. Допуск подбирают под реальный разброс времён (INTERVAL '5' SECOND на стенде — это про демо-темп).
  4. Остановка джобы без savepoint. Просто убить джобу и запустить заново — потерять состояние окон; на стыке возникнут пропуски или дубли. Апгрейд и rescale — только через flink stop --savepoint.
  5. Апгрейд SQL с несовместимым состоянием. Поменяли ключи джойна/окна/агрегаты и ждёте бесшовного restore из старого savepoint — а раскладка состояния изменилась, восстановление падает. Изменения stateful-части проверяйте на стейджинге, планируйте миграцию.
  6. exactly-once-сток без понимания цены. sink.delivery-guarantee = exactly-once привязывает видимость записей к завершению чекпоинта: потребитель с read_committed увидит данные только после коммита транзакции. Интервал чекпоинта = нижняя граница end-to-end-латентности; это осознанный размен, а не бесплатная гарантия.
  7. Диагностика backpressure «на глаз». «Джоба тормозит» без взгляда в Dashboard ведёт к слепому наращиванию слотов. Сначала найдите узкий оператор (цвет busy/backpressured), потом лечите причину — параллелизм, перекос ключа или сток.

Checklist: эксплуатационная готовность

  1. У каждого стримингового джойна явно определено, чем ограничено состояние — окном (interval), TTL или версионностью (temporal)?
  2. Тип результата запроса (append vs changelog) соответствует выбранному стоку?
  3. У event-time-таблиц объявлен watermark, а его допуск на опоздание подобран под реальный разброс, не под демо-темп?
  4. Апгрейд и rescale идут только через savepoint (flink stop --savepoint), а не через убийство джобы?
  5. Изменения запроса, затрагивающие состояние, прогнаны на стейджинге на предмет совместимости savepoint?
  6. Интервал чекпоинта осознанно связан с требуемой end-to-end-латентностью exactly-once-стока?
  7. Есть доступ к Flink Dashboard и понимание, как по нему найти узкий оператор и растущее состояние?
  8. Каталоги checkpoints/savepoints лежат на надёжном хранилище (в проде — S3/HDFS, не локальный /tmp)?

Если на большинство — «да», Flink SQL-джоба у вас не только считает правильный результат на демо, но и переживает апгрейд, масштабирование и нагрузку в проде.

Демо и версии

  • Demo: digital-cookbook/flink/02-sql-ops/ — стриминговый interval join двух Kafka-топиков на Flink SQL (обогащение заказа методом платежа в окне ±1 минута) плюс эксплуатационная операция savepoint + rescale (смена параллелизма) без потери состояния. Всё через SQL client, без Java: feeders.sql (datagen → топики orders/payments) и join.sql (interval join → топик orders-enriched, exactly-once). Что попробовать самому: поднять rows-per-second и уменьшить слоты, чтобы увидеть backpressure в Dashboard; убрать условие по ts, чтобы получить regular join с растущим состоянием.
  • Версии (пиновать — поведение SQL-планировщика и коннекторов чувствительно к минорам): Apache Flink 1.20.1 (Java 17), Apache Kafka 3.9.0 (KRaft), Flink Kafka SQL connector 3.3.0-1.20. State backend — RocksDB, execution.checkpointing.interval: 10s, taskmanager.numberOfTaskSlots: 8 (восьми, а не четырёх: фидер-джоба держит слот, а джойн после rescale просит четыре).

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии