SparkSession

Точка входа для программирования Spark с помощью API набора данных и кадра данных. SparkSession можно использовать для создания кадров данных, регистрации кадров данных в качестве таблиц, выполнения SQL над таблицами, таблиц кэша и чтения файлов parquet.

Синтаксис

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Свойства

Недвижимость Описание
builder Интерфейс для построения конфигурации сессии.
catalog Интерфейс, с помощью которого пользователь может создавать, удалять, изменять или запрашивать базовые базы данных, таблицы, функции и т. д.
client Даёт доступ к клиенту Spark Connect. Только Spark Connect.
conf Интерфейс конфигурации среды выполнения для Spark.
dataSource Возвращает dataSourceRegistration для регистрации источника данных.
profile Возвращает профиль для профилирования производительности и памяти.
read Возвращает dataFrameReader, который можно использовать для чтения данных в виде кадра данных.
readStream Возвращает DataStreamReader, который можно использовать для чтения потоковых потоков данных в качестве кадра потоковой передачи.
sparkContext Возвращает базовый SparkContext. Только классический режим.
streams Возвращает streamQueryManager, который позволяет управлять всеми активными запросами потоковой передачи.
tvf Возвращает значение TableValuedFunction для вызова табличных функций (TVFs).
udf Возвращает UDFRegistration для регистрации UDF.
udtf Возвращает UDTFRegistration для регистрации UDTF.
version Версия Spark, в которой работает это приложение.

Методы

Метод Описание
createDataFrame(data, schema, samplingRatio, verifySchema) Создает кадр данных из RDD, списка, кадра данных pandas, число ndarray или таблицы pyarrow.
sql(sqlQuery, args, **kwargs) Возвращает кадр данных, представляющий результат заданного запроса.
table(tableName) Возвращает указанную таблицу в виде кадра данных.
range(start, end, step, numPartitions) Создает кадр данных с одним столбцом LongType с именем id, содержащими элементы в диапазоне.
newSession() Возвращает новую версию SparkSession с отдельными SQLConf, зарегистрированными временными представлениями и пользовательскими файлами, но общим кэшем SparkContext и таблиц. Только классический режим.
getActiveSession() Возвращает активное решение SparkSession для текущего потока.
active() Возвращает активное или стандартное решение SparkSession для текущего потока.
stop() Останавливает базовый SparkContext.
addArtifacts(*path, pyfile, archive, file) Добавляет артефакты в сеанс клиента.
interruptAll() Прерывает все операции этого сеанса, запущенного на сервере.
interruptTag(tag) Прерывает все операции этого сеанса с заданным тегом.
interruptOperation(op_id) Прерывает операцию этого сеанса с заданным идентификатором операции.
addTag(tag) Добавляет тег, назначенный всем операциям, запущенным этим потоком в этом сеансе.
removeTag(tag) Удаляет тег, добавленный ранее для операций, запущенных этим потоком.
getTags() Возвращает теги, которые в настоящее время присваиваются всем операциям, запущенным этим потоком.
clearTags() Очищает теги операций текущего потока.