Pular para o conteúdo principal

Janelas rolantes no Structured Streaming

Uma janela deslizante compute agregações baseadas em tempo sobre uma fronteira em avanço contínuo. Ao contrário das janelas fixas (tumbling) ou deslizantes (sliding), que atribuem cada evento a um bloco de tempo fixo, uma janela deslizante emite agregados atualizados à medida que a fronteira avança e eventos mais antigos saem do intervalo. Use janelas deslizantes para Feature Serving em tempo real e monitoramento de baixa latência, onde você precisa do valor atual de um agregado em um período de análise recente.

As janelas deslizantes oferecem suporte ao seguinte:

  • Janelas delimitadas que agregam dados ao longo de uma duração de histórico fixa, como os últimos 10 minutos.
  • Janelas de tempo de vida que acumulam todos os dados desde que a transmissão começou.
  • Múltiplas janelas em um operador , para que você possa compute diferentes durações de lookback em uma única chamada.

Requisitos​

  • Databricks Runtime 19e acima.
  • Modo de micro-lotes: compute clássico (modo de acesso Standard ou Dedicated) ou compute Serverless.
  • tempo real mode: classic compute (Standard or Dedicated access mode). In LakeFlow Pipelines, tempo real mode also execução on Serverless compute through pipeline configuration. See Use o modo de tempo-real nos Lakeflow pipelines.

Para obter mais informações, consulte Modos de execução.

Como funcionam as janelas deslizantes​

Uma query de janela deslizante agrupa eventos por uma ou mais colunas de partição, as ordena por uma coluna de timestamp e compute agregados em relação a uma fronteira. Cada linha de saída inclui o valor da fronteira como uma coluna de timestamp.

Frontier​

A fronteira define o "now" atual para o operador. Ele avança com base no tempo de processamento (tempo real) e pode, opcionalmente, apresentar atraso por um delay configurável. Um atraso é útil para absorver pequenos atrasos de ingestão para que os eventos em andamento sejam contados antes que a fronteira os ultrapasse.

Defina a fronteira com FrontierSpec:

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

Janelas delimitadas​

Uma janela limitada agrega dados dentro de uma duração fixa, olhando para trás a partir da fronteira. Defina o lookback com .over():

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

Janelas de tempo de vida​

Uma janela de tempo de vida acumula dados desde o início da transmissão sem expiração. Para criar uma janela de tempo de vida, omita .over():

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

Modos de execução​

As janelas deslizantes são executadas tanto no modo de micro-lotes quanto no modo de tempo real. A fronteira se comporta de maneira diferente em cada uma:

  • Modo de micro-lotes : a fronteira é fixa para cada micro-lotes. Cada key de partição emite no máximo uma linha de saída por lote.
  • Modo em tempo real : a fronteira avança continuamente durante um lote. Uma key de partição pode emitir várias linhas de saída dentro de um único lote conforme a fronteira avança.

Eventos tardios e marcas d'água​

Uma janela deslizante define o limite da janela a partir da fronteira do tempo de processamento, mas você pode aplicar uma marca d'água baseada no tempo do evento para descartar eventos que chegam tarde. Adicione .withWatermark() antes de rollingWindow(), usando a mesma coluna que você passa para 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")))],
)

An event whose event time is at or abaixo the watermark is dropped, even when it falls within the window's time range. Without a watermark, every event that falls within the window is aggregated, including late arrivals. Watermarks apply in both micro-lotes and tempo real mode, and the watermark position is checkpointed, so it's restored when a query restarts.

Sintaxe​

Chame rollingWindow() em um DataFrame de transmissão:

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
],
)

The rollingWindow() operator takes the following parameters:

Parâmetro

Descrição

partitionBy

Obrigatório. Uma ou mais colunas para agrupar eventos.

orderBy

Obrigatório. A coluna de Timestamp que ordena os eventos. Deve ser uma coluna TimestampType.

frontierSpec

A definição de fronteira. Defina delay para controlar o quanto a fronteira fica atrás do tempo de processamento e alias para nomear a coluna de fronteira de saída.

measures

Uma lista de expressões agregadas. Adicione .over() para uma janela limitada ou omita-o para uma janela de tempo de vida.

Parâmetro

Descrição

partitionBy

Obrigatório. Uma ou mais colunas para agrupar eventos.

orderBy

Obrigatório. A coluna de Timestamp que ordena os eventos. Deve ser uma coluna TimestampType.

frontierSpec

A definição de fronteira. Defina delay para controlar o quanto a fronteira fica atrás do tempo de processamento e alias para nomear a coluna de fronteira de saída.

measures

Uma lista de expressões agregadas. Adicione .over() para uma janela limitada ou omita-o para uma janela de tempo de vida.

Supported aggregate functions​

As janelas móveis oferecem suporte às seguintes funções agregadas:

  • sum
  • count
  • avg
  • min
  • max
  • stddev, incluindo variantes de amostra (stddev_samp) e população (stddev_pop)
  • variance, incluindo variantes de amostra (var_samp) e população (var_pop)
  • approx_count_distinct
  • approx_percentile
  • UDAFs de Scala Aggregator registradas com functions.udaf. Suportado somente no compute clássico. Para ver um exemplo, consulte Exemplo 3: UDAF em Scala.

Exemplos​

Exemplo 1: receita móvel por usuário na última hora​

O exemplo a seguir compute uma soma de receita contínua por usuário na última hora e grava atualizações no 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()
)

Exemplo 2: Janelas limitadas e de tempo de vida mistas​

O exemplo a seguir compute dois agregados delimitados e um agregado de tempo de vida em um único operador:

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
],
)

O esquema de saída é (region, product, window_ts, sum(revenue), avg(price), count(value)).

Example 3: Scala UDAF​

nota

Os UDAFs Scala Aggregator são compatíveis apenas com o compute clássico.

O exemplo a seguir define uma média personalizada como um Scala Aggregator, faz o registro com udaf e o usa como uma medida em uma janela delimitada de 10 minutos:

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")
)
)

Limitações​

  • Use o modo de saída update. Os modos de saída append e complete não são aceitos.
  • Cada medida pode ter no máximo uma cláusula .over(). Não é possível encadear especificações de janela.
  • Não é possível combinar janelas deslizantes com funções de janela regulares do Spark. Uma expressão de coluna única não pode usar .over(WindowSpec) e .over(RollingWindow...) ao mesmo tempo.
  • O avanço da fronteira é feito apenas com base no tempo de processamento. Fronteiras baseadas no horário do evento não são aceitas.
  • Os UDAFs Scala Aggregator são compatíveis apenas com o compute clássico.
  • O Databricks recomenda o uso de durações de janela rolante menores ou iguais a 7 dias. Durações maiores que 7 dias podem exigir armazenamento substancial do executor e causar desempenho lento ou falhas no cluster.