Aller au contenu principal

SparkSession

Le point d'entrée pour la programmation Spark avec l'API Dataset et DataFrame. Une SparkSession peut être utilisée pour créer des DataFrames, enregistrer des DataFrames comme tables, exécuter des requêtes SQL sur des tables, mettre en cache des tables et lire des fichiers Parquet.

Syntaxe

Python
from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Propriétés

Propriété

Description

version

La version de Spark sur laquelle cette application est exécutée.

conf

Interface de configuration Runtime pour Spark.

catalog

Interface par laquelle l'utilisateur peut créer, supprimer, modifier ou query les bases de données, tables, fonctions sous-jacentes, etc.

udf

Renvoie un UDFRegistration pour l'enregistrement UDF.

udtf

Renvoie une UDTFRegistration pour l'enregistrement UDTF.

dataSource

Renvoie une DataSourceRegistration pour l'enregistrement de la source de données.

profile

Renvoie un Profile pour l'analyse des performances et de la mémoire.

sparkContext

Renvoie le SparkContext sous-jacent. Mode classique uniquement.

read

Renvoie un DataFrameReader qui peut être utilisé pour lire des données en tant que DataFrame.

readStream

Renvoie un DataStreamReader qui peut être utilisé pour lire des flux de données sous la forme d'un DataFrame en streaming.

streams

Renvoie un StreamingQueryManager qui permet de gérer toutes les requêtes de streaming actives.

tvf

Renvoie une TableValuedFunction pour l'appel de fonctions à valeur de table (TVF).

Propriété

Description

version

La version de Spark sur laquelle cette application est exécutée.

conf

Interface de configuration Runtime pour Spark.

catalog

Interface par laquelle l'utilisateur peut créer, supprimer, modifier ou query les bases de données, tables, fonctions sous-jacentes, etc.

udf

Renvoie un UDFRegistration pour l'enregistrement UDF.

udtf

Renvoie une UDTFRegistration pour l'enregistrement UDTF.

dataSource

Renvoie une DataSourceRegistration pour l'enregistrement de la source de données.

profile

Renvoie un Profile pour l'analyse des performances et de la mémoire.

sparkContext

Renvoie le SparkContext sous-jacent. Mode classique uniquement.

read

Renvoie un DataFrameReader qui peut être utilisé pour lire des données en tant que DataFrame.

readStream

Renvoie un DataStreamReader qui peut être utilisé pour lire des flux de données sous la forme d'un DataFrame en streaming.

streams

Renvoie un StreamingQueryManager qui permet de gérer toutes les requêtes de streaming actives.

tvf

Renvoie une TableValuedFunction pour l'appel de fonctions à valeur de table (TVF).

Méthodes

Méthode

Description

createDataFrame(data, schema, samplingRatio, verifySchema)

Crée un DataFrame à partir d'un RDD, d'une liste, d'un DataFrame pandas, d'un ndarray numpy ou d'une Table pyarrow.

sql(sqlQuery, args, **kwargs)

Retourne un DataFrame représentant le résultat de la query donnée.

table(tableName)

Renvoie la table spécifiée en tant que DataFrame.

range(start, end, step, numPartitions)

Crée un DataFrame avec une seule colonne de type LongType nommée id, contenant des éléments dans une plage.

newSession()

Renvoie une nouvelle SparkSession avec des SQLConf, des vues temporaires enregistrées et des UDF séparées, mais un SparkContext et un cache de table partagés. Mode classique uniquement.

getActiveSession()

Renvoie la SparkSession active pour le thread actuel.

active()

Renvoie la SparkSession active ou par default pour le thread actuel.

stop()

Arrête le SparkContext sous-jacent.

addArtifacts(*path, pyfile, archive, file)

Ajoute des artefacts à la session client.

interruptAll()

Interrompt toutes les opérations de cette session en cours d'exécution sur le serveur.

interruptTag(tag)

Interrompt toutes les opérations de cette session avec le tag donné.

interruptOperation(op_id)

Interrompt une opération de cette session avec l'operationId donné.

addTag(tag)

Ajoute un tag à attribuer à toutes les opérations démarrées par ce fil dans cette session.

removeTag(tag)

Supprime une balise précédemment ajoutée pour les opérations lancées par ce fil.

getTags()

Récupère les tags actuellement définis à attribuer à toutes les Opérations start initiées par ce fil.

clearTags()

Efface les tags d'opération du thread actuel.

Méthode

Description

createDataFrame(data, schema, samplingRatio, verifySchema)

Crée un DataFrame à partir d'un RDD, d'une liste, d'un DataFrame pandas, d'un ndarray numpy ou d'une Table pyarrow.

sql(sqlQuery, args, **kwargs)

Retourne un DataFrame représentant le résultat de la query donnée.

table(tableName)

Renvoie la table spécifiée en tant que DataFrame.

range(start, end, step, numPartitions)

Crée un DataFrame avec une seule colonne de type LongType nommée id, contenant des éléments dans une plage.

newSession()

Renvoie une nouvelle SparkSession avec des SQLConf, des vues temporaires enregistrées et des UDF séparées, mais un SparkContext et un cache de table partagés. Mode classique uniquement.

getActiveSession()

Renvoie la SparkSession active pour le thread actuel.

active()

Renvoie la SparkSession active ou par default pour le thread actuel.

stop()

Arrête le SparkContext sous-jacent.

addArtifacts(*path, pyfile, archive, file)

Ajoute des artefacts à la session client.

interruptAll()

Interrompt toutes les opérations de cette session en cours d'exécution sur le serveur.

interruptTag(tag)

Interrompt toutes les opérations de cette session avec le tag donné.

interruptOperation(op_id)

Interrompt une opération de cette session avec l'operationId donné.

addTag(tag)

Ajoute un tag à attribuer à toutes les opérations démarrées par ce fil dans cette session.

removeTag(tag)

Supprime une balise précédemment ajoutée pour les opérations lancées par ce fil.

getTags()

Récupère les tags actuellement définis à attribuer à toutes les Opérations start initiées par ce fil.

clearTags()

Efface les tags d'opération du thread actuel.

Générateur

Méthode

Description

config(key, value)

Définit une option de configuration. Les options sont automatiquement propagées à SparkConf et aux propres configurations de SparkSession.

master(master)

Définit l'URL maître Spark à laquelle se connecter.

remote(url)

Définit l'URL distante Spark pour se connecter via Spark Connect.

appName(name)

Définit un nom pour l'application, qui sera affiché dans l'interface utilisateur web de Spark.

enableHiveSupport()

Active le support Hive, y compris la connectivité à un Hive metastore persistant.

getOrCreate()

Obtient une SparkSession existante ou, s'il n'y en a pas, en crée une nouvelle basée sur les options définies dans ce générateur.

create()

Crée une nouvelle SparkSession.

Méthode

Description

config(key, value)

Définit une option de configuration. Les options sont automatiquement propagées à SparkConf et aux propres configurations de SparkSession.

master(master)

Définit l'URL maître Spark à laquelle se connecter.

remote(url)

Définit l'URL distante Spark pour se connecter via Spark Connect.

appName(name)

Définit un nom pour l'application, qui sera affiché dans l'interface utilisateur web de Spark.

enableHiveSupport()

Active le support Hive, y compris la connectivité à un Hive metastore persistant.

getOrCreate()

Obtient une SparkSession existante ou, s'il n'y en a pas, en crée une nouvelle basée sur les options définies dans ce générateur.

create()

Crée une nouvelle SparkSession.

Exemples

Python
spark = (
SparkSession.builder
.master("local")
.appName("Word Count")
.config("spark.some.config.option", "some-value")
.getOrCreate()
)
Python
spark.sql("SELECT * FROM range(10) where id > 7").show()
Output
+---+
| id|
+---+
| 8|
| 9|
+---+
Python
spark.createDataFrame([('Alice', 1)], ['name', 'age']).show()
Output
+-----+---+
| name|age|
+-----+---+
|Alice| 1|
+-----+---+
Python
spark.range(1, 7, 2).show()
Output
+---+
| id|
+---+
| 1|
| 3|
| 5|
+---+