メインコンテンツまでスキップ

Structured Streamingにおけるローリングウィンドウ

ローリングウィンドウは、継続的に前進する フロンティア にわたって時間ベースの集計をコンピュートします。各イベントを固定された時間バケットに割り当てるタンブリングウィンドウやスライディングウィンドウとは異なり、ローリングウィンドウは、フロンティアが前進し、古いイベントが範囲外になると、更新された集計を出力します。直近のルックバック期間における集計値の現在の値が必要な場合は、ローリングウィンドウを使用して、リアルタイムのFeature Servingと低レイテンシのモニタリングを行います。

ローリングウィンドウでは、以下がサポートされます。

  • 過去 10 分間など、固定されたルックバック期間で集計する 境界付きウィンドウ 。
  • ストリームの起動以降のすべてのデータを蓄積する ライフタイムウィンドウ 。
  • 1つのオペレーター内の複数のウィンドウ 。これにより、1回の呼び出しで異なるルックバック期間をコンピュートできます。

要件​

  • Databricks Runtime 19以降。
  • マイクロバッチ モード:クラシック コンピュート(標準または専用アクセス モード)または Serverless コンピュート。
  • リアルタイムモード:クラシックコンピュート(標準または専用アクセスモード)。LakeFlow Pipelinesでは、パイプライン構成により、Serverlessコンピュート上でリアルタイムモードをランすることもできます。「LakeFlow Pipelinesでのリアルタイムモードの使用」を参照してください。

情報の詳細については、実行モードを参照してください。

ローリングウィンドウの仕組み​

A rolling window クエリー グループ events by one or more partition columns, orders them by a Timestamp column, and コンピュート aggregates relative to a frontier.Each output row includes the frontier value as a timestamp column.

Frontier​

フロンティアは、オペレーターの現在の「今」を定義します。これは処理時間(経過実時間)に基づいて進み、オプションで設定可能な delay だけ遅らせることができます。わずかな取り込み遅延を吸収するために遅延が役立ち、フロンティアがイベントを通過する前に処理中のイベントが確実にカウントされます。

FrontierSpec でフロンティアを定義します:

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

Bounded windows​

バウンドウィンドウは、フロンティアからさかのぼって、固定期間内のデータを集約します。.over()を使用してルックバックを定義します:

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

有効期間ウィンドウ​

ライフタイムウィンドウは、期限なしでストリームの起動時点からデータを蓄積します。ライフタイムウィンドウを作成するには、.over()を省略します:

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

実行モード​

ローリングウィンドウは、マイクロバッチモードとリアルタイムモードの両方で実行されます。フロンティアの動作はそれぞれ異なります。

  • マイクロバッチ モード :境界は各マイクロバッチで固定されます。各パーティションキーは、バッチごとに最大 1 行の出力行を出力します。
  • リアルタイムモード :バッチの進行に伴ってフロンティアが継続的に進みます。パーティションキーは、フロンティアの進行に伴い、1つのバッチ内で複数の出力行を生成できます。

レイトイベントとウォーターマーク​

ローリングウィンドウは処理時間フロンティアからウィンドウ境界を設定しますが、イベント時間ウォーターマークを適用して、遅れて到着するイベントをドロップすることができます。orderBy に渡すのと同じ列を使用して、rollingWindow() の前に .withWatermark() を追加します。

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

イベント時間がウォーターマークの時間以下であるイベントは、ウィンドウの時間範囲内に収まる場合でも削除されます。ウォーターマークがない場合、遅延到着を含め、ウィンドウに該当するすべてのイベントが集計されます。ウォーターマークはマイクロバッチモードとリアルタイムモードの両方で適用され、ウォーターマークの位置はチェックポイントに保存されるため、クエリーが再起動したときに復元されます。

構文​

ストリーミング DataFrameでrollingWindow()を呼び出します:

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

rollingWindow()演算子は、次のパラメーターを受け取ります。

パラメーター

説明

partitionBy

必須。イベントをグループ化する基準となる 1 つ以上の列。

orderBy

必須。イベントを順序付けるTimestamp列。「TimestampType」列である必要があります。

frontierSpec

フロンティアの定義。フロンティアが処理時間からどの程度遅れるかを制御するには delay を設定し、出力のフロンティア列に名前を付けるには alias を設定します。

measures

集計式のリスト。制限付きウィンドウの場合は .over() を追加し、ライフタイムウィンドウの場合は省略します。

パラメーター

説明

partitionBy

必須。イベントをグループ化する基準となる 1 つ以上の列。

orderBy

必須。イベントを順序付けるTimestamp列。「TimestampType」列である必要があります。

frontierSpec

フロンティアの定義。フロンティアが処理時間からどの程度遅れるかを制御するには delay を設定し、出力のフロンティア列に名前を付けるには alias を設定します。

measures

集計式のリスト。制限付きウィンドウの場合は .over() を追加し、ライフタイムウィンドウの場合は省略します。

サポートされている集計関数​

ローリングウィンドウでは、次の集計関数がサポートされています。

  • sum
  • count
  • avg
  • min
  • max
  • stddev、サンプル(stddev_samp)および母集団(stddev_pop)のバリアントを含む
  • variance、サンプル(var_samp)および母集団(var_pop)のバリアントを含む
  • approx_count_distinct
  • approx_percentile
  • functions.udafに登録されているScala Aggregator UDAFs。クラシックコンピュートでのみサポートされます。例については、例 3: Scala UDAFを参照してください。

例​

例 1: 過去 1 時間のユーザーあたりのローリング収益​

次の例では、過去 1 時間のユーザーごとのローリング収益合計をコンピュートし、更新を 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()
)

例2:有界ウィンドウとライフタイムウィンドウの混在​

次の例では、単一の演算子で2つのバインドされた集計と1つのライフタイム集計をコンピュートします。

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

出力スキーマは (region, product, window_ts, sum(revenue), avg(price), count(value)) です。

例 3:Scala UDAF​

注記

Scala Aggregator UDAFs は、クラシックコンピュートでのみサポートされています。

次の例では、カスタム平均をScalaのAggregatorとして定義し、udafに登録し、10分間の制限付きウィンドウのメジャーとして使用します。

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

制限事項​

  • update出力モードを使用します。append および complete 出力モードはサポートされていません。
  • 各メジャーには、最大で1つの .over() 句を含めることができます。ウィンドウ仕様をチェーンすることはできません。
  • ローリング ウィンドウを通常の Spark ウィンドウ関数と組み合わせることはできません。単一の列式で .over(WindowSpec) と .over(RollingWindow...) の両方を使用することはできません。
  • フロンティアは処理時間のみに基づいて進行します。イベント時間のフロンティアはサポートされていません。
  • Scala Aggregator UDAFs は、クラシックコンピュートでのみサポートされています。
  • Databricksでは、ローリングウィンドウの期間を7日以下に設定することをお勧めします。7日を超える期間では、多量のエグゼキューターのストレージが必要になり、パフォーマンスの低下やクラスター障害の原因となる可能性があります。