SparkSession

Point d’entrée pour programmer Spark avec le jeu de données et l’API DataFrame. Une session SparkSession peut être utilisée pour créer des DataFrames, inscrire des DataFrames en tant que tables, exécuter SQL sur des tables, mettre en cache des tables et lire des fichiers Parquet.

Syntaxe

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Propriétés

Propriété Description
builder Interface pour construire la configuration de la session.
catalog Interface via laquelle l’utilisateur peut créer, supprimer, modifier ou interroger des bases de données sous-jacentes, des tables, des fonctions, etc.
client Cela donne accès au client Spark Connect. Spark Connect uniquement.
conf Interface de configuration du runtime pour Spark.
dataSource Retourne une DataSourceRegistration pour l’inscription de la source de données.
profile Retourne un profil pour le profilage des performances/de la mémoire.
read Retourne un DataFrameReader qui peut être utilisé pour lire des données en tant que DataFrame.
readStream Retourne un DataStreamReader qui peut être utilisé pour lire des flux de données en tant que DataFrame de streaming.
sparkContext Retourne le sparkContext sous-jacent. Mode classique uniquement.
streams Retourne un StreamingQueryManager qui permet de gérer toutes les requêtes de streaming actives.
tvf Retourne une tableValuedFunction pour appeler des fonctions table (TVFs).
udf Retourne une UDFRegistration pour l’inscription UDF.
udtf Retourne une valeur UDTFRegistration pour l’inscription UDTF.
version Version de Spark sur laquelle cette application s’exécute.

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 requête donnée.
table(tableName) Retourne la table spécifiée en tant que DataFrame.
range(start, end, step, numPartitions) Crée un DataFrame avec une seule colonne LongType nommée id, contenant des éléments dans une plage.
newSession() Retourne une nouvelle session SparkSession avec sqlConf distinct, des vues temporaires inscrites et des fonctions définies par l’utilisateur, mais un cache sparkContext et un cache de table partagés. Mode classique uniquement.
getActiveSession() Retourne la session SparkSession active pour le thread actif.
active() Retourne la session SparkSession active ou par défaut pour le thread actif.
stop() Arrête le SparkContext sous-jacent.
addArtifacts(*path, pyfile, archive, file) Ajoute des artefacts à la session cliente.
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 la balise donnée.
interruptOperation(op_id) Interrompt une opération de cette session avec l’id d’opération donné.
addTag(tag) Ajoute une balise à affecter à toutes les opérations démarrées par ce thread dans cette session.
removeTag(tag) Supprime une balise précédemment ajoutée pour les opérations démarrées par ce thread.
getTags() Obtient les balises actuellement définies pour être affectées à toutes les opérations démarrées par ce thread.
clearTags() Efface les balises d’opération du thread actuel.