Petit pipeline rapide
Avant d'analyser des données plus complexes, votre gestionnaire souhaite voir un exemple simple de pipeline avec les étapes de base. Pour cet exemple, vous allez ingérer un fichier de données, filtrer quelques lignes, ajouter une colonne d'identifiant, puis l'écrire en données JSON.
Le contexte spark est défini, et la bibliothèque pyspark.sql.functions est importée avec l'alias F, comme c'est l'usage.
Cette activité fait partie du cours
Nettoyer des données avec PySpark
Instructions de l’exercice
- Importez le fichier
2015-departures.csv.gzdans un DataFrame. Notez que l'en-tête est déjà défini. - Filtrez le DataFrame pour ne conserver que les vols d'une durée supérieure à 0 minute. Utilisez l'index de la colonne, et non son nom (n'oubliez pas d'utiliser
.printSchema()pour voir les noms et l'ordre des colonnes). - Ajoutez une colonne d'identifiant.
- Écrivez le fichier sous forme de document JSON nommé
output.json.
Exercice interactif pratique
Essayez cet exercice en complétant ce code d’exemple.
# Import the data to a DataFrame
departures_df = spark.____(____, header=____)
# Remove any duration of 0
departures_df = departures_df.____(____)
# Add an ID column
departures_df = departures_df.____('id', ____)
# Write the file out to JSON format
____.write.____(____, mode='overwrite')