Что такое потоковая обработка данных (stream processing)?

«Что такое потоковая обработка данных (stream processing)?» — вопрос из категории Потоковая обработка и Kafka, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Потоковая обработка данных — это парадигма обработки, при которой данные обрабатываются непрерывно, по мере их поступления из источников (потоков), в режиме, близком к реальному времени. В отличие от пакетной обработки (batch), которая работает с конечными, статичными наборами данных, потоковые системы имеют дело с теоретически бесконечными потоками событий.

Ключевые концепции:

  • Неограниченный поток данных: Последовательность событий, которые поступают постоянно (логи приложений, телеметрия IoT, клики пользователей, финансовые транзакции).
  • Низкая задержка: Цель — обработать событие за миллисекунды или секунды с момента его генерации.
  • Обработка событий (Event-driven): Приложение реагирует на каждое новое событие или микропакеты событий.

Архитектурный паттерн и пример с Apache Kafka и Kafka Streams: Популярная архитектура — использование Apache Kafka как надежного, распределенного брокера сообщений, который буферизует потоки событий в топиках. Обработка может происходить с помощью фреймворков вроде Kafka Streams, Apache Flink или Spark Structured Streaming.

// Пример простого приложения Kafka Streams для агрегации в реальном времени
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

public class ClickstreamCounter {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "clickstream-aggregator");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

        StreamsBuilder builder = new StreamsBuilder();
        // Источник: чтение потока кликов из Kafka-топика
        KStream<String, String> clickStream = builder.stream("user-clicks");

        // Обработка: подсчет кликов по пользователю за окно времени
        KTable<Windowed<String>, Long> clicksPerUser = clickStream
            .groupBy((key, value) -> extractUserId(value)) // Группировка по user_id
            .windowedBy(TimeWindows.of(Duration.ofMinutes(5))) // Скользящее окно 5 минут
            .count(); // Агрегация - подсчет событий

        // Приемник: отправка результатов в другой топик
        clicksPerUser.toStream()
            .map((windowedKey, count) -> new KeyValue<>(windowedKey.key(), count.toString()))
            .to("user-clicks-per-5min");

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();
    }
}

Типичные сценарии использования:

  • Мониторинг и алертинг в реальном времени: Обнаружение аномалий в метриках серверов или падение количества успешных транзакций.
  • Обогащение данных: Добавление контекста к событию (например, информации о пользователе) перед сохранением в хранилище.
  • Построение актуальных дашбордов: Агрегация данных для отображения текущих KPI.
  • Обработка конвейеров событий (Event Sourcing/CQRS): Обновление материализованных представлений на основе потока событий.

Основные вызовы: Гарантии доставки (exactly-once, at-least-once), управление состоянием приложения, обработка поздних данных (late events) и обеспечение отказоустойчивости.