Открыть сервис

Apache Spark Streaming

Apache Spark Streaming — это компонент экосистемы Apache Spark, предназначенный для обработки потоковых данных (streaming data) в реальном времени или близком к реальному времени. Он представляет собой высокоуровневый API для построения масштабируемых, отказоустойчивых и распределённых приложений, работающих с непрерывно поступающими данными. Spark Streaming входит в состав платформы Apache Spark, разработанной в Калифорнийском университете в Беркли (AMPLab) и впоследствии переданной в Apache Software Foundation.

История и развитие

Первая версия Spark Streaming была представлена в 2013 году как часть Apache Spark 0.8. Изначально она базировалась на модели дискретизированных потоков (DStreams — Discretized Streams), где поток данных разбивался на небольшие пакеты (микропакеты) фиксированной длины, которые затем обрабатывались как микро-пакетные задания (micro-batches). Этот подход обеспечивал высокую пропускную способность и отказоустойчивость за счёт использования механизма линеаризации (lineage) Spark RDD.

В 2016 году с выходом Apache Spark 2.0 был представлен Structured Streaming — более высокоуровневый и декларативный API, построенный поверх Spark SQL. Structured Streaming позволял обрабатывать потоковые данные как бесконечные таблицы (unbounded tables), используя знакомые операции DataFrame/Dataset. Это значительно упростило разработку потоковых приложений и устранило многие ограничения модели DStreams.

К 2018 году (Spark 2.4) Structured Streaming стал основным рекомендуемым способом обработки потоков в Spark, а DStreams были переведены в режим поддержки (maintenance mode). В последующих версиях (Spark 3.x) развитие сосредоточено на Structured Streaming, включая поддержку событийно-управляемой обработки (event-time processing), водяных знаков (watermarks) и улучшенную интеграцию с Kafka.

Архитектура и ключевые концепции

Модель обработки: микропакеты (Micro-Batch)

Основная архитектурная особенность Spark Streaming — использование микропакетной обработки. Входящий поток данных разбивается на небольшие интервалы (batch interval), обычно от 500 миллисекунд до нескольких секунд. Каждый такой интервал формирует один микропакет, который обрабатывается как обычное задание Spark (Spark job) на кластере. Это обеспечивает:

  • Отказоустойчивость: при сбое узла Spark может перевычислить потерянные микропакеты на основе их lineage.
  • Масштабируемость: микропакеты распределяются по узлам кластера, как и обычные RDD.
  • Точность-один-раз (exactly-once semantics): при правильной настройке источников (например, Kafka) и приёмников (sinks) гарантируется, что каждое событие будет обработано ровно один раз.

Structured Streaming API

Structured Streaming использует декларативную модель: пользователь определяет, что нужно сделать с данными (фильтрация, агрегация, join), а Spark оптимизирует выполнение. Ключевые элементы:

  • Бесконечная таблица (Unbounded Table): поток данных рассматривается как таблица, которая непрерывно пополняется новыми строками.
  • Результат (Result Table): результат обработки (агрегации, фильтрации) также представляется как таблица, которая обновляется с каждым новым микропакетом.
  • Режимы вывода (Output Modes):
  • Append: новые строки добавляются в результат (подходит для фильтрации или простых трансформаций).
  • Update: изменяются существующие строки (например, при агрегации с группировкой).

Complete: весь результат перезаписывается каждый раз (подходит для агрегаций без группировки или с фиксированным набором ключей).

  • Водяные знаки (Watermarks): механизм для обработки событий с задержкой (late data). Водяной знак определяет временную границу, после которой данные считаются слишком старыми и отбрасываются. Это необходимо для корректной работы оконных агрегаций (windowed aggregations) и предотвращения бесконечного накопления состояния.

Обработка времени событий (Event-Time Processing)

В отличие от времени обработки (processing time — когда данные поступили в Spark), Spark Streaming поддерживает обработку на основе времени события (event time — когда событие произошло на источнике). Это критично для приложений, где важна хронология (например, логирование, IoT). Для этого используется параметр withWatermark и оконные функции (например, window("timestamp", "10 minutes")).

Источники и приёмники данных

Источники (Sources)

Spark Streaming поддерживает множество источников данных. Основные:

  • Apache Kafka: наиболее популярный источник. Обеспечивает высокую пропускную способность и строгую гарантию доставки. Используется интеграция через readStream().format("kafka").
  • Файловые системы: чтение из HDFS, S3, Azure Blob Storage, локальной файловой системы. Поддерживаются форматы Parquet, ORC, JSON, CSV, Avro.
  • Сокеты (Socket): для тестирования и разработки (TCP-сокеты).
  • Delta Lake: чтение из таблиц Delta Lake (формат хранения, обеспечивающий ACID-транзакции).
  • Amazon Kinesis, Azure Event Hubs: через отдельные библиотеки.

Приёмники (Sinks)

Вывод результатов может производиться в:

  • Kafka: запись обратно в топики Kafka.
  • Файловые системы: запись в Parquet, ORC, JSON, CSV.
  • Консоль: для отладки (console sink).
  • Memory: для тестирования (memory sink).
  • Foreach/foreachBatch: пользовательский код для записи в любую внешнюю систему (базы данных, API, объектные хранилища). foreachBatch позволяет выполнять пакетную обработку каждого микропакета, что даёт доступ к полному набору Spark API.
  • Delta Lake: запись в таблицы Delta Lake с поддержкой транзакций.

Применение

Apache Spark Streaming широко используется в областях, требующих обработки непрерывных потоков данных:

  • Мониторинг и аналитика в реальном времени: агрегация метрик с серверов, датчиков, систем логирования (ELK-стек, Prometheus + Grafana).
  • Обработка событий: выявление аномалий, мошенничества (fraud detection) в финансовых транзакциях, обработка кликов на веб-сайтах.
  • IoT (Интернет вещей): сбор и обработка данных с миллионов устройств (температура, вибрация, GPS-координаты).
  • ETL (Extract, Transform, Load): непрерывная загрузка данных из источников (Kafka, файлы) в хранилища данных (Data Lakes, Data Warehouses) с трансформацией.
  • Машинное обучение: обучение моделей на потоковых данных (online learning) или применение обученных моделей для классификации/предсказаний в реальном времени (streaming inference). Spark Streaming интегрируется с MLlib.

Сравнение с другими технологиями

Spark Streaming занимает нишу между чисто потоковыми системами (Apache Flink, Apache Storm) и пакетными системами (Hadoop MapReduce). Его преимущества:

  • Единая платформа: возможность совмещать потоковую и пакетную обработку в одном кластере и одном коде (через общий API DataFrame/Dataset).
  • Простота: декларативный API Structured Streaming значительно проще для разработки, чем низкоуровневые API Flink или Storm.
  • Экосистема: интеграция с Spark SQL, MLlib, GraphX.

Недостатки:

  • Задержка (Latency): из-за микропакетной обработки минимальная задержка составляет несколько сотен миллисекунд (обычно 500 мс — 2 секунды). Для приложений, требующих миллисекундной задержки (например, высокочастотная торговля), Flink или Storm более подходят.
  • Управление состоянием: хотя Structured Streaming поддерживает управление состоянием (state management), оно менее гибкое, чем в Flink (например, нет поддержки сложных состояний с пользовательскими таймерами).

Критика и ограничения

Основные критические замечания в адрес Spark Streaming:

  • Задержка: микропакетная модель принципиально не может обеспечить задержку меньше batch interval. Для некоторых сценариев (например, оповещения о критических событиях) это неприемлемо.
  • Сложность настройки: для достижения высокой производительности требуется тонкая настройка параметров (размер микропакета, параллелизм, управление памятью).
  • Ограниченная поддержка событийно-управляемой обработки: хотя Structured Streaming поддерживает event-time, реализация водяных знаков и оконных функций менее гибкая, чем в Flink.
  • Потребление ресурсов: микропакетная обработка может приводить к избыточному потреблению памяти и CPU при большом количестве мелких микропакетов.

Интересные факты

  • Spark Streaming изначально был создан для обработки логов веб-серверов в Facebook (продукт Meta, признанной экстремистской и запрещённой в РФ), но затем был адаптирован для общего применения.
  • В 2014 году Spark Streaming установил рекорд по скорости обработки данных на конкурсе Sort Benchmark, обработав 100 ТБ данных за 23 минуты.
  • Structured Streaming использует оптимизатор Catalyst (тот же, что и Spark SQL) для автоматической оптимизации планов выполнения потоковых запросов.

Источники

  • Apache Spark Documentation: Structured Streaming Programming Guide (Spark 3.5.0)
  • Matei Zaharia et al. "Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing" (2012)
  • Michael Armbrust et al. "Structured Streaming: A Declarative API for Real-Time Applications in Apache Spark" (2018)
  • "Learning Spark", 2nd Edition, O'Reilly Media (2020)
  • "Spark: The Definitive Guide", Bill Chambers, Matei Zaharia, O'Reilly Media (2018)

BFOmetr — база данных и аналитика по компаниям России.

На главную BFOmetr →