Streaming SQL на собеседовании Data Engineer
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 для быстрой стриминговой агрегации без отдельного кластера обработки».
Flink 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 без ручного управления состоянием.
Окна и 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-базе, которые обновляются по расписанию. Разница в том, что стриминговый результат поддерживается инкрементально и непрерывно, а не пересчитывается целиком по крону.
Связанные темы
- Kafka на собесе DE
- Spark Structured Streaming для DE
- Apache Flink для DE
- Kafka Streams для DE
- Подготовка к собесу Data Engineer
FAQ
Чем streaming SQL отличается от обычного SQL?
Обычный SQL выполняется один раз и возвращает снапшот данных. Streaming SQL — это continuous query: он живёт постоянно и инкрементально обновляет результат по мере поступления событий. Результат всегда актуален без повторного запуска.
Что выбрать — ksqlDB, Flink или Materialize?
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+ вопросами для собесов.