CommencezCommencez gratuitement

Ajoutez un peu de Spark à vos données

Dans l'exercice précédent, vous avez vu comment transférer des données de Spark vers pandas. Mais vous voudrez peut-être faire l'inverse et charger un DataFrame pandas dans un cluster Spark! La classe SparkSession offre aussi une méthode pour ça.

La méthode .createDataFrame() prend un DataFrame pandas et retourne un DataFrame Spark.

La sortie de cette méthode est stockée localement, et non dans le catalogue de SparkSession. Cela signifie que vous pouvez utiliser toutes les méthodes de DataFrame Spark dessus, mais vous ne pouvez pas y accéder dans d'autres contextes.

Par exemple, une requête SQL (au moyen de la méthode .sql()) qui fait référence à votre DataFrame lèvera une erreur. Pour y accéder de cette façon, vous devez l'enregistrer comme table temporaire.

Vous pouvez le faire avec la méthode de DataFrame Spark .createTempView(), qui prend comme seul argument le nom de la table temporaire que vous souhaitez enregistrer. Cette méthode inscrit le DataFrame comme table dans le catalogue, mais comme il s'agit d'une table temporaire, elle n'est accessible qu'à partir de la SparkSession précise qui a servi à créer le DataFrame Spark.

Il existe aussi la méthode .createOrReplaceTempView(). Elle crée sans risque une nouvelle table temporaire si aucune n'existe, ou met à jour une table existante si elle est déjà définie. Vous utiliserez cette méthode pour éviter les problèmes de tables en double.

Consultez le schéma pour voir toutes les façons dont vos structures de données Spark interagissent entre elles.

Une SparkSession nommée spark est déjà disponible dans votre espace de travail, numpy a été importé sous np et pandas sous pd.

Cette activité fait partie du cours

Fondements de PySpark

Voir le cours

Instructions de l’exercice

  • Le code pour créer un DataFrame pandas de nombres aléatoires a déjà été fourni et est enregistré sous pd_temp.
  • Créez un DataFrame Spark appelé spark_temp en appelant la méthode Spark .createDataFrame() avec pd_temp comme argument.
  • Examinez la liste des tables de votre cluster Spark et vérifiez que le nouveau DataFrame n'y figure pas. Rappelez-vous que vous pouvez utiliser spark.catalog.listTables() pour ce faire.
  • Enregistrez le DataFrame spark_temp que vous venez de créer comme table temporaire à l'aide de la méthode .createOrReplaceTempView(). La table temporaire doit s'appeler "temp". N'oubliez pas que le nom de la table est défini en l'incluant comme unique argument de votre méthode!
  • Examinez de nouveau la liste des tables.

Exercice interactif pratique

Essayez cet exercice en complétant ce code d’exemple.

# Create pd_temp
pd_temp = pd.DataFrame(np.random.random(10))

# Create spark_temp from pd_temp
spark_temp = ____

# Examine the tables in the catalog
print(____)

# Add spark_temp to the catalog
spark_temp.____

# Examine the tables in the catalog again
print(____)
Modifier et exécuter le code