Aller au contenu principal

Fenêtres glissantes dans Structured Streaming

Une fenêtre glissante compute des agrégations temporelles sur une frontière qui avance en continu. Contrairement aux fenêtres fixes ou glissantes qui attribuent chaque événement à un compartiment temporel fixe, une fenêtre glissante émet des agrégats mis à jour à mesure que la frontière avance et que les événements plus anciens sortent de la plage. Utilisez des fenêtres glissantes pour le feature serving en temps réel et le monitoring à faible latence, lorsque vous avez besoin de la valeur actuelle d’un agrégat sur une période de rétrospection récente.

Les fenêtres glissantes prennent en charge les éléments suivants :

  • Fenêtres bornées qui agrègent sur une durée de rétrospective fixe, telle que les 10 dernières minutes.
  • Fenêtres de durée de vie qui cumulent toutes les données depuis le démarrage du stream.
  • Fenêtres multiples dans un seul opérateur , ce qui vous permet de compute différentes durées rétrospectives en un seul appel.

Configuration requise​

  • Databricks Runtime 19 et versions supérieures.
  • Mode micro-batch : compute classique (mode d'accès Standard ou Dedicated) ou compute serverless.
  • Mode temps réel : compute classique (mode d'accès standard ou dédié). Dans Lakeflow Pipelines, le mode temps réel s'exécute également sur un compute Serverless via la configuration du pipeline. Voir Utiliser le mode temps réel dans les pipelines LakeFlow Pipelines.

Pour plus d’information, voir Modes d'exécution.

Fonctionnement des fenêtres glissantes​

Une query avec fenêtre glissante regroupe les événements par une ou plusieurs colonnes de partition, les trie par colonne de Timestamp et compute des agrégats par rapport à une frontière. Chaque ligne de sortie inclut la valeur de la frontière en tant que colonne de timestamp.

Frontier​

La frontière définit le « now » actuel pour l’opérateur. Il avance en fonction du temps de traitement (horloge), et peut éventuellement accuser un retard d’une valeur configurable delay. Un délai est utile pour absorber de légers retards d’ingestion afin que les événements en cours de traitement soient comptabilisés avant que la frontière ne les dépasse.

Définissez la frontière avec FrontierSpec:

Python
# Frontier with a 5-second delay, emitted as a column named "window_time"
RollingWindow.FrontierSpec(delay="5 seconds", alias="window_time")

Fenêtres délimitées​

Une fenêtre bornée agrège les données sur une durée fixe, en remontant à partir de la frontière. Définir la période rétrospective avec .over():

Python
# Sum of "value" over the last hour
F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("1 hour")))

Fenêtres de durée de vie​

A lifetime window accumulates data from the start of the stream with no expiration. Pour créer une fenêtre de durée de vie, omettez .over():

Python
# Running count of non-null values since the stream started
F.count("value")

Modes d'exécution​

Les fenêtres glissantes s’exécutent à la fois en mode micro-batch et en mode temps réel. La frontière se comporte différemment dans chacun d’eux :

  • Mode micro-batch : La frontière est fixe pour chaque micro-batch. Chaque clé de partition émet au plus une ligne de sortie par batch.
  • Mode temps réel : La frontière progresse de manière continue pendant un batch. Une clé de partition peut émettre plusieurs lignes de sortie au sein d’un même batch à mesure que la frontière se déplace.

Événements tardifs et filigranes​

Une fenêtre glissante définit la limite de la fenêtre à partir de la frontière du temps de traitement, mais vous pouvez appliquer un repère temporel d’événement pour ignorer les événements arrivant en retard. Ajoutez .withWatermark() avant rollingWindow() en utilisant la même colonne que celle transmise à orderBy:

Python
df.withWatermark("event_time", "5 seconds").rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="0 seconds"),
measures=[F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("10 minutes")))],
)

Un événement dont l'heure de l'événement est inférieure ou égale au filigrane est ignoré, même lorsqu'il se situe dans la plage horaire de la fenêtre. Sans filigrane, chaque événement se trouvant dans la fenêtre est agrégé, y compris les arrivées tardives. Les filigranes s'appliquent en mode micro-batch et en temps réel, et la position du filigrane est enregistrée par point de contrôle, de sorte qu'elle est restaurée lorsqu'une query redémarre.

Syntaxe​

Appelez rollingWindow() sur un DataFrame en streaming :

Python
from pyspark.sql.streaming.rolling_window import RollingWindow
from pyspark.sql import functions as F

df.rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(
delay="0 seconds",
alias="frontier",
),
measures=[
F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("10 minutes"))),
F.count("value").alias("non_null_value_count"), # No .over() = lifetime window
],
)

L'opérateur rollingWindow() prend les paramètres suivants :

parameter

Description

partitionBy

Obligatoire. Une ou plusieurs colonnes sur lesquelles regrouper les événements.

orderBy

Obligatoire. La colonne de timestamp qui ordonne les événements. Doit être une colonne TimestampType.

frontierSpec

La définition de la frontière. Définissez delay pour contrôler le retard de la frontière par rapport au temps de traitement, et alias pour nommer la colonne de la frontière de sortie.

measures

Une liste d'expressions d'agrégation. Ajoutez .over() pour une fenêtre bornée, ou omettez-le pour une fenêtre à vie.

parameter

Description

partitionBy

Obligatoire. Une ou plusieurs colonnes sur lesquelles regrouper les événements.

orderBy

Obligatoire. La colonne de timestamp qui ordonne les événements. Doit être une colonne TimestampType.

frontierSpec

La définition de la frontière. Définissez delay pour contrôler le retard de la frontière par rapport au temps de traitement, et alias pour nommer la colonne de la frontière de sortie.

measures

Une liste d'expressions d'agrégation. Ajoutez .over() pour une fenêtre bornée, ou omettez-le pour une fenêtre à vie.

Fonctions d’agrégation prises en charge​

Les fenêtres glissantes prennent en charge les fonctions d’agrégation suivantes :

  • sum
  • count
  • avg
  • min
  • max
  • stddev, y compris les variantes d'échantillon (stddev_samp) et de population (stddev_pop)
  • variance, y compris les variantes d'échantillon (var_samp) et de population (var_pop)
  • approx_count_distinct
  • approx_percentile
  • UDAFs Scala Aggregator enregistrées avec functions.udaf. Pris en charge uniquement sur le compute classique. Pour obtenir un exemple, consultez Example 3: Scala UDAF.

Exemples​

Exemple 1 : Revenu cumulé glissant par utilisateur au cours de la dernière heure​

L’exemple suivant compute une somme mobile des revenus par utilisateur sur la dernière heure et écrit les mises à jour dans Kafka :

Python
from pyspark.sql.streaming.rolling_window import RollingWindow
from pyspark.sql import functions as F

result = events_df.rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="0 seconds"),
measures=[
F.sum("revenue")
.over(RollingWindow.preceding(RollingWindow.Range("1 hour")))
.alias("rolling_revenue"),
],
)

(
result.select(
F.col("user_id").cast("string").alias("key"),
F.to_json(F.struct("__frontier", "rolling_revenue")).alias("value"),
)
.writeStream.format("kafka")
.option("kafka.bootstrap.servers", "<host1:port1,host2:port2>")
.option("topic", "<topic-name>")
.outputMode("update")
.start()
)

Exemple 2 : Fenêtres mixtes limitées et de durée de vie​

L'exemple suivant compute deux agrégats bornés et un agrégat de durée de vie dans un seul opérateur :

Python
result = events_df.rollingWindow(
partitionBy=["region", "product"],
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="10 seconds", alias="window_ts"),
measures=[
F.sum("revenue").over(RollingWindow.preceding(RollingWindow.Range("1 hour"))),
F.avg("price").over(RollingWindow.preceding(RollingWindow.Range("10 minutes"))),
F.count("value"), # Lifetime: non-null values since the stream started
],
)

Le schéma des résultats est (region, product, window_ts, sum(revenue), avg(price), count(value)).

Exemple 3 : Scala UDAF​

remarque

Scala Aggregator UDAFs are supported only on classic compute.

L'exemple suivant définit une moyenne personnalisée en tant que Scala Aggregator, l'enregistre dans udaf et l'utilise en tant que mesure sur une fenêtre bornée de 10 minutes :

Scala
import org.apache.spark.sql.{Encoder, Encoders}
import org.apache.spark.sql.expressions.Aggregator
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.RollingWindow.{FrontierSpec, Range, preceding}

case class AvgBuffer(sum: Long, count: Long)

object CustomAverage extends Aggregator[Long, AvgBuffer, java.lang.Double] {
override def zero: AvgBuffer = AvgBuffer(0L, 0L)
override def reduce(buffer: AvgBuffer, value: Long): AvgBuffer =
AvgBuffer(buffer.sum + value, buffer.count + 1)
override def merge(left: AvgBuffer, right: AvgBuffer): AvgBuffer =
AvgBuffer(left.sum + right.sum, left.count + right.count)
override def finish(buffer: AvgBuffer): java.lang.Double =
if (buffer.count == 0) null else buffer.sum.toDouble / buffer.count
override def bufferEncoder: Encoder[AvgBuffer] = Encoders.product
override def outputEncoder: Encoder[java.lang.Double] = Encoders.DOUBLE
}

val customAvg = udaf(CustomAverage)

val result = events.rollingWindow(
partitionBy = Seq("user_id"),
orderBy = "event_time",
frontierSpec = FrontierSpec(delay = "0 seconds", alias = "frontier"),
measures = Seq(
customAvg(col("value").cast("long"))
.over(preceding(Range("10 minutes")))
.as("custom_avg_10m")
)
)

Limitations​

  • Utilisez le mode de sortie update. Les modes de sortie append et complete ne sont pas pris en charge.
  • Chaque mesure peut comporter au plus une clause .over(). Vous ne pouvez pas enchaîner les spécifications de fenêtre.
  • Vous ne pouvez pas combiner les fenêtres glissantes avec les fonctions de fenêtre Spark standard. Une expression de colonne unique ne peut pas utiliser à la fois .over(WindowSpec) et .over(RollingWindow...).
  • La frontière progresse uniquement par temps de traitement. Les frontières basées sur l'heure de l'événement ne sont pas prises en charge.
  • Scala Aggregator UDAFs are supported only on classic compute.
  • Databricks recommande d'utiliser des durées de fenêtre glissante inférieures ou égales à 7 jours. Les durées supérieures à 7 jours peuvent nécessiter un stockage d'exécuteur important et entraîner des performances lentes ou des échecs de clusters.