Apache Spark
Что это
Apache Spark — фреймворк с открытым исходным кодом для реализации распределённой обработки неструктурированных
и слабоструктурированных данных через параллельные вычисления на кластере, входящий в экосистему проектов
Hadoop.
Spark состоит:
- Spark Core - базовый движок. Управление памятью, Планирование задач, Отказоустойчивость, Работа с RDD.
- SQL. Работа со структурированными данными, DataFrame и Dataset API, SQL-запросы.
- STREAMING DATA - обработка потоковых данных в реальном времени.
- MACHINE LEARNING - Библиотека машинного обучения, классификация, регрессия, кластеризация и т.д.
- GRAPH ANALYTICS - Обработка графов и сетей связей
С архитектурной точки зрения Spark-приложение состоит из:
- Driver — запускает приложение, строит план выполнения.
- Cluster Manager — распределяет ресурсы(Standalone, YARN, Kubernetes и др.).
- Executors — процессы на рабочих узлах, которые выполняют задачи.
- Worker Nodes — серверы, на которых работают executors.
Типы данных
| Характеристика |
RDD |
DataFrame |
Dataset |
| Типизация |
✔ |
✖ |
✔ |
| Оптимизация(Catalys) |
✖ |
✔ |
✔ |
| Производительность |
Ниже |
Высокая |
Высокая |
| Удобство SQL |
✖ |
✔ |
✔ |
| Контроль над данными |
Максимальный |
Средний |
Высокий |
RDD(Resilient Distributed Datasets)
RDD - основной базовая абстракция Spark, представляющая неизменный набор элементов, разделенных по узлам
кластера, что позволяет выполнять параллельные вычисления.
Его можно получить из:
- Файла
- Памяти
- Другого RDD
SparkContext или JavaSparkContext(для JAVA) - место откуда мы получает RDD.
Пример конфигурации JavaSparkContext:
SparkConf conf = new SparkConf();
conf.setAppName("my spark application");
conf.setAppName("local[*]");
JavaSparkContext sc = new JavaSparkContext(conf);
DataFrame
DataFrame - состоит из RDD. По сути это таблица со схемой, аналог SQL-таблицы.
Datasets
Datasets - это типизированные DataFrame. Комбинация RDD и DataFrame.
RDD Operations:
Имеет методы, которые работают аналогично с Stream API(Transformations - Intermediate, Actions - Terminal).
Также как в JAVA Stream API если метод возвращает RDD то это Transformations иначе это Actions.
Transformations:
- map
- flatMap
- filter
- mapPartitions, mapPartitionsWithIndex
- sample
- union, intersection, join, cogroup, cartesian
- distinct
- reduceByKey, aggregateByKey, sortByKey
- pipe
- coalesce, repartition, repartitionAndSortWithinPartitions
Actions:
- reduce
- collect
- count, countByKey, countByValue
- first
- take, takeSample, takeOrdered
- saveAsTextFile, saveAsSequenceFile, saveAsObjectFile
- foreach
На Spark все Transformations методы выполняются на кластере и Actions на драйвере(тот кто отвечает за
таксу и распределяет ее по кластеру).
Мы можем сохранять промежуточные операции методом Transformations persist(StorageLevel).
DataFrame Operations:
- show() - отображение контента DataFrame.
show(boolean) - truncate.
show(int) - number of rows.
Стратегии выполнения Join в Spark
- Broadcast Hash Join - маленькая таблица отправляется на все executors.
- Shuffle Hash Join - обе таблицы перераспределяются по сети (shuffle), затем строится hash-таблица.
- Sort Merge Join - наиболее распространен для больших таблиц. Shuffle -> Sort -> Merge.
- Broadcast Nested Loop Join - используется для некоторых неравенств.
Shuffle
Spark shuffle - это операция перемешивания(перераспределения) данных между Executors. Побочный эффект
таких аналитических преобразований, как join(), groupBy(), orderBy(), reduceByKey(), union() и тд.
Способы снижения отрицательного эффекта Shuffle-операций
- Сократить число Executors(меньше обмена между машинами)
- Уменьши объем данных(сбрасывания лишних столбцов, фильтрация)
Catalyst Optimizer
Catalyst Optimizer - оптимизатор для DataSet.
Примеры кода
Проверка колонки на максимальное значение(Date или Numeric)
.groupBy(col(AF.OfficeId)).agg(min(AF.Date))