Aller au contenu principal

fenêtre

Ventilez les lignes en une ou plusieurs fenêtres temporelles étant donné une colonne spécifiant le Timestamp. Les débuts de fenêtre sont inclusifs, mais les fins de fenêtre sont exclusives, par ex. : 12:05 sera dans la fenêtre [12:05,12:10) mais pas dans [12:00,12:05). Windows peut prendre en charge une précision à la microseconde. Windows dans l'ordre des mois ne sont pas prises en charge.

La colonne d'heure doit être de pyspark.sql.types.TimestampType.

Les durées sont fournies sous forme de chaînes, par exemple. '1 seconde', '1 jour 12 heures', '2 minutes'. Les chaînes d'intervalle valides sont 'semaine', 'jour', 'heure', 'minute', 'seconde', 'milliseconde', 'microseconde'. Si le slideDuration n’est pas fourni, les fenêtres seront des fenêtres basculantes.

Le startTime est le décalage par rapport au 1970-01-01 00:00:00 UTC avec lequel start les intervalles de fenêtre. Par exemple, pour avoir des fenêtres glissantes horaires qui start 15 minutes après l'heure, par exemple. 12 :15 - 13 :15, 13 :15 - 14 :15... fournissez startTime en tant que 15 minutes.

La colonne de sortie sera une structure appelée « window » par default avec les colonnes imbriquées « start » et « end », où « start » et « end » seront de pyspark.sql.types.TimestampType.

Pour la fonction Databricks SQL correspondante, consultez l'expression de regroupementwindow.

Syntaxe

Python
from pyspark.sql import functions as dbf

dbf.window(timeColumn=<timeColumn>, windowDuration=<windowDuration>, slideDuration=<slideDuration>, startTime=<startTime>)

parameter

parameter

Type

Description

timeColumn

pyspark.sql.Column OU str

La colonne ou l'expression à utiliser comme timestamp pour le fenêtrage temporel. La colonne d'horodatage doit être de type TimestampType ou TimestampNTZType.

windowDuration

literal string

Une chaîne spécifiant la largeur de la fenêtre, par exemple, 10 minutes, 1 second. Vérifiez org.apache.spark.unsafe.types.CalendarInterval pour des identifiants de durée valides. Notez que la durée est une période de temps fixe et ne varie pas au fil du temps selon un calendrier. Par exemple, 1 day signifie toujours 86 400 000 millisecondes, et non un jour calendaire.

slideDuration

literal string, optional

Une nouvelle fenêtre sera générée tous les slideDuration. Doit être inférieur ou égal au windowDuration. Vérifiez org.apache.spark.unsafe.types.CalendarInterval pour des identifiants de durée valides. Cette durée est également absolue et ne varie pas selon un calendrier.

startTime

literal string, optional

Le décalage par rapport au 1970-01-01 00:00:00 UTC pour start les intervalles de fenêtre. Par exemple, pour avoir des fenêtres basculantes horaires qui start 15 minutes après l'heure, par exemple. 12 :15 - 13 :15, 13 :15 - 14 :15... fournissez startTime en tant que 15 minutes.

parameter

Type

Description

timeColumn

pyspark.sql.Column OU str

La colonne ou l'expression à utiliser comme timestamp pour le fenêtrage temporel. La colonne d'horodatage doit être de type TimestampType ou TimestampNTZType.

windowDuration

literal string

Une chaîne spécifiant la largeur de la fenêtre, par exemple, 10 minutes, 1 second. Vérifiez org.apache.spark.unsafe.types.CalendarInterval pour des identifiants de durée valides. Notez que la durée est une période de temps fixe et ne varie pas au fil du temps selon un calendrier. Par exemple, 1 day signifie toujours 86 400 000 millisecondes, et non un jour calendaire.

slideDuration

literal string, optional

Une nouvelle fenêtre sera générée tous les slideDuration. Doit être inférieur ou égal au windowDuration. Vérifiez org.apache.spark.unsafe.types.CalendarInterval pour des identifiants de durée valides. Cette durée est également absolue et ne varie pas selon un calendrier.

startTime

literal string, optional

Le décalage par rapport au 1970-01-01 00:00:00 UTC pour start les intervalles de fenêtre. Par exemple, pour avoir des fenêtres basculantes horaires qui start 15 minutes après l'heure, par exemple. 12 :15 - 13 :15, 13 :15 - 14 :15... fournissez startTime en tant que 15 minutes.

Renvoie

pyspark.sql.Column: la colonne pour les résultats calculés.

Exemples

Python
import datetime
from pyspark.sql import functions as dbf
df = spark.createDataFrame([(datetime.datetime(2016, 3, 11, 9, 0, 7), 1)], ['dt', 'v'])
df2 = df.groupBy(dbf.window('dt', '5 seconds')).agg(dbf.sum('v'))
df2.show(truncate=False)
df2.printSchema()