Apache Spark: гайд для новичков
Специалисты компании Databricks, основанной создателями Spark, собрали лучшее о функционале Apache Spark в своей книге Gentle Intro to Apache Spark (очень рекомендую прочитать):
“Apache Spark — это целостная вычислительная система с набором библиотек для параллельной обработки данных на кластерах компьютеров. На данный момент Spark считается самым активно разрабатываемым средством с открытым кодом для решения подобных задач, что позволяет ему быть полезным инструментом для любого разработчика или исследователя-специалиста, заинтересованного в больших данных. Spark поддерживает множество широко используемых языков программирования (Python, Java, Scala и R), а также библиотеки для различных задач, начиная от SQL и заканчивая стримингом и машинным обучением, а запустить его можно как с ноутбука, так и с кластера, состоящего из тысячи серверов. Благодаря этому Apache Spark и является удобной системой для начала самостоятельной работы, перетекающей в обработку больших данных в невероятно огромных масштабах.”
Что такое большие данные?
Посмотрим-ка на популярное определение больших данных по Гартнеру. Это поможет разобраться в том, как Spark способен решить множество интересных задач, которые связаны с работой с большими данными в реальном времени:
“Большие данные — это информационные активы, которые характеризуются большим объёмом, высокой скоростью и/или многообразием, а также требуют экономически эффективных инновационных форм обработки информации, что приводит к усиленному пониманию, улучшению принятия решений и автоматизации процессов.”
Сложный мир больших данных
Заметка: Ключевой вывод — слово “большие” в больших данных относится не только к объёму. Вы не просто получаете много данных, они поступают в реальном времени очень быстро и в различных комплексных форматах, а ещё — из большого многообразия источников. Вот откуда появились 3-V больших данных: Volume (Объём), Velocity (Скорость), Variety (Многообразие).
Причины использовать Spark
Основываясь на самостоятельном предварительном исследовании этого вопроса, я пришёл к выводу, что у Apache Spark есть три главных компонента, которые делают его лидером в эффективной работе с большими данными, а это мотивирует многие крупные компании работать с большими наборами неструктурированных данных, чтобы Apache Spark входил в их технологический стек.
- Spark — всё-в-одном для работы с большими данными. “Spark создан для того, чтобы помогать решать широкий круг задач по анализу данных, начиная с простой загрузки данных и SQL-запросов и заканчивая машинным обучением и потоковыми вычислениями, при помощи одного и того же вычислительного инструмента с неизменным набором API. Главный инсайт этой программной многозадачности в том, что задачи по анализу данных в реальном мире — будь они интерактивной аналитикой в таком инструменте, как Jupyter Notebook, или же обычным программированием для выпуска приложений — имеют тенденцию требовать сочетания множества разных типов обработки и библиотек. Целостная природа Spark делает решение этих заданий проще и эффективнее.” (Из книги Databricks). Например, если вы загружаете данные при помощи SQL-запроса и потом оцениваете модель машинного обучения при помощи библиотеки Spark ML, движок может объединить все эти шаги в один проход по данным. Более того, для исследователей данных может быть выгодно применять объединённый набор библиотек (например, Python или R) при моделировании, а веб-разработчикам пригодятся унифицированные фреймворки, такие как Node.js или Django.
- Spark оптимизирует своё машинное ядро для эффективных вычислений — “то есть Spark только управляет загрузкой данных из систем хранения и производит вычисления над ними, но сам не является конечным постоянным хранилищем. Со Spark можно работать, когда имеешь дело с широким разнообразием постоянных систем хранения, включая системы облачного типа по примеру Azure Storage и Amazon S3, распределенные файловые системы, такие как Apache Hadoop, пространства для хранения ключей, как Apache Cassandra, и последовательностей сообщений, как Apache Kafka. И всё же, Spark не сохраняет данные сам по себе надолго и не поддерживает ни одну из этих систем. Главная причина здесь в том, что большинство данных уже находится в нескольких системах хранения. Перемещать данные дорого, поэтому Spark только обрабатывает данные при помощи вычислительных операций, не важно, где они при этом находятся.” (из книги Databricks). Сфокусированность Sparks на вычислениях отличает его от более ранних программных платформ по обработке больших данных, например от Apache Hadoop. Это ПО включает в себя и систему хранения (HFS, сделанную для недорогих хранилищ на кластерах продуктовых серверов Defining Spark 4) и вычислительную систему (MapReduce). Между собой они интегрируются достаточно хорошо. И всё же это тяжело реализовать с участием только одной части без применения второй или, что важнее, написать приложения, которые имеют доступ к данным, хранящимся где-то еще. Spark также широко применяется сейчас в средах, где в архитектуре Hadoop нет смысла. Например, на публичном облаке (где хранение можно купить отдельно от обработки) или в потоковых приложениях.
- Библиотеки Spark дарят очень широкую функциональность — сегодня стандартные библиотеки Spark являются главной частью этого проекта с открытым кодом. Ядро Spark само по себе не слишком сильно изменялось с тех пор, как было выпущено, а вот библиотеки росли, чтобы добавлять ещё больше функциональности. И так Spark превратился в мультифункциональный инструмент анализа данных. В Spark есть библиотеки для SQL и структурированных данных (Spark SQL), машинного обучения (MLlib), потоковой обработки (Spark Streaming и более новый Structured Streaming) и аналитики графов (GraphX). Кроме этих библиотек есть сотни открытых сторонних библиотек, начиная от тех, что работают с коннекторами и до вариантов для различных систем хранения и алгоритмов машинного обучения.
Apache Spark или Hadoop MapReduce…Что вам подходит больше?
Если отвечать коротко, то выбор зависит от конкретных потребностей вашего бизнеса, естественно. Подытоживая свои исследования, скажу, что Spark выбирают в 7-ми из 10-ти случаев. Линейная обработка огромных датасетов — преимущество Hadoop MapReduce. Ну а Spark знаменит своей быстрой производительностью, итеративной обработкой, аналитикой в режиме реального времени, обработкой графов, машинным обучением и это ещё не всё.
Хорошие новости в том, что Spark полностью совместим с экосистемой Hadoop и работает замечательно с Hadoop Distributed File System (HDFS — Распределённая файловая система Hadoop), а также с Apache Hive и другими похожими системами. Так что, когда объёмы данных слишком огромные для того, чтобы Spark мог удержать их в памяти, Hadoop может помочь преодолеть это затруднение при помощи возможностей его файловой системы. Привожу ниже пример того, как эти две системы могут работать вместе:
![]()
Это изображение наглядно показывает, как Spark использует в работе лучшее от Hadoop: HDFS для чтения и хранения данных, MapReduce — для дополнительной обработки и YARN — для распределения ресурсов.
Дальше я пробую сосредоточиться на множестве преимуществ Spark перед Hadoop MapReduce. Для этого я сделаю краткое поверхностное сравнение.
![]()
Скорость
- Apache Spark —это вычислительный инструмент, работающий со скоростью света. Благодаря уменьшению количества чтения-записи на диск и хранения промежуточных данных в памяти, Spark запускает приложения в 100 раз быстрее в памяти и в 10 раз быстрее на диске, чем Hadoop.
- Hadoop MapReduce— MapReduce читает и записывает на диск, а это снижает скорость обработки и эффективность в целом.
Просто пользоваться
- Apache Spark— многие библиотеки Spark облегчают выполнение большого количества основных высокоуровневых операций при помощи RDD (Resilient Distributed Dataset/эластичный распределённый набор данных).
- Hadoop — в MapReduce разработчикам нужно написать вручную каждую операцию, что только усложняет процесс при масштабировании сложных проектов.
Обработка больших наборов данных
- Apache Spark— так как, Spark оптимизирован относительно скорости и вычислительной эффективности при помощи хранения основного объёма данных в памяти, а не на диске, он может показывать более низкую производительность относительно Hadoop MapReduce в случаях, когда размеры данных становятся такими огромными, что недостаточность RAM становится проблемой.
- Hadoop —Hadoop MapReduce позволяет обрабатывать огромные наборы данных параллельно. Он разбивает большую цепочку на небольшие отрезки, чтобы обрабатывать каждый отдельно на разных узлах данных. Если итоговому датасету необходимо больше, чем имеется в доступе RAM, Hadoop MapReduce может сработать лучше, чем Spark. Поэтому Hadoop стоит выбрать в том случае, когда скорость обработки не критична и решению задач можно отвести ночное время, чтобы утром результаты были готовы.
Функциональность
Apache Spark — неизменный победитель в этой категории. Ниже я даю список основных задач по анализу больших данных, в которых Spark опережает Hadoop по производительности:
- Итеративная обработка. Если по условию задачи нужно обрабатывать данные снова и снова, Spark разгромит Hadoop MapReduce. Spark RDD активирует многие операции в памяти, в то время как Hadoop MapReduce должен записать промежуточные результаты на диск.
- Обработка в почти что реальном времени. Если бизнесу нужны немедленные инсайты, тогда стоит использовать Spark и его обработку прямо в памяти.
- Обработка графов. Вычислительная модель Spark хороша для итеративных вычислений, которые часто нужны при обработке графов. И в Apache Spark есть GraphX — API для расчёта графов.
Машинное обучение. В Spark есть MLlib — встроенная библиотека машинного обучения, а вот Hadoop нужна третья сторона для такого же функционала. MLlib имеет алгоритмы “out-of-the-box” (возможность подключения устройства сразу после того, как его достали из коробки, без необходимости устанавливать дополнительное ПО, драйверы и т.д.), которые также реализуются в памяти.
- Объединение датасетов. Благодаря скорости Spark может создавать все комбинации быстрее, а вот Hadoop показывает себя лучше в объединении очень больших наборов данных, которым нужно много перемешивания и сортировки.
А вот и визуальный итог множества возможностей Spark и его совместимости с другими инструментами обработки больших данных и языками программирования:
![]()
- Spark Core — это базовый инструмент для крупномасштабной параллельной и распределённой обработки данных. Кроме того, есть дополнительные библиотеки, встроенные поверх ядра. Они позволяют разделить рабочие нагрузки для стриминга, SQL и машинного обучения. Отвечают за управление памятью и восстановление после ошибок, планирование, распределение и мониторинг задач в кластере, а также взаимодействие с системами хранения.
- Cluster management (управление кластером) — контроль кластера используется для получения кластерных ресурсов, необходимых для решения задач. Spark Core работает на разных кластерных контроллерах, включая Hadoop YARN, Apache Mesos, Amazon EC2 и встроенный кластерный менеджер Spark. Такая служба контролирует распределение ресурсов между приложениями Spark. Кроме того, Spark может получать доступ к данным в HDFS, Cassandra, HBase, Hive, Alluxio и любом хранилище данных Hadoop.
- Spark Streaming — это компонент Spark, который нужен для обработки потоковых данных в реальном времени.
- Spark SQL — это новый модуль в Spark. Он интегрирует реляционную обработку с API функционального программирования в Spark. Поддерживает извлечение данных, как через SQL, так и через Hive Query Language. API DataFrame и Dataset в Spark SQL обеспечивают самый высокий уровень абстракции для структурированных данных.
- GraphX — API Spark для графов и параллельных вычислений с графами. Так что он является расширением Spark RDD с графом устойчивого распределения свойств (Resilient Distributed Property Graph).
- MLlib (Машинное обучение): MLlib расшифровывается как библиотека машинного обучения. Нужна для реализации машинного обучения в Apache Spark.
Заключение
Вместе со всем этим массовым распространением больших данных и экспоненциально растущей скоростью вычислительных мощностей инструменты вроде Apache Spark и других программ, анализирующих большие данные, скоро будут незаменимы в работе исследователей данных и быстро станут стандартом в индустрии реализации аналитики больших данных и решении сложных бизнес-задач в реальном времени.
Для тех, кому интересно погрузиться глубоко в технологию, которая стоит за всеми этими внешними функциями, почитайте книгу Databricks — “A Gentle Intro to Apache Spark” или “Big Data Analytics on Apache Spark”.
Практика использования Spark SQL, или Как не наступить на грабли
Если вы работаете с SQL, то вам это будет нужно очень скоро. Apache Spark – это один из инструментов, входящих в экосистему Hadoop, который обрабатывает данные в оперативной памяти. Одним из его расширений является Spark SQL, позволяющий выполнять SQL-запросы над данными. Spark SQL удобно использовать для работы посредством SQL-запросов с большими объемами данных и в системах с высокой нагрузкой.
Ниже вы найдёте некоторые нехитрые приёмы по работе со Spark SQL:
- Как с помощью сбора статистики и использования хинтов оптимизировать план выполнения запроса.
- Как, оставаясь в рамках SQL, эффективно обрабатывать соединения по ключам с неравномерным распределением значений (skewed joins).
- Как организовать broadcast join таблицы, если её размер слишком велик.
- Как средствами Spark SQL понять, сколько приложение Spark реально использовало памяти и ядер кластера в развёртке по времени.
Spark SQL помогает в обработке данных многим компаниям из Global 100. При простоте разработки Spark SQL даёт на выходе высокую производительность, если использовать его правильно. Spark SQL можно использовать в различных ETL процессах при обработке и загрузке данных. Альтернативами Spark SQL при работе с данными можно назвать использование Impala SQL или Hive (c обработкой SQL-запросов в mapreduce). Если же выходить за рамки SQL, то альтернатива – использование в Spark языка DSL. Не приходится сомневаться в перспективности языка SQL, существующего десятилетия и имеющего миллионы пользователей. Немалые перспективы развития имеет и Spark SQL. Например, с новыми версиями Spark появляется возможность использовать новые оптимизационные хинты, способы обработки данных становятся эффективнее. Spark развивается от версии к версии в погоне за повышением эффективности обработки данных и интеграцией с DeepLearning. Меняются в сторону улучшения подходы: RDD, DataFrame, DataSet. Тем самым, изучение и применение Spark SQL – актуально и многообещающе. Для компаний же использование Spark SQL ведёт в конечном итоге к накоплению знаний о клиентах, их обработке и построении новых бизнесов на новых знаниях.
Для более тонкой настройки загрузки данных посредством Spark потребуется выйти за рамки Spark SQL и использовать, например, DSL.
Уточним, что речь в статье идёт о работе с версией Apache Spark 2.3, загрузка данных осуществляется из Hadoop в Hadoop, управление ресурсами осуществляется посредством Yarn, используется операционная система Linux, а также СУБД Hive с сохранением данных в hdfs в формате parquet.
Код приложения, запускающего переданный ему в качестве параметра файл с командой SQL полностью приводить не будем, напишем лишь, что в приложении нужно открыть переданный файл, загрузить его содержимое в переменную sqlStr и вызвать:
SparkSession spark = SparkSession .builder() .appName("") .enableHiveSupport() .getOrCreate(); Dataset sql = spark.sql(sqlStr);
Получится приложение *.jar, которому можно передавать в качестве параметра файл с командой SQL и выполнять его посредством spark-submit:
spark-submit --class --name --queue --executor-cores 1 --executor-memory 1g --driver-cores 1 --driver-memory 1g --num-executors 1 --master yarn --deploy-mode cluster sqlFile=;
При этом в передаваемом файле sqlFile можно указывать команды SQL, например:
insert overwrite table target_scheme.target_table select s.* from source_scheme.source_table s;
Сбор и просмотр статистики для Spark, хинт BROADCAST
Как известно, одной из наиболее сложных операций для Spark является соединение (join) таблиц или датасетов. При этом в Spark существуют различные алгоритмы реализации join-ов: SortMergeJoin, BroadcastHashJoin, CartesianProduct и др. Чтобы помочь Spark автоматически понять, большие или маленькие таблицы он соединяет и какой тип соединения оптимален, нужно собрать статистику по этим таблицам. Для этого посредством spark-submit вызываем команду вида:
ANALYZE TABLE scheme_name.table_name COMPUTE STATISTICS;
Эту команду сбора статистики мы можем передать вышеописанному модулю sqlrunner.jar в качестве параметра вместо команды SQL. Для выполнения этой команды требуются права на чтение и запись в таблицу, по которой собирается статистика.
Проверить, что статистика собрана, можно в среде hive командой вида:
show create table scheme_name.table_name;
Нужно посмотреть, появились ли в конце описания в блоке TBLPROPERTIES свойства ‘spark.sql.statistics.numRows’ и ‘spark.sql.statistics.totalSize’:
CREATE EXTERNAL TABLE `scheme_name.table_name`( … TBLPROPERTIES ( … 'spark.sql.statistics.numRows'='363852167', 'spark.sql.statistics.totalSize'='82589603650', …
Таким образом, целесообразно иметь собранную статистику по таблицам-источникам, по используемым в расчётах промежуточным таблицам, а также по таблицам витрин, которые тоже могут быть использованы в соединениях в пользовательских запросах.
Теперь скажем пару слов о том, что такое broadcast в Spark. При присоединении маленькой таблицы к большому датасету часто оказывается целесообразным не распределять маленькую таблицу на различные ноды по ключам соединения в виде порций данных, а поместить её целиком на все ноды кластера, использующиеся в приложении. При этом говорят, что осуществляется broadcast маленькой таблицы на все ноды и тип присоединения маленькой таблицы к основному массиву данных – broadcast join.
Зачастую, broadcast join нужен в Spark SQL в тех случаях, когда в реляционных базах данных требуется nested loop join. Особенно необходим broadcast при неравномерной статистике распределения данных в участвующих в соединении ключах, например, при присоединении маленького справочника валют к большой таблице проводок, в которой 99% строк относится к рублям.
Для того, чтобы указать Spark, что надо сделать broadcast какой-то небольшой таблицы или датасета (до ~1Gb), можно в SQL-запросе указывать хинт /*+ BROADCAST(t)*/, где t – алиас таблицы или датасета.
insert overwrite table target_scheme.target_table select /*+ BROADCAST(t) */ big.field1, big.field2, t.field3 from source_scheme.big_table as big left join source_scheme.small_table as t on big.field1 = t.field1;
В том числе, если по какой-то таблице статистика не собрана или мы имеем дело с подзапросом, результат которого формируется “на лету” и Spark заранее не знает, большой ли он, то хинт /*+ BROADCAST(t)*/ особенно целесообразен, если в таблице или подзапросе мало данных.
Если broadcast какого-то подзапроса завершается позже чем, через 5 минут с начала выполнения запроса, то стоит вызывать spark-submit с нижеприведённым параметром (здесь 36000 – время в секундах), потому как выделенных по умолчанию 5 минут может не хватить:
--conf spark.sql.broadcastTimeout=36000
Например, такой параметр пригодится в нижеприведённом запросе, где делается broadcast результата подзапроса /*+ BROADCAST(small) */, который становится готов не с самого начала работы запроса, а на более поздних этапах:
select /*+ BROADCAST(small) */ big.field3, small.field1, small.field2 from source_scheme.big_table as big inner join (select t1.field1, t2.field2 from source_scheme.table1 as t1 inner join source_scheme.table2 as t2 on t1.field1 = t2.field1) small on big.field2 = small.field2;
Убедиться, что хинт учтён оптимизатором и broadcast реально осуществлён можно, посмотрев в Yarn по ссылке в столбце Tracking UI у соответствующего приложения Spark план выполнения запроса на закладке SQL:

Отметим, что в Spark SQL есть и другие хинты, в т.ч. с версии 2.4 появляются хинты /*+ COALESCE(n) */, где n – количество партиций, на которые будет разбит результат, и /* + REPARTITION (n) */, где n – количество партиций при repartition.
Обход SKEWED JOIN явным указанием SKEWED ключей
Рассмотрим две таблицы:
- Таблица A, имеющая поле AID
- Таблица B, имеющая поле BID
select A.*, B.* from A inner join B on A.AID = B.BID;
Рассмотрим случай, когда значения ключей, которые встречаются в одной из таблиц, распределены неравномерно.
Пусть в таблице B на значения ключей BID 4089, 4107, 4468 и 6802 приходится в сумме 70% строк, в то время как на миллион прочих значений ключей BID приходится лишь оставшиеся 30% строк.
При этом будем считать, что в таблице A значения ключей AID распределены более-менее равномерно, при этом значения 4089, 4107, 4468 и 6802 встречаются по одному разу.
В таком случае это соединение таблиц называется skewed join, т.е. “перекошенное соединение”, а слишком частые ключи называются skewed keys.
По умолчанию задача соединения двух таблиц, выполненная посредством spark-submit, будет разбита на 200 партиций. Изменить эту настройку по умолчанию можно, задав другое значение конфигурации spark.sql.shuffle.partitions, например:
--conf spark.sql.shuffle.partitions=1000
Как Spark будет обрабатывать skewed join?
Spark без подсказки не знает, что данные распределены неравномерно, и распределит задачу соединения таблиц так: большинство значений ключей попадёт в “лёгкие” партиции, которые быстро обработают небольшой объём данных; при этом обработка skewed keys попадёт в “тяжёлые” партиции.
Картинки ниже показывают пример такого случая: 999 из 1000 партиций отработали быстро, выдав небольшое количество итоговых строк, а одна партиция, в которую попали почти все данные, работала несколько часов.


Как сделать распределение соединяемых данных по партициям более равномерным?
Как вариант, нужно соединять таблицы не по одному ключу A.AID = B.BID, а по паре ключей, при этом распределение данных по паре ключей сделать уже более равномерным.
Сделаем допущение, что в таблице B имеется также числовое поле BID2, данные в котором распределены равномерно (в отличие от BID). Если такового равномерно распределенного числового поля нет, то его можно сгенерировать работающей в hive и spark функцией hash() из конкатенации других полей. В таком случае функция hash() выдаст равномерно распределённые данные.
Нижеприведённый работающий под spark запрос решает задачу. Этот запрос целесообразно применить вместо “select A.*, B.* from A inner join B on A.AID = B.BID;”
В секции with приведён пример вышеописанных данных для таблиц A и B. Имеет место перекошенность данных в ключе BID. Формируется подзапрос ANTISKEW, в котором перекошенные значения ключей 4089, 4107, 4468 и 6802 размножены до 10 экземпляров каждое. Данные из таблицы А соединяются с подзапросом ANTISKEW, тем самым сформировав в массиве A_ для ключей AID, равных 4089, 4107, 4468 и 6802, по 10 экземпляров с различными значениями ANTISKEWKEY от 0 до 9. При этом для прочих значений AID в массиве A_ будет по одному экземпляру с ANTISKEWKEY = 0. Далее массив A_ соединяется с таблицей B не только по условию A_.AID = B.BID, но и по второй паре ключей. А именно, для перекошенных значений ключей BID многочисленные записи из таблицы B соединятся с одним из значений ANTISKEWKEY от 0 до 9 в зависимости от остатка от деления BID2 на 10. Неперекошенные значения BID соединятся один к одному с ANTISKEWKEY = 0. Тем самым будет достигнуто более равномерное распределение перекошенных ключей BID по партициям – они пойдут не в одну, а в 10 партиций.
with A as (select 4089 AID union all select 4468 AID union all select 6802 AID union all select 5 AID union all select 8 AID union all select 14 AID), B as (select 4089 BID, 1 BID2 union all select 4089 BID, 2 BID2 union all select 4089 BID, 3 BID2 union all select 4107 BID, 4 BID2 union all select 4107 BID, 5 BID2 union all select 4107 BID, 6 BID2 union all select 4468 BID, 7 BID2 union all select 4468 BID, 8 BID2 union all select 4468 BID, 9 BID2 union all select 6802 BID, 10 BID2 union all select 6802 BID, 11 BID2 union all select 6802 BID, 12 BID2 union all select 1 BID, 13 BID2 union all select 2 BID, 14 BID2 union all select 3 BID, 15 BID2 union all select 4 BID, 16 BID2 union all select 5 BID, 17 BID2 union all select 6 BID, 18 BID2 union all select 7 BID, 19 BID2 union all select 8 BID, 20 BID2 union all select 9 BID, 21 BID2 union all select 10 BID, 22 BID2 union all select 11 BID, 23 BID2 union all select 12 BID, 24 BID2 union all select 13 BID, 25 BID2 union all select 14 BID, 26 BID2 union all select 15 BID, 27 BID2) select A_.*, B.* from (select A.*, coalesce(ANTISKEW.ANTISKEWKEY, 0) ANTISKEWKEY from A left join (select Y.ID, Z.ANTISKEWKEY from (select 4089 ID union all select 4107 ID union all select 4468 ID union all select 6802 ID) Y cross join (select 0 ANTISKEWKEY union all select 1 ANTISKEWKEY union all select 2 ANTISKEWKEY union all select 3 ANTISKEWKEY union all select 4 ANTISKEWKEY union all select 5 ANTISKEWKEY union all select 6 ANTISKEWKEY union all select 7 ANTISKEWKEY union all select 8 ANTISKEWKEY union all select 9 ANTISKEWKEY) Z) ANTISKEW on A.AID = ANTISKEW.ID) A_ inner join B on A_.AID = B.BID and A_.ANTISKEWKEY = case when A_.AID in (4089, 4107, 4468, 6802) then B.BID2%10 else 0 end;
Для ещё более равномерного распределения данных нужно подбирать в подзапросе ANTISKEW, на сколько именно партиций хочется разбить каждый ключ. Описанный метод требует знания статистики – по каким перекошенным ключам сколько процентов строк следует ожидать.
Партиционирование помогает сделать BROADCAST
Ниже описан способ, как в Spark SQL делать broadcast join вместо shuffle join в случае, если размер меньшей таблицы слишком велик для broadcast.
Допустим, нужно соединить таблицу big_table (100Gb)
create table big_table( id decimal(18,0), key1 decimal(18,0) ) stored as parquet;
с таблицей small_table (4Gb)
create table small_table( id decimal(18,0), key2 decimal(18,0) ) stored as parquet;
При этом предполагается, что размер таблицы small_table слишком велик, чтобы сделать её broadcast, о чём выдаётся ошибка.
Целью является заполнить результирующую таблицу:
truncate table result_table; insert into table result_table select big_table.id, big_table.key1, small_table.key2 from big_table left join small_table on big_table.id = small_table.id;
Соединение по ключу id при этом предполагается нормально распределённым.
Сделаем обе соединяемые таблицы партиционированными:
create table big_table( id decimal(18,0), key1 decimal(18,0) ) partitioned by (part_mod decimal(6,0)) stored as parquet; create table small_table( id decimal(18,0), key2 decimal(18,0) ) partitioned by (part_mod decimal(6,0)) stored as parquet;
При заполнении обеих партиционированных таблиц big_table и small_table ключ партиционирования part_mod сделаем равным остатку от деления id на 4:
Заведём псевдотаблицу, смотрящую своим location на одну партицию big_table (part_mod = 0):
create table big_table_0( id decimal(18,0), key1 decimal(18,0) ) stored as parquet location '…/big_table/part_mod=0';
Также заведём псевдотаблицу, смотрящую своим location на одну партицию small_table (part_mod = 0):
create table small_table_0( id decimal(18,0), key1 decimal(18,0) ) stored as parquet location '…/small_table/part_mod=0';
Теперь мы можем записать в результирующую таблицу одну четверть требуемых данных из соединения вышеописанных псевдотаблиц.
При этом проверено, что хинт BROADCAST успешно срабатывает:
truncate table result_table; insert into table result_table select /*+ BROADCAST(small_table_0)*/ big_table_0.id, big_table_0.key1, small_table_0.key2 from big_table_0 left join small_table_0 on big_table.id = small_table.id;
Следующими тремя шагами мы можем аналогично дописать в результирующую таблицу оставшиеся три четверти данных с part_mod = 1, part_mod = 2, part_mod = 3.
Итоговая продолжительность загрузки результирующей таблицы «методом поэтапного broadcast партиций» обычно получается на одинаковых ресурсах меньше, чем методом shuffle join. Однако такой метод требует, чтобы таблицы с исходными данными были заранее партиционированы одинаковым образом, т.е. по остатку от деления id на какое-либо число, по которому соединяемся. Если id не целочисленное, то можно брать остаток от деления hash(id) на какое-либо число. Особенно легко организовать, чтобы заранее партиционированными были промежуточные (стейджинговые) таблицы, дизайн которых находится на усмотрении разработчика.
Мониторинг ресурсов
Как известно, при работе с озером данных важно правильно выделять ресурсы задачам по загрузке данных. В противном случае имеются риски падения задачи от нехватки ресурсов, слишком длительного выполнения задачи, либо неоправданно высокого выделения ресурсов в ущерб другим задачам.
Как правильно выделять ресурсы – отдельная тема. В этой же статье мы поговорим о том, как понять:
- Сколько ресурсов выделил Yarn на задачу или группу задач в развёртке по времени.
- Сколько ресурсов реально использовал Spark на задачу из числа выделенных ему менеджером ресурсов Yarn в развёртке по времени.
Сначала опишем первый вид мониторинга (“мониторинг yarn”).
while [ ! -f complete.flg ]; do for app in $(yarn application -list -appStates RUNNING | grep | awk ''); do ( command1 () < date +"%Y-%m-%d %H:%M:%S"; >; command2 () < yarn application -status $app; >; x="|""|$(command1) $(command2)"; echo $x >> app.txt; ) & done; wait; done
осуществляется следующее. В период работы мониторинга, когда установлен флаг в виде файла complete.flg, в цикле по списку находящихся в статусе RUNNING приложений Yarn, имеющих в названии шаблон , в текстовый файл app.txt дозаписывается результат вызова команды yarn application –status $app. Таким образом, в текстовом файле app.txt для каждого момента мониторинга содержится информация о затраченных за время работы каждого приложения мегабайт*секундах (MB-seconds) и ядро*секундах (vcore-seconds) накопительным итогом. По окончании работы мониторинга ставится флаг завершения complete.flg и итоговый файл app.txt переносится в hdfs в текстовую таблицу yarn_monitoring:
CREATE EXTERNAL TABLE IF NOT EXISTS scheme_name.yarn_monitoring( loading_date string, loading_id string, txt string ) ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe' WITH SERDEPROPERTIES ( 'field.delim'='|', 'serialization.format'='|') STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat' LOCATION '…/yarn_monitoring';
Строки в этой таблице имеют вид:
2020-02-18|2235599|2020-02-18 01:00:44 Application Report : Application-Id : application_1580486634374_3632578 Application-Name : Application-Type : SPARK User : Queue : Start-Time : 1581976835176 Finish-Time : 0 Progress : 10% State : RUNNING Final-State : UNDEFINED Tracking-URL : RPC Port : 0 AM Host : Aggregate Resource Allocation : 18663 MB-seconds, 9 vcore-seconds Log Aggregation Status : NOT_START Diagnostics :
Затем можно на том же Spark SQL написать запрос к этой текстовой таблице, в котором распарсить значения показателей использованных ресурсов (MB-seconds, vcore-seconds) на каждый момент осуществления мониторинга. Далее в этом запросе можно визуализировать эти показатели в виде графиков в формате векторной графики *.svg:
… 13:47:48 …
Результат запроса можно сохранить в файл *.svg и посмотреть в браузере. Пример картинки с результатами мониторинга ресурсов, использованных группой наблюдаемых приложений Yarn, приведен ниже. Красным цветом показана динамика использованной памяти, синим цветом — динамика использованных процессорных ядер:

На основании такой визуализации можно принимать решения об изменении расписания выполнения определённых приложений Yarn с целью более равномерного использования ресурсов в течение суток.
Теперь опишем второй вид мониторинга (“мониторинг spark”).
Итак, интересна ситуация, когда Yarn выделил задаче Spark, например, 200 ядер и не отдаёт эти ресурсы другим задачам. При этом в реальности 199 из 200 ядер уже не используются в приложении Spark, а работу реально продолжает только одно ядро. Этот вид мониторинга среди прочего нацелен на поиск таких ситуаций.
Считаем, что подвергаемые мониторингу приложения Spark вызываются из некого управляющего механизма. В случае Сбербанка это, например, Oozie или Informatica BDM. Таким образом, управляющий механизм может по завершению работы приложения Spark скопировать его лог в hdfs и запустить маленькое приложение на Spark SQL по парсингу этого лога (или группы логов) и визуализации динамики реально использовавшихся Spark ядер в формате *.svg.
Итак, с помощью команды нижеприведённого вида мы:
- Сохраняем в переменную $app идентификатор приложения в yarn, например, «application_1576046768499_16691»
- С помощью команды command1 записываем в файл $app».txt» строку-заголовок со статусом завершившегося приложения Spark
- С помощью команды «yarn logs -applicationId $app»дозаписываем в файл $app».txt» полученный из Yarn лог завершившегося приложения Spark
- Переносим сформированный файл $app».txt» со статусом и логом в hdfs
'app=$(grep "(state: FINISHED)" /tmp/'' | grep -o "application_[^ ]*" | tee /tmp/applicationid_''.txt); command1 () < yarn application -status $app; >; x ; echo $x > "/tmp/"$app".txt"; yarn logs -applicationId $app | nl -ba -s"|" | sed "s/^/''|''|"$app"|/" >> "/tmp/"$app".txt"; hdfs dfs -copyFromLocal -f "/tmp/"$app".txt" …/spark_monitoring; rm "/tmp/"$app".txt"'
В результате в hdfs получаем строки вида:
2019-12-13|1227125|application_1576046768499_16691|0|Application Report : Application-Id : application_1576046768499_16691 Application-Name : Application-Type : SPARK User : Queue : Start-Time : 1576229137128 Finish-Time : 1576229189406 Progress : 100% State : FINISHED Final-State : SUCCEEDED Tracking-URL : http://. sbrf.ru:18089/history/application_1576046768499_16691/1 RPC Port : 0 AM Host : … Aggregate Resource Allocation : 93763306 MB-seconds, 9957 vcore-seconds Log Aggregation Status : NOT_START Diagnostics : 2019-12-13|1227125|application_1576046768499_16691| 1| 2019-12-13|1227125|application_1576046768499_16691| 2| 2019-12-13|1227125|application_1576046768499_16691| 3|Container: container_1576046768499_16691_01_000162 on …sbrf.ru_8041 2019-12-13|1227125|application_1576046768499_16691| 4|=========================================================================================== 2019-12-13|1227125|application_1576046768499_16691| 5|LogType:container-localizer-syslog 2019-12-13|1227125|application_1576046768499_16691| 6|Log Upload Time:Fri Dec 13 12:26:32 +0300 2019 2019-12-13|1227125|application_1576046768499_16691| 7|LogLength:0 2019-12-13|1227125|application_1576046768499_16691| 8|Log Contents: 2019-12-13|1227125|application_1576046768499_16691| 9| … 2019-12-13|1227125|application_1576046768499_16691| 53|19/12/13 12:26:02 INFO executor.Executor: Running task 26.0 in stage 0.0 (TID 19) … 2019-12-13|1227125|application_1576046768499_16691| 82|19/12/13 12:26:07 INFO executor.Executor: Finished task 26.0 in stage 0.0 (TID 19). 5735 bytes result sent to driver … 2019-12-13|1227125|application_1576046768499_16691| 84|19/12/13 12:26:07 INFO executor.Executor: Running task 441.0 in stage 0.0 (TID 364) … 2019-12-13|1227125|application_1576046768499_16691| 102|19/12/13 12:26:08 INFO executor.Executor: Finished task 441.0 in stage 0.0 (TID 364). 5543 bytes result sent to driver …
При парсинге в логе приложения Spark ищутся события начала использования и высвобождения ядер кластера тем или иным контейнером, т.е. мы ищем время записи в лог событий по шаблонам:
INFO executor.Executor: Running task INFO executor.Executor: Finished task
Таким образом, на любой момент времени известно, сколько процессорных ядер было активно при работе приложения Spark, т.е. начали и еще не завершили свою работу.
Для парсинга собранных в hdfs логов формируется таблица spark_monitoring в текстовом формате:
CREATE EXTERNAL TABLE IF NOT EXISTS scheme_name.spark_monitoring( loading_date string, loading_id string, applicationid string, rn string, txt string ) ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.lazy.LazySimpleSerDe' WITH SERDEPROPERTIES ( 'field.delim'='|', 'serialization.format'='|') STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat' LOCATION '…/spark_monitoring';
Как было написано выше, парсинг логов осуществляется силами Spark SQL. Т.е. написан запрос к этой таблице spark_monitoring, который на выходе даёт картинку в формате *.svg с динамикой реального использования ядер кластера приложением Spark.
Результат запроса можно сохранить в файл *.svg и посмотреть в браузере. Пример картинки приведен ниже:

На основании таких графиков можно понять, эффективно ли приложение Spark использует выделенные ему процессорные ядра кластера. Можно смотреть приложения массово. Резкие скачки вверх-вниз, видные на графиках – это переходы между этапами (stages) приложений Spark. Когда заканчивается какой-то этап по чтению или соединению массивов данных и по одному освобождаются все ядра, получается снижение графика к нулю. А вначале выполнения новых этапов идёт резкий рост используемых ядер. Ситуация, когда график притянут к максимуму выделенных ядер и скачки между этапами резкие – это признак оптимального использования ресурсов. Медленные же снижения могут свидетельствовать, например, о skew joins, о которых рассказывалось в этой статье выше.
Автор: Михаил Гричик, эксперт профессионального сообщества Сбербанка SberProfi DWH/BigData.
Профессиональное сообщество SberProfi DWH/BigData отвечает за развитие компетенций в таких направлениях, как экосистема Hadoop, Teradata, Oracle DB, GreenPlum, а также BI инструментах Qlik, SAP BO, Tableau и др.
- Блог компании Сбер
- Программирование
- SQL
- Администрирование баз данных
Аналитика в режиме реального времени с помощью Spark SQL

В этой статье мы рассмотрим Spark SQL и выясним, почему этот инструмент является предпочтительным, когда речь идет об аналитике реального времени. Spark SQL – это модуль Apache Spark, интегрирующий реляционную обработку данных и процедурный API Spark. Spark SQL является частью ядра Spark с версии 1.0. Он может работать совместно с Hive (HiveQL/SQL) или замещать его. Кроме того, модуль способен взаи модействовать с инструментами бизнес-аналитики.
Для работы с модулем можно использовать Python, Scala и Java. Благодаря Spark SQL , функционал фреймворка получает два ключевых дополнения. Во-первых, модуль обеспечивает тесную интеграцию между реляционной и процедурной обработкой данных посредством интеграции декларативного DataFrame API и процедурного API Spark . Во-вторых, он включает в себя расширяемый оптимизатор, созданный на языке Scala, обладающем широкими возможностями сопоставления с образцом (pattern matching), что позволяет легко формировать правила, управлять генерацией кода и создавать расширения.
Предназначение Spark SQL
Несмотря на то, что реляционный подход может быть использован для решения задач в области больших данных, применительно ко многим из них одного этого подхода недостаточно. До недавнего времени реляционный и процедурный подходы существовали независимо друг от друга, вынуждая разработчиков выбирать из них какой-либо один. Теперь, благодаря Spark SQL , оба подхода могут использоваться совместно.
Spark SQL поддерживает реляционную обработку как в рамках программ Spark (посредством RDD ), так и применительно ко внешним источникам данных. Он также способен взаимодействовать с новыми источниками данных, включая слабоструктурированные данные и внешние базы данных, поддерживающие федеративные запросы (federated query).
Как говорится: «Самый быстрый способ прочитать данные – НЕ читать их вообще». Spark SQL реализует данную парадигму с помощью следующих подходов:
- Преобразование данных в более эффективные форматы (с точки зрения хранилища, сети и операций ввода/вывода), в частности, в различные форматы, ориентированные на столбцы (columnar format).
- Секционирование данных.
- Уменьшение количества операций чтения на основе статистики.
- Оптимизация предикат.
- Выполнение оптимизации как можно позже, когда доступна вся информация о конвейерах данных.
Spark SQL и DataFrame используют оптимизатор запросов Catalyst для интеллектуального планирования выполнения запросов.
Spark Streaming + Spark SQL и DataFrame
Spark SQL может поддерживать пакетный и потоковый SQL . Ядро Spark обеспечивает обработку пакетных нагрузок посредством RDD . RDD могут ссылаться на статические наборы данных, а с помощью обширного API Spark можно манипулировать RDD в оперативной памяти с применение «ленивых» вычислений.
Давайте кратко рассмотрим, что такое RDD и DStream. Понимание этих базовых концепций понадобится нам для дальнейшего изложения.
Spark оперирует данными в форме RDD . RDD ( resilient distributed dataset , отказоустойчивый распределенный набор данных) – это распределенная структура данных, размещаемая в оперативной памяти. Каждый RDD представляет собой фрагмент данных, распределенных по узлам кластера. RDD являются неизменяемыми структурами, поэтому после преобразований создаются новые RDD . RDD обрабатываются параллельно с помощью таких преобразований/действий, как отображение, фильтрация и др. Эти операции выполняются одновременно во всех разделах (partition). RDD являются отказоустойчивыми: если раздел теряется в результате сбоя узла, он может быть восстановлен из исходных источников.

Spark Streaming реализует абстракцию под названием DStream (discretized stream, дискретизированный поток), представляющую собой непрерывный поток данных. DStream может быть создан следующим образом: из потока входных данных; на основе таких источников, как Kafka или Flume; или посредством выполнения операций с другими DStream. По сути, DStream – это последовательность RDD .

RDD , созданный посредством DStream, можно преобразовать в DataFrame и выполнять к нему запросы с помощью SQL. Доступ к потоку можно предоставить для любого внешнего приложения, поддерживающего SQL , с помощью JDBC-драйвера. Пакеты потоковых данных хранятся в памяти узла. Эти данные можно интерактивно запрашивать, используя SQL или API Spark.
StreamSQL – это компонент Spark, объединяющий Catalyst и Spark Streaming для выполнения SQL-запросов к DStream. StreamSQL расширяет SQL, обеспечивая поддержку следующих потоковых операций:
- Выборка (SELECT) из потока для вычисления функций или фильтрации ненужных данных (с помощью условия WHERE).
- Соединение (JOIN) потока с одним или несколькими наборами данных для создания нового потока.
- Применение оконных функций и выполнение агрегаций. Поток можно настроить таким образом, чтобы он создавал наборы данных ограниченного размера. С помощью оконных функций можно выполнять сложный отбор сообщений на основе значений полей. После создания ограниченного пакета можно выполнять аналитику.
Компоненты Spark SQL
- Выполнение запросов.
- Чтение данных в различных форматах: Parquet, JSON, Avro и др.
- Чтение SQL — и NoSQL -источников данных.
- HQL, MetaStore, сериализация/десериализация (SerDes), пользовательские функции (UDF).
- Оптимизация реляционной алгебры и выражений.
- Оптимизация запросов.
Преимущества Spark SQL
Spark SQL реализует концепцию интегрированной платформы, в рамках которой нет необходимости в перемещении данных вне кластера. Также не требуется установка дополнительных модулей. Spark SQL обеспечивает единый интерфейс загрузки и сохранения данных вне зависимости от источника данных и языка программирования.
Пример ниже демонстрирует, насколько легко можно загрузить данные из Avro и преобразовать их в Parquet.
val df = sqlContext.load("mydata.avro", "com.databricks.spark.avro") df.save("mydata.parquet", "parquet")
Spark (включая Spark SQL и Spark Streaming ) представляет собой единый фреймворк для решения аналитических задач, связанных как с пакетными, так и с потоковыми данными. Подобный инструмент в течение долгого времени был Святым Граалем в области обработки данных. Ранее существовали фреймворки, которые решали задачи лишь в одной плоскости. И хотя они обеспечивали хорошую масштабируемость, производительность и функциональность, о едином фреймворке все могли только мечтать, пока не появился Spark. Благодаря Spark один и тот же код (логика) может работать как с пакетными (RDD), так и с потоковыми (DStream) данными. DStream – это просто последовательность RDD. Такое представление стирает грань между пакетными и потоковыми нагрузками. Благодаря этому значительно сокращаются накладные расходы на обслуживание кода и обучение разработчиков, которые больше не должны осваивать два различных набора навыков.
Чтение из JDBC-источников данных
Spark SQL имеет встроенную поддержку чтения из JDBC-источников данных. Модуль позволяет извлекать данные из любых реляционных баз данных, поддерживающих JDBC, например, MySQL, PostgreSQL, H2 и др. Чтение данных из таких систем выполняется крайне просто и сводится к созданию виртуальной таблицы, ссылающейся на внешнюю таблицу. Затем данные из этой таблицы можно легко читать и соединять с другими источниками, поддерживаемыми Spark SQL.
Spark SQL и DataFrame
DataFrame – это распределенная коллекция данных, организованных посредством именованных столбцов. Данная абстракция предназначена для выборки, фильтрации, агрегации и визуализации структурированных данных. Ранее эта структура называлась SchemaRDD.
DataFrame API позволяет выполнять реляционные операции как с внешними источниками данных, так и со встроенными распределенными коллекциями Spark.
DataFrame поддерживает глубокую реляционную/процедурную интеграцию в рамках программ Spark и позволяет манипулировать данными как с помощью процедурного API Spark, так и посредством нового реляционного API, обеспечивающего более эффективную оптимизацию. DataFrame может быть создан непосредственно из RDD , что обеспечивает возможность реляционной обработки уже имеющихся данных.
DataFrame предоставляет более удобные и эффективные средства обработки данных, чем процедурный API Spark. В частности, можно вычислить несколько агрегаций за один проход с мощью SQL -инструкции, что достаточно сложно реализовать посредством традиционного процедурного API.
В отличие от RDD, DataFrame отслеживает свою схему и поддерживает различные реляционные операции, что обеспечивает более оптимизированное выполнение. DataFrame формирует схему посредством отражения (reflection).
DataFrame является «ленивой» структурой данных, то есть содержит логический план для вычисления набора данных, при этом вычисления не выполняются до тех пор, пока пользователь не запросит специальную «операцию вывода», например, сохранение. Такой подход обеспечивает эффективную оптимизацию всех операций.
Концепция DataFrame расширяет модель RDD. В результате, благодаря упрощенным методам фильтрации и агрегации, Spark-разработчики получают возможность быстрее и эффективнее работать с большими наборами структурированных данных. Для работы с DataFrame доступны API на Java, Scala и Python.
API Spark SQL для работы с источниками данных позволяет читать/записывать DataFrame из/в различных источников и форматов: Avro, Parquet, ORC, JSON, H2.
Пример, демонстрирующий краткость кода, которую обеспечивает DataFrame по сравнению с RDD:
Ниже представлены два эквивалентных фрагмента кода на Scala. В первом из них используется RDD API, а во втором – DataFrame API. Для примера рассмотрим набор данных, содержащий информацию о людях. Атрибутами каждого человека являются: имя, фамилия и возраст. Наша цель заключается в том, чтобы вычислить базовые статистические характеристики возраста людей, сгруппированных по имени.
case class People(firstname: String, lastname: String, age: Intger) val people = rdd.map(p => (people.firstname, people.age)).cache() // RDD Code val minAgeByFN = people.reduceByKey( scala.math.min(_, _) ) val maxAgeByFN = people.reduceByKey( scala.math.max(_, _) ) val avgAgeByFN = people.mapValues(x => (x, 1)) .reduceByKey((x, y) => (x._1 + y._1, x._2 + y._2)) val countByFN = people.mapValues(x => 1).reduceByKey(_ + _) // Data Frame Code df = people.toDF people = df.groupBy("firstname").agg( min("age"), max("age"), avg("age"), count("*"))
Благодаря оптимизатору Catalyst , DataFrame позволяет получить значительное преимущество в скорости по сравнению с RDD . DataFrame поддерживает те же операции, что и реляционные языки, такие как SQL и Pig .
Подход к реализации аналитики реального времени с помощью Spark SQL
На рисунке ниже представлена логическая схема реализации аналитики реального времени с помощью Spark SQL. В основе данного подхода лежит лямбда-архитектура ( lambda architecture ), применяемая для создания аналитических систем реального времени в контексте больших потоковых данных.

Ограничения Spark SQL
Spark SQL имеет те же ограничения, что и любой другой инструмент, работающий в кластере Hadoop. Скорость обработки зависит не столько от скорости системы, сколько от количества пользователей, совместно использующих кластер.
Кроме того, при первом знакомстве код Spark SQL может показаться слишком лаконичным и непростым для понимания. Необходим некоторый опыт и практика, чтобы привыкнуть к стилю программирования и особенностям кода.
Заключение
Spark SQL – результат эволюционного развития базового API Spark. Базовый процедурный API Spark является достаточно общим и предоставляет лишь ограниченные возможности для автоматической оптимизации. Spark SQL существенно расширяет эти возможности.
¿Qué es Spark SQL en Apache Spark?

El proceso de Spark SQL en Apache Spark se ha instaurado como una ventaja para el desarrollo de un estudio de macrodatos efectivo por medio de este sistema de computación y sus servicios.
Es por ello que, en este artículo, te ponemos al tanto de qué es Spark SQL en Apache Spark, de manera que puedas emplearlo en tu procesamiento de datos.
En este post encontrarás: ocultar
¿Qué es Spark SQL en Apache Spark?
Spark SQL en Apache Spark es uno de los módulos para el procesamiento de la información que ofrece Apache Spark y que trabaja con datos estructurados.
Spark SQL en Apache Spark cuenta con las interfaces que proporcionan mayor información sobre la estructura y la computación de los datos que la API de Spark RDD.
Por otra parte, existen varias formas para trabajar con Spark SQL:
- Dataframe: es el conjunto de datos organizados en columnas del Spark SQL, también es equivalente a
una tabla relacional. Estos pueden construirse desde ficheros estructurados, tablas
en Hive, base de datos externas o RRDs existentes. Por último, trabaja lenguajes como Scala, Java, Python o R. - Dataset:con conjunto de datos distribuidos, este permite usar Spark SQL en Apache Spark con los beneficios de los RDDs: tipado fuerte, transformaciones funcionales, etc. Trabaja con los lenguajes Scala y Java.
- SQL:este se emplea con un lenguaje con sintaxis SQL y Spark SQL.
Spark SQL: Base Project
Para empezar un proyecto de Spark SQL, necesitas añadir las dependencias en tu proyecto de sbt en el IDE.
libraryDependencies ++= Seq(
«org.apache.spark» %% «spark-core» % «3.0.1”,
«org.apache.spark» %% «spark-sql» % «3.0.1”
)
Una vez cuentes con las dependencias, podrás crear un objeto principal con un método main,
donde crearás un nuevo contexto de SparkSQL.
Desde este punto podrás ejecutar en el play del IDE y ejecutar SparkSQL en modo local, usando la CPU de la máquina.
Para el siguiente ejemplo usaremos una instancia de SparkSession en lugar de SparkContext:
package io.keepcoding.spark.sql
import org.apache.spark.sql.SparkSession
object SparkSqlBaseProject def main(args: Array[String]): Unit = val spark = SparkSession
.builder()
.appName(«Spark SQL KeepCoding»)
.getOrCreate()
import spark.implicits._
//
spark.close()
>
>
SparkSQL: Apache Parquet
Dentro de Spark SQL en Apache Spark podrás hallar a Apache Parquet. Esta herramienta soporta de manera nativa la lectura y escritura en ficheros Parquet. Leyendo Parquet, las columnas se marcan como nullable por razones de compatibilidad.
Por otra parte, Apache Parquet permite la partición de datos para optimizar sus lecturas y los datos se almacenan en distintos directorios, codificando los valores en el nombre del directorio para cada partición del SparkSQL.

SparkSQL: Apache Avro
Spark SQL en Apache Spark.SQL necesita de una extensión para poder trabajar con Avro y poder hacer .format(«avro”).
libraryDependencies += «org.apache.spark»
%% «spark-avro» % «3.0.1»
Avro permite indicar el tema por medio de options a la hora de hacer los read/write, mediante la propiedad avroSchema.

Por otra parte, Avro también tiene tipos lógicos para soportar estructuras de datos más complejas.

Por último, el tipo unión en Avro se utiliza para indicar que un campo puede tener distintos valores de tipo. Spark.SQL realiza las siguientes traducciones:
- Unión (int, long) a LongType.
- Unión(float, double) a DoubleType.
- Unión (any, null) a Nullable con Spark.SQL.
- Si la unión es de dos tipos distintos, son consideradas complejas y se traducen en una estructura con distintos miembros, así: union(int, string) à member0 (int), member1 (string).
Conoce más del Big Data con KeepCoding
Por medio de este post, has podido identificar qué es Spark SQL en Apache Spark. No obstante, esto exige continuar practicando para ganar experiencia. Si no tienes claro cómo puedes empezar, ¡en KeepCoding te ofrecemos la mejor opción!
Nuestro Bootcamp Full Stack Big Data, Inteligencia Artificial & Machine Learning cuenta con once módulos que te prepararán y pondrán a prueba tus destrezas con las principales herramientas desarrolladas para el procesamiento de los macrodatos. Para ello, también contarás con el apoyo de una serie de expertos en Big Data que te guiarán en los procesos tanto teóricos como prácticos. ¡No esperes más e inscríbete ahora!
