Потоковая обработка данных
Поверх лога работают движки - Flink, Spark Structured Streaming, Kafka Streams, и они отличаются моделью состояния и восстановления. Собес проверяет, понимаешь ли ты, что чекпоинт это единственная память джобы, а replay - суперсила лога.
Стержень: состояние живёт в state store и чекпоинтится, а исправления и новые пайплайны прогоняются через replay истории топика.
// Формулировки: «чем Flink отличается от Spark Streaming?», «что такое чекпоинт?», «что такое replay и зачем retention сырья?».
Классы движков
Flink - нативный event-time стриминг с состоянием, exactly-once и низкой латентностью. Spark Structured Streaming работает микробатчами и даёт единый API с батчем. Kafka Streams - библиотека внутри приложения, без отдельного кластера.
В Structured Streaming поток это бесконечная таблица, а запросы пишутся как над DataFrame; trigger задаёт темп микробатчей. Чекпоинт в надёжном сторадже - единственная память джобы между запусками.
// Выбор по задаче: нужна миллисекундная латентность и сложное состояние - Flink; уже есть Spark и хватает микробатчей - Structured Streaming; лёгкая обработка внутри сервиса - Kafka Streams без кластера.
- micro-batch
- обработка потока маленькими батчами (Spark SS)
- trigger
- темп микробатчей в Structured Streaming
Состояние и чекпоинты
Состояние - агрегаты, окна, джойны - живёт в state store (RocksDB-класс) и чекпоинтится. Восстановление после сбоя это состояние плюс offset'ы, поднятые из чекпоинта; потерять или удалить чекпоинт значит обнулить и агрегаты, и позицию чтения.
Stream-stream join требует буферизации обеих сторон с ограничением по времени (watermark плюс интервал), иначе состояние растёт бесконечно до OOM (out of memory).
// Stream-static join (поток на справочник) дёшев; если справочник обновляется, его подают как compacted topic-changelog или периодически перезагружают.
- state store
- локальное состояние оператора с чекпоинтами
- checkpoint
- снапшот состояния и offset'ов для восстановления
Режимы вывода и replay
Выходные режимы согласуют с sink'ом: append (только финализированные строки), update (обновления по ключу), complete (вся таблица). Append для обновляемых окон даёт дубли строк.
Эволюция пайплайна опасна: смена логики агрегации ломает совместимость состояния, поэтому запуск нового кода со старым чекпоинтом падает. Логику версионируют, а пересчёт из сырого топика (replay) планируют как штатную операцию.
// Replay - суперсила лога: новый пайплайн или исправление бага прогоняется по всей истории топика. Но это работает, только если retention сырых событий выбран с расчётом на пересчёты - retention в 3 дня при цикле багфиксов в неделю делает replay невозможным.
- output mode
- append / update / complete - что выдаём в sink
- replay
- повторная обработка истории лога новым кодом
Как отвечать: «Что такое чекпоинт и replay в стриминге?»
Чекпоинт это снапшот памяти стриминговой джобы: её состояния (агрегаты, окна, содержимое джойнов) и позиций чтения, offset'ов, сохранённый в надёжный сторадж. Он нужен, потому что stateful-обработка держит состояние локально, и без чекпоинта рестарт означал бы пересчёт агрегатов с нуля и потерю позиции. При сбое джоба поднимает состояние и offset'ы из последнего чекпоинта и продолжает, будто ничего не было. Поэтому потерять чекпоинт - катастрофа, а сменить схему агрегации при живом старом чекпоинте нельзя, состояние несовместимо. Replay это другая сила, идущая от лога: раз события хранятся в топике, я могу прогнать по всей истории новый пайплайн или исправленную версию, пересчитав результат с нуля. Это штатный способ выкатывать багфиксы и новые метрики. Но replay возможен, только если retention сырых событий выбран с запасом на такие пересчёты.
Почему это сильный ответ: чекпоинт объяснён как память джобы (состояние + offset) с последствиями потери и несовместимости, replay - как пересчёт из лога с ключевым условием retention; связаны две разные механики.
На чём валят
- −Удалить или потерять чекпоинт - состояние и позиция чтения обнулились.
- −Stream-stream join без временных ограничений - state растёт до OOM.
- −Поменять схему агрегации и запуститься со старым чекпоинтом - несовместимость состояния.
- −Append-режим для обновляемых окон - дубли строк окна в sink.
- −Retention сырья 3 дня при цикле исправления багов в неделю - replay невозможен.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 11, остальные разбираются в тренажёре.
- Зачем в стриминговой архитектуре используют Schema Registry?A)Централизованно хранить и версионировать схемы сообщений, проверяя совместимость, чтобы продюсер не сломал консьюмеровB)Schema Registry хранит сами сообщения топиков, выступая заменой брокеру Kafka для их персистентностиC)Реестр схем нужен лишь для красивой документации и на совместимость форматов никак не влияетD)Он позволяет продюсеру менять формат как угодно, автоматически подстраивая всех консьюмеров под него
показать ответ и разбор
+A)Централизованно хранить и версионировать схемы сообщений, проверяя совместимость, чтобы продюсер не сломал консьюмеров// разбор: В потоке продюсеры и консьюмеры развязаны во времени, и изменение формата сообщения легко ломает читателей. Schema Registry хранит версии схем (Avro/Protobuf/JSON Schema) и при регистрации новой проверяет совместимость (backward/forward): нельзя выкатить продюсера, чей формат не прочитают существующие консьюмеры. В сообщении едет id схемы, а не вся схема — компактно. Это контракт данных для эволюции формата без простоев.
- Что такое CDC (Change Data Capture) и зачем он нужен дата-инженеру?A)CDC — это периодический полный SELECT всей таблицы источника по расписанию раз в суткиB)Поток изменений из WAL/binlog без полного перечитывания источникаC)CDC захватывает операции INSERT, игнорируя изменения и удаления строкD)CDC работает внутри одной БД и не может передавать изменения во внешние системы
показать ответ и разбор
+B)Поток изменений из WAL/binlog без полного перечитывания источника// разбор: CDC читает журнал транзакций БД-источника (WAL/binlog) и превращает каждое изменение строки в событие в потоке (часто через Debezium в Kafka). Это даёт near-real-time репликацию и синхронизацию в DWH/озеро без тяжёлых периодических full-scan выгрузок и без нагрузки на источник запросами. Ключевые заботы — начальный снапшот, порядок, обработка update/delete и схемной эволюции. Основа современных near-real-time пайплайнов из OLTP.
- Стрим-джоб перегружен: источник шлёт быстрее, чем джоб успевает. Что такое backpressure и как система реагирует?A)Backpressure — это когда джоб просто молча отбрасывает лишние события, которые не успевает обработатьB)Backpressure заставляет источник слать данные ещё быстрее, чтобы очередь разгрузилась скорееC)Backpressure — сигнал перегрузки, замедляющий чтение из источника до темпа обработки, чтобы не переполнить память/очередьD)Это ошибка конфигурации, backpressure не должен возникать в правильно настроенной системе
показать ответ и разбор
+C)Backpressure — сигнал перегрузки, замедляющий чтение из источника до темпа обработки, чтобы не переполнить память/очередь// разбор: Backpressure (обратное давление) — механизм, которым перегруженный оператор сообщает вверх по конвейеру «читай медленнее». В pull-моделях (Kafka + Flink/Spark) обработчик просто вытягивает новые записи по мере готовности, поэтому темп чтения естественно ограничивается скоростью обработки, а лаг растёт в надёжном логе Kafka, а не в памяти джоба. Без backpressure перегруженный джоб копил бы данные в памяти до OOM. Устойчивый разрыв лечат масштабом, не буфером.
- Зачем брать стрим-фреймворк (Flink, Spark Structured Streaming) вместо голого консьюмера Kafka?A)Стрим-фреймворк работает быстрее голого консьюмера Kafka на нагрузкеB)Фреймворк заменяет собой брокер Kafka, так что сам Kafka в архитектуре не нуженC)Голый консьюмер не способен прочитать ни одного сообщения из топика без стрим-фреймворкаD)Он даёт из коробки состояние, окна, чекпоинты и exactly-once — не переписывать это руками
показать ответ и разбор
+D)Он даёт из коробки состояние, окна, чекпоинты и exactly-once — не переписывать это руками// разбор: Голый Kafka-консьюмер читает сообщения, но всё остальное — состояние агрегатов, окна по времени, watermark'и, чекпоинты для восстановления, exactly-once — писать и отлаживать самому. Стрим-фреймворк (Flink, Spark Structured Streaming, Kafka Streams) даёт это встроенными примитивами: объявляешь оконную агрегацию и семантику доставки, а фреймворк берёт на себя надёжное состояние и восстановление после сбоя. Голый консьюмер оправдан для простой stateless-обработки (фильтр, роутинг).
- Джоб читает поток событий и отбрасывает те, где сумма меньше ста. Какая это обработка?A)stateless: решение принимается по одному событиюB)stateful: джоб помнит все прошедшие событияC)оконная: события группируются по интервалам времениD)пакетная: события накапливаются перед обработкой
показать ответ и разбор
+A)stateless: решение принимается по одному событию// разбор: Фильтрация смотрит на событие и больше ни на что: чтобы решить его судьбу, прошлое не нужно. Такой джоб легко масштабируется и переживает перезапуск без хлопот. Состояние появляется, когда результат зависит от предыдущих событий — счётчик по ключу, дедупликация, соединение двух потоков. Вот тогда и нужны чекпоинты, иначе перезапуск потеряет накопленное.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.