Что такое Apache Spark и где его используют
Apache Spark — это open-source фреймворк для распределенной обработки больших данных. Его используют для batch-обработки, потоковой аналитики, SQL-запросов, машинного обучения и обработки данных в кластерах Hadoop, Kubernetes, облаках и автономных окружениях.
Spark стал популярным благодаря скорости, работе с данными в памяти, удобным API и поддержке нескольких языков: Scala, Java, Python и R. В статье разберем, что такое Apache Spark, из каких компонентов он состоит, чем Spark Core отличается от Spark SQL и Spark Streaming, зачем нужны MLlib и GraphX, а также в каких задачах Spark действительно полезен.
Что внутри:
- что такое Apache Spark простыми словами;
- зачем нужен фреймворк для распределенной обработки данных;
- основные преимущества Spark;
- Spark Core и управление вычислениями;
- Spark SQL и DataFrame;
- Spark Streaming и обработка потоковых данных;
- MLlib для машинного обучения;
- GraphX для графовых вычислений;
- поддерживаемые языки: Scala, Java, Python и R;
- где запускают Spark: Hadoop YARN, Kubernetes, облака и standalone;
- примеры применения Spark в аналитике и data engineering.
- Производительность
- Удобный интерфейс
- Отказоустойчивость
- Поддержка четырех языков: Scala, Java, Python и R
Компоненты архитектуры Apache Spark
1. Ядро Apache Spark (Spark Core) — основа фреймворка. Это базовый движок для параллельной и распределенной обработки данных. Ядро отвечает за:
- Управление памятью и восстановление системы после отказов
- Планирование, распределение и отслеживание заданий в кластере
- Взаимодействие с хранилищем данных
2. Apache Spark SQL — одна из четырех библиотек фреймворка, которая использует структуру данных DataFrames и может выступать в роли распределенного механизма запросов SQL. Библиотека возникла как порт Apache Hive для работы поверх Apache Spark (вместо MapReduce), а сейчас уже интегрирована со стеком Spark. Apache Spark SQL обеспечивает поддержку различных источников данных и позволяет переплетать SQL-запросы с трансформациями кода.
3. Apache Spark Streaming — инструмент для обработки потоковых данных, который легко интегрируется с широким спектром популярных источников данных: HDFS, Flume, Kafka, ZeroMQ, Kinesis и Twitter. Apache Spark Streaming обрабатывает данные в реальном в режиме micro-batch, минимальное время обработки каждого micro-batch 0,5 секунды. Apache Spark Streaming получает входные потоки данных и разбивает данные на пакеты. Далее они обрабатываются движком Apache Spark, после чего генерируется конечный поток данных также в пакетной форме. API Spark Streaming точно соответствует API Spark Core, поэтому можно одновременно работать как с пакетными, так и с потоковыми данными.
- Банки и страховые компании (прогноз востребованности услуг)
- Поисковики и социальные сети (выявление фейковых аккаунтов, оптимизация таргетинга и т.д.)
- Службы такси (анализ времени и геолокаций, прогноз спроса и цен)
- Транспортные и авиакомпании (модели для прогнозирования задержек рейсов)

Spark предоставляет быструю и универсальную платформу для обработки данных. По сравнению с Hadoop Spark ускоряет работу программ в памяти более чем в 100 раз, а на диске – более чем в 10 раз. Spark дает больше возможностей для работы с данными. Его синтаксис не так сложен, чтобы начать погружение, для сравнения приведу пример из Pandas.
Для работы с Spark, нужно создать сессию.
``` spark = SparkSession.builder.getOrCreate() ```
Во время создания сессии, происходит кластеризация.
Pandas
```
data = pd.read_csv('data.csv')
````
Spark
````
data = spark.read.csv(path=’data.csv’, header=True, sep=’,’)
````Далее, сгруппируем данные и «сместим» в колонке на одну позицию. В Pandas это делается так:
``` data[group1] = pandas_df.groupby(group2)[group3].shift(-1) ```
В Spark
```
w = Window().partitionBy("group2").orderBy("group3")
data = data.withColumn("group2", lag("group2", -1, 0).over(w))
```Можно использовать оконную функцию, где partitionBy отвечает за группировку данных, а orderBy сортировка. Функция lag принимает 3 параметра: это колонка, шаг смещения и значения, которые будет на месте шага. Или для группировки можно использовать обычную функцию groupBy, которая тоже есть в Spark. Разница в том, что с окном каждая строка будет связана с результатом агрегирования, вычисленным для всего окна. Однако при группировке каждая группа будет связана с результатом агрегации в этой группе (группа строк становится только одной строкой).
```
dataframe = spark.range(6).withColumn("key", 'id % 2)
dataframe.show
windowing = Window.partitionBy("key")
dataframe.withColumn("sum", sum(col("id")).over(windowing).show
dataframe.groupBy("key").agg(sum('id)).show
К сожалению, некоторых функций может не быть в Spark (например, factorize).
```
labels_start, uniques = pd.factorize(anomaly_time['activity_start']) anomaly_time['activity_start_code'] = labels_start
```
Spark
````
win_func = Window().partitionBy().orderBy(lit(' '))
data = data.select('name_column').distinct().withColumn('name_column', row_number().over(win_func) - 1)
````Функция factorize закодирует объект как перечислимый тип или категориальную переменную, или присвоит объекту идентификатор.
``` codes, uniques = pd.factorize(['b', 'b', 'a', 'c', 'b']) codes array([0, 0, 1, 2, 0]...) ```
Для выполнения подобного функционала в Spark, берется колонка select (‘name_column’) и выбираются все уникальные значения, с помощью функции distinct. Далее с помощью функции withColumn создается колонка и присваивается номер строки (чтобы начиналось с 0 — я отнимаю 1).
Вывод
Apache Spark это огромная система, с множеством инструментов для разных типов задач от SQL до машинного обучения. В этой статье был показан лишь маленький кусочек от всего Spark, но даже этого хватит, чтобы начать обрабатывать данные.
FAQ
Вопрос: Что такое Apache Spark?
Ответ: Apache Spark — это фреймворк с открытым исходным кодом для распределенной обработки больших данных. Он позволяет выполнять вычисления на кластере и работать с batch-данными, потоками, SQL и машинным обучением.
Вопрос: Для чего используют Spark?
Ответ: Spark используют для ETL, обработки больших данных, аналитики, потоковой обработки, машинного обучения, подготовки витрин данных, расчета признаков и интеграции данных из разных источников.
Вопрос: Что такое Spark SQL?
Ответ: Spark SQL — компонент Apache Spark для работы со структурированными данными через SQL и DataFrame API. Он позволяет выполнять распределенные SQL-запросы и совмещать их с кодом на Python, Scala, Java или R.
Вопрос: Чем Spark Streaming отличается от обычной batch-обработки?
Ответ: Batch-обработка работает с накопленными наборами данных, а Spark Streaming обрабатывает входящие потоки небольшими пакетами. Это позволяет строить почти real-time пайплайны данных.
Вопрос: Что такое MLlib?
Ответ: MLlib — библиотека машинного обучения в Apache Spark. В ней есть алгоритмы классификации, регрессии, кластеризации, рекомендаций и инструменты подготовки признаков.
- Apache Airflow DAGs: основы и расписания
- ETL и ELT: 5 основных отличий
- DWH: зачем компании хранилище данных
- Что такое Parquet: преимущества и случаи использования
- Алгоритмы кластеризации в машинном обучении



