сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Kafka и стриминг

Потоковая обработка данных

Движки обработки потока

Поверх лога работают движки - 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, остальные разбираются в тренажёре.

  1. #stream_processing1 / 5
    Зачем в стриминговой архитектуре используют Schema Registry?
    A)Централизованно хранить и версионировать схемы сообщений, проверяя совместимость, чтобы продюсер не сломал консьюмеров
    B)Schema Registry хранит сами сообщения топиков, выступая заменой брокеру Kafka для их персистентности
    C)Реестр схем нужен лишь для красивой документации и на совместимость форматов никак не влияет
    D)Он позволяет продюсеру менять формат как угодно, автоматически подстраивая всех консьюмеров под него
    показать ответ и разбор
    +A)Централизованно хранить и версионировать схемы сообщений, проверяя совместимость, чтобы продюсер не сломал консьюмеров

    // разбор: В потоке продюсеры и консьюмеры развязаны во времени, и изменение формата сообщения легко ломает читателей. Schema Registry хранит версии схем (Avro/Protobuf/JSON Schema) и при регистрации новой проверяет совместимость (backward/forward): нельзя выкатить продюсера, чей формат не прочитают существующие консьюмеры. В сообщении едет id схемы, а не вся схема — компактно. Это контракт данных для эволюции формата без простоев.

  2. #stream_processing2 / 5
    Что такое 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.

  3. #stream_processing3 / 5
    Стрим-джоб перегружен: источник шлёт быстрее, чем джоб успевает. Что такое backpressure и как система реагирует?
    A)Backpressure — это когда джоб просто молча отбрасывает лишние события, которые не успевает обработать
    B)Backpressure заставляет источник слать данные ещё быстрее, чтобы очередь разгрузилась скорее
    C)Backpressure — сигнал перегрузки, замедляющий чтение из источника до темпа обработки, чтобы не переполнить память/очередь
    D)Это ошибка конфигурации, backpressure не должен возникать в правильно настроенной системе
    показать ответ и разбор
    +C)Backpressure — сигнал перегрузки, замедляющий чтение из источника до темпа обработки, чтобы не переполнить память/очередь

    // разбор: Backpressure (обратное давление) — механизм, которым перегруженный оператор сообщает вверх по конвейеру «читай медленнее». В pull-моделях (Kafka + Flink/Spark) обработчик просто вытягивает новые записи по мере готовности, поэтому темп чтения естественно ограничивается скоростью обработки, а лаг растёт в надёжном логе Kafka, а не в памяти джоба. Без backpressure перегруженный джоб копил бы данные в памяти до OOM. Устойчивый разрыв лечат масштабом, не буфером.

  4. #stream_processing4 / 5
    Зачем брать стрим-фреймворк (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-обработки (фильтр, роутинг).

  5. #stream_processing5 / 5
    Джоб читает поток событий и отбрасывает те, где сумма меньше ста. Какая это обработка?
    A)stateless: решение принимается по одному событию
    B)stateful: джоб помнит все прошедшие события
    C)оконная: события группируются по интервалам времени
    D)пакетная: события накапливаются перед обработкой
    показать ответ и разбор
    +A)stateless: решение принимается по одному событию

    // разбор: Фильтрация смотрит на событие и больше ни на что: чтобы решить его судьбу, прошлое не нужно. Такой джоб легко масштабируется и переживает перезапуск без хлопот. Состояние появляется, когда результат зависит от предыдущих событий — счётчик по ключу, дедупликация, соединение двух потоков. Вот тогда и нужны чекпоинты, иначе перезапуск потеряет накопленное.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.