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
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
Propriétés
Propriété | Description |
|---|---|
La version de Spark sur laquelle cette application est exécutée. | |
Interface de configuration Runtime pour Spark. | |
Interface par laquelle l'utilisateur peut créer, supprimer, modifier ou query les bases de données, tables, fonctions sous-jacentes, etc. | |
Renvoie un UDFRegistration pour l'enregistrement UDF. | |
Renvoie une UDTFRegistration pour l'enregistrement UDTF. | |
Renvoie une DataSourceRegistration pour l'enregistrement de la source de données. | |
Renvoie un Profile pour l'analyse des performances et de la mémoire. | |
Renvoie le SparkContext sous-jacent. Mode classique uniquement. | |
Renvoie un DataFrameReader qui peut être utilisé pour lire des données en tant que DataFrame. | |
Renvoie un DataStreamReader qui peut être utilisé pour lire des flux de données sous la forme d'un DataFrame en streaming. | |
Renvoie un StreamingQueryManager qui permet de gérer toutes les requêtes de streaming actives. | |
Renvoie une TableValuedFunction pour l'appel de fonctions à valeur de table (TVF). |
Méthodes
Méthode | Description |
|---|---|
Crée un DataFrame à partir d'un RDD, d'une liste, d'un DataFrame pandas, d'un ndarray numpy ou d'une Table pyarrow. | |
Retourne un DataFrame représentant le résultat de la query donnée. | |
Renvoie la table spécifiée en tant que DataFrame. | |
Crée un DataFrame avec une seule colonne de type LongType nommée | |
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. | |
Renvoie la SparkSession active pour le thread actuel. | |
Renvoie la SparkSession active ou par default pour le thread actuel. | |
Arrête le SparkContext sous-jacent. | |
Ajoute des artefacts à la session client. | |
Interrompt toutes les opérations de cette session en cours d'exécution sur le serveur. | |
Interrompt toutes les opérations de cette session avec le tag donné. | |
Interrompt une opération de cette session avec l'operationId donné. | |
Ajoute un tag à attribuer à toutes les opérations démarrées par ce fil dans cette session. | |
Supprime une balise précédemment ajoutée pour les opérations lancées par ce fil. | |
Récupère les tags actuellement définis à attribuer à toutes les Opérations start initiées par ce fil. | |
Efface les tags d'opération du thread actuel. |
Générateur
Méthode | Description |
|---|---|
| Définit une option de configuration. Les options sont automatiquement propagées à SparkConf et aux propres configurations de SparkSession. |
| Définit l'URL maître Spark à laquelle se connecter. |
| Définit l'URL distante Spark pour se connecter via Spark Connect. |
| Définit un nom pour l'application, qui sera affiché dans l'interface utilisateur web de Spark. |
| Active le support Hive, y compris la connectivité à un Hive metastore persistant. |
| 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. |
| Crée une nouvelle SparkSession. |
Exemples
spark = (
SparkSession.builder
.master("local")
.appName("Word Count")
.config("spark.some.config.option", "some-value")
.getOrCreate()
)
spark.sql("SELECT * FROM range(10) where id > 7").show()
+---+
| id|
+---+
| 8|
| 9|
+---+
spark.createDataFrame([('Alice', 1)], ['name', 'age']).show()
+-----+---+
| name|age|
+-----+---+
|Alice| 1|
+-----+---+
spark.range(1, 7, 2).show()
+---+
| id|
+---+
| 1|
| 3|
| 5|
+---+