Streaming SQL на собеседовании Data Engineer

Проверь себя · 1/3разбор после ответа
Что делает оператор DISTINCT в SELECT-запросе?

Зачем нужен streaming SQL

Классическая потоковая обработка требует писать код на Java, Scala или Python и вручную управлять состоянием, окнами и чекпоинтами. Streaming SQL снимает этот порог: аналитик или инженер описывает потоковую логику привычным SQL, а движок сам разворачивает её в распределённый стриминговый пайплайн.

Ключевое отличие от обычного SQL — запрос не выполняется один раз и не завершается. Это continuous query: он живёт постоянно и обновляет результат по мере поступления новых событий.

-- потоковая агрегация: результат пересчитывается на каждом новом событии
CREATE STREAM orders_by_country AS
SELECT country, COUNT(*) AS cnt, SUM(amount) AS revenue
FROM orders_stream
GROUP BY country;

В батче вы гоняете такой запрос раз в час по расписанию. В стриминге результат живой: пришло новое событие — счётчик по стране обновился в тот же момент. На собесе важно уметь объяснить именно этот сдвиг мышления: от «запрос вернул снапшот» к «запрос поддерживает всегда-актуальный результат».

ksqlDB

ksqlDB — это SQL-движок поверх Kafka от Confluent. Он читает и пишет прямо в топики Kafka, поэтому естественно ложится на экосистему, где Kafka уже есть.

CREATE STREAM orders WITH (KAFKA_TOPIC='orders', VALUE_FORMAT='JSON');

CREATE TABLE country_revenue AS
SELECT country, SUM(amount) AS total
FROM orders
GROUP BY country
EMIT CHANGES;

Ключевое различие в ksqlDB — stream против table. Stream — это неограниченный лог событий (каждая запись самостоятельна), а table — это текущее состояние по ключу (последнее значение). EMIT CHANGES превращает запрос в push-запрос: он не отдаёт разовый снапшот, а постоянно шлёт обновления результата в даунстрим. На собесе про Kafka-стек ksqlDB спрашивают как «SQL для быстрой стриминговой агрегации без отдельного кластера обработки».

Apache Flink — самый мощный из стриминговых движков, и его SQL-слой поддерживает полноценную обработку по событийному времени с watermark и строгими гарантиями exactly-once.

CREATE TABLE orders (
  order_id BIGINT,
  amount DECIMAL,
  ts TIMESTAMP(3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH ('connector' = 'kafka', ...);

SELECT TUMBLE_START(ts, INTERVAL '1' HOUR) AS window_start,
       SUM(amount) AS hourly_revenue
FROM orders
GROUP BY TUMBLE(ts, INTERVAL '1' HOUR);

Это стандартный SQL с расширениями под окна (TUMBLE, HOP, SESSION) и объявлением watermark прямо в DDL. Flink выбирают, когда нужны точная семантика по событийному времени, корректная обработка опоздавших данных и отказоустойчивость на больших объёмах. Именно Flink чаще всего всплывает на DE-собесах как «серьёзный» стриминг в отличие от лёгкого ksqlDB.

Materialize

Materialize — потоковая база, совместимая с протоколом Postgres. Главная идея — инкрементально поддерживаемые материализованные представления: вы пишете обычный SQL с джойнами и агрегациями, а Materialize держит результат всегда свежим, пересчитывая только то, что изменилось.

CREATE SOURCE orders FROM KAFKA BROKER '...' TOPIC 'orders';

CREATE MATERIALIZED VIEW country_revenue AS
SELECT country, SUM(amount) FROM orders GROUP BY country;

SELECT * FROM country_revenue;  -- всегда актуальный результат

Ценность в том, что не нужно вручную гонять обновления представления — оно поддерживается инкрементально. Плюс совместимость с Postgres: можно подключаться привычными BI-инструментами и драйверами. Materialize берут, когда хочется сложных стриминговых джойнов на языке обычного SQL без ручного управления состоянием.

Готовишься к собесу Data Engineer?
Spark, Airflow, ClickHouse, SQL для DE — вопросы с разборами в Telegram
Тренировать DE в Telegram

Окна и watermark

Окна — центральная тема стриминга, и на собесе её спрашивают почти всегда. Поток бесконечен, поэтому агрегировать «всё» нельзя — данные бьют на окна.

  • Tumbling (перекатывающиеся) — непересекающиеся окна фиксированного размера: каждую минуту, каждый час. Каждое событие попадает ровно в одно окно.
  • Hopping / sliding (скользящие) — окна фиксированного размера с шагом меньше размера, поэтому они перекрываются. Например, окно 10 минут с шагом 1 минута.
  • Session (сессионные) — окна переменной длины, которые закрываются после паузы без событий. Удобны для группировки активности пользователя.

Отдельно спрашивают про event-time против processing-time. Event-time — время, когда событие реально произошло (проставлено в самом сообщении); processing-time — когда его обработал движок. Корректная аналитика почти всегда требует event-time, иначе задержки в сети исказят агрегаты.

С event-time связан watermark — это метка «событий раньше такого-то времени больше не ждём». Watermark позволяет закрывать окна и решать, что делать с опоздавшими данными. Понимание watermark отличает кандидата, который реально трогал стриминг, от того, кто читал про него.

Сценарии применения

  • Real-time дашборды. Метрики обновляются в реальном времени, а не раз в час батчем. Данные на экране всегда свежие.
  • Операционная аналитика. Аномалии видны через секунды после того, как случились, а не на следующий день в отчёте.
  • Триггеры и алерты. Поток + окно + условие → уведомление. Например, «за 5 минут больше 100 неудачных платежей с одного IP» — сразу алерт.
  • Streaming ETL. Топик Kafka → на лету очистить и обогатить → записать в другой топик или витрину. Трансформация без промежуточного батч-джоба.

Как это спрашивают на собесе

Streaming SQL редко бывает отдельной секцией — обычно он всплывает внутри system design или блока про стриминг. На что смотрит интервьюер:

  • Понимаете ли разницу батч vs стриминг. Ждут, что вы объясните continuous query и инкрементальный результат, а не просто «SQL по потоку».
  • Окна и время. Практически гарантированный вопрос — типы окон и почему event-time важнее processing-time. Без watermark ответ будет неполным.
  • Гарантии доставки. Как движок обеспечивает exactly-once, что происходит с опоздавшими и дублированными событиями.
  • Выбор инструмента. «Когда ksqlDB, а когда Flink?» Сильный ответ — ksqlDB для лёгких агрегаций поверх Kafka, Flink для сложной event-time логики и больших объёмов, Materialize для стриминговых джойнов на языке Postgres.

Типичная ошибка — путать streaming SQL с обычными материализованными представлениями в OLAP-базе, которые обновляются по расписанию. Разница в том, что стриминговый результат поддерживается инкрементально и непрерывно, а не пересчитывается целиком по крону.

Связанные темы

FAQ

Чем streaming SQL отличается от обычного SQL?

Обычный SQL выполняется один раз и возвращает снапшот данных. Streaming SQL — это continuous query: он живёт постоянно и инкрементально обновляет результат по мере поступления событий. Результат всегда актуален без повторного запуска.

ksqlDB хорош для лёгких агрегаций прямо поверх Kafka без отдельного кластера. Flink — для сложной обработки по событийному времени, watermark и exactly-once на больших объёмах. Materialize — для стриминговых джойнов на языке обычного Postgres-совместимого SQL.

Зачем нужен watermark?

Watermark — это метка времени, после которой движок считает, что более ранних событий уже не будет. Он позволяет закрывать окна и принимать решение об опоздавших данных. Без watermark в event-time обработке окна нельзя корректно завершить.

В чём разница event-time и processing-time?

Event-time — время, когда событие реально произошло (записано в сообщении). Processing-time — когда движок его обработал. Из-за сетевых задержек эти времена расходятся, поэтому корректная аналитика почти всегда строится на event-time.

Это официальная информация?

Нет. Статья основана на документации ksqlDB, Apache Flink и Materialize, а также на опыте прохождения собеседований. Конкретный стек зависит от компании.


Тренируйте Data Engineering — откройте тренажёр с 1500+ вопросами для собесов.