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:
# 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():
# 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():
# 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:
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
- Scala
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
],
)
import org.apache.spark.sql.streaming.RollingWindow
import org.apache.spark.sql.streaming.RollingWindow.{FrontierSpec, Range, preceding}
import org.apache.spark.sql.functions._
df.rollingWindow(
partitionBy = Seq("user_id"),
orderBy = "event_time",
frontierSpec = FrontierSpec(delay = "0 seconds", alias = "frontier"),
measures = Seq(
sum("value").over(preceding(Range("10 minutes"))),
count("value")
)
)
L'opérateur rollingWindow() prend les paramètres suivants :
parameter | Description |
|---|---|
| Obligatoire. Une ou plusieurs colonnes sur lesquelles regrouper les événements. |
| Obligatoire. La colonne de timestamp qui ordonne les événements. Doit être une colonne |
| La définition de la frontière. Définissez |
| Une liste d'expressions d'agrégation. Ajoutez |
Fonctions d’agrégation prises en charge
Les fenêtres glissantes prennent en charge les fonctions d’agrégation suivantes :
sumcountavgminmaxstddev, 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_distinctapprox_percentile- UDAFs Scala
Aggregatorenregistrées avecfunctions.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 :
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 :
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
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 :
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 sortieappendetcompletene 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
AggregatorUDAFs 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.