В чем разница между micro-batching и streaming?

«В чем разница между micro-batching и streaming?» — вопрос из категории Потоковая обработка и Kafka, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Micro-batching — это подход, при котором данные накапливаются небольшими порциями (батчами) в течение короткого фиксированного интервала (например, 500 мс или 1 секунда), а затем обрабатываются как единый блок. Streaming — это обработка каждого события или записи по мере их поступления, практически без задержки.

Ключевые различия на примере Apache Kafka и фреймворков обработки:

Аспект Micro-batching (например, Spark Structured Streaming) True Streaming (например, Apache Flink, Kafka Streams)
Задержка (Latency) Выше (сотни миллисекунд - секунды), так как данные ждут закрытия окна батча. Минимальная (миллисекунды), события обрабатываются сразу.
Модель обработки Дискретная, по тактам. Приложение «просыпается» для обработки накопленного батча. Непрерывная, операторы обрабатывают поток событий постоянно.
Семантика доставки Проще обеспечить exactly-once благодаря детерминированным интервалам и checkpoint'ам на уровне батча. Требует более сложных механизмов (водяные знаки, распределенные снапшоты) для гарантии exactly-once.
Пропускная способность Может быть выше за счет оптимизации обработки групп записей и снижения накладных расходов. Нагрузка более равномерная, но overhead на обработку каждого события может быть выше.

Пример концептуального кода для Kafka:

  • Micro-batching (Spark):

    val df = spark.readStream
      .format("kafka")
      .option("subscribe", "topic")
      .load() // Данные читаются микробатчами
  • Streaming (Kafka Streams):

    KStream<String, String> stream = builder.stream("input-topic");
    stream.mapValues(value -> value.toUpperCase()) // Преобразование применяется к каждому событию сразу
          .to("output-topic");

Выбор подхода зависит от требований: streaming для низкой задержки (мониторинг, алертинг), micro-batching для высокой пропускной способности при приемлемой задержке (ETL, агрегации).