CommencerCommencez gratuitement

Création de colonnes

Dans ce chapitre, vous apprendrez à utiliser les méthodes définies par la classe DataFrame de Spark pour effectuer des opérations courantes sur les données.

Examinons les opérations effectuées colonne par colonne. Dans Spark, vous pouvez réaliser cela en utilisant la méthode .withColumn(), qui prend deux arguments. Tout d'abord, une chaîne contenant le nom de votre nouvelle colonne, puis la nouvelle colonne elle-même.

La nouvelle colonne doit être un objet de la classe Column. Pour créer l'un de ces éléments, il suffit d'extraire une colonne de votre DataFrame avec df.colName.

La mise à jour d'un Spark DataFrame diffère quelque peu du travail dans un pandas, car le Spark DataFrame est immuable. Cela signifie qu'il ne peut pas être modifié, et donc que les colonnes ne peuvent pas être mises à jour sur place.

Par conséquent, toutes ces méthodes renvoient un nouveau DataFrame. Pour remplacer le DataFrame d'origine, il est nécessaire de réaffecter le DataFrame renvoyé à l'aide de la méthode suivante :

df = df.withColumn("newCol", df.oldCol + 1)

Le code ci-dessus crée un DataFrame avec les mêmes colonnes que df, ainsi qu'une nouvelle colonne, newCol, où chaque entrée est égale à l'entrée correspondante de oldCol, plus un.

Pour remplacer une colonne existante, veuillez simplement transmettre le nom de la colonne en tant que premier argument.

N’oubliez pas : un SparkSession nommé spark se trouve déjà dans votre espace de travail.

Cet exercice fait partie du cours

<cours>Principes fondamentaux de PySpark</cours>
Voir le cours

Instructions de l’exercice

  • Utilisez la méthode spark.table() avec l'argument "flights" pour créer un DataFrame contenant les valeurs de la table flights dans .catalog. Sauvegardez le résultat sous flights.
  • Veuillez afficher le head de flights en utilisant flights.show(). Vérifiez la sortie : la colonne air_time contient la durée du vol en minutes.
  • Mettez à jour flights afin d'y inclure une nouvelle colonne intitulée duration_hrs, qui contient la durée de chaque vol en heures (il sera nécessaire de diviser air_time par le nombre de minutes dans une heure).

Exercice interactif pratique

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

# Create the DataFrame flights
flights = spark.table(____)

# Show the head
____.____()

# Add duration_hrs
flights = flights.withColumn(____)
Modifier et exécuter le code