CommencezCommencez gratuitement

Convertir une fonction de fenêtre de la notation par points vers SQL

Nous allons ajouter une colonne à un horaire de train pour que chaque rangée indique le nombre de minutes nécessaires au train pour atteindre son prochain arrêt.

  • Nous avons un dataframe dfdf.columns == ['train_id', 'station', 'time'].
  • df est enregistré comme table SQL nommée schedule.
  • La fonction de fenêtre suivante utilise la notation par points. Elle produit un nouveau dataframe dot_df.
window = Window.partitionBy('train_id').orderBy('time')
dot_df = df.withColumn('diff_min', 
                    (unix_timestamp(lead('time', 1).over(window),'H:m') 
                     - unix_timestamp('time', 'H:m'))/60)

Notez l'utilisation de la fonction unix_timestamp, qui est l'équivalent de la fonction SQL UNIX_TIMESTAMP.

Veuillez tenir compte de l'échafaudage dans l'exemple de code. Si vous formatez la réponse en fonction de cet échafaudage, votre soumission ne sera pas rejetée par erreur à cause d'un problème de formatage.

Cette activité fait partie du cours

Introduction à Spark SQL en Python

Voir le cours

Instructions de l’exercice

  • Créez une requête SQL qui produit un résultat identique à dot_df. Veuillez formater la requête selon l'échafaudage (c.-à-d. les traits de soulignement _____).

Exercice interactif pratique

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

# Create a SQL query to obtain an identical result to dot_df
query = """
SELECT *, 
(____(____(time, 1) ____ (____ BY train_id ____ BY time),'H:m') 
 - ____(time, 'H:m'))/60 AS diff_min 
FROM schedule 
"""
sql_df = spark.sql(query)
sql_df.show()
Modifier et exécuter le code