リアルタイム モードのクエリ パフォーマンスを最適化および監視する

このページでは、コンピューティングのチューニング、エンドツーエンドの待機時間を短縮する手法、およびリアルタイム モードでクエリのパフォーマンスを測定する方法について説明します。

計算の最適化

コンピューティングを構成するときは、次の点を考慮してください。

  • マイクロバッチ モードとは異なり、リアルタイム タスクはデータの待機中にアイドル状態を維持できるため、リソースの無駄を避けるためには適切なサイズ設定が不可欠です。
  • 次の調整により、ターゲット クラスターの使用率レベル (50%など) を目指します。
    • maxPartitions (Kafka の場合)
    • spark.sql.shuffle.partitions (シャッフル ステージの場合)
  • Databricks では、オーバーヘッドを削減するために、各タスクが複数の Kafka パーティションを処理するように、 maxPartitions を設定することをお勧めします。
  • 単純な 1 段階のジョブのワークロードに合わせて、ワーカーごとのタスク スロットを調整します。
  • シャッフルが多いジョブの場合、バックログを回避するためのシャッフルパーティションの最小数を実験で見つけ、その結果に基づいて調整してください。 十分なスロットがない場合、コンピューターはジョブをスケジュールしません。

注

Databricks Runtime 16.4 LTS 以降では、すべてのリアルタイム パイプラインでチェックポイント v2 を使用して、リアルタイムモードとマイクロバッチ モードをシームレスに切り替えることができます。

待機時間の最適化

構造化ストリーミング リアルタイム モードには、エンドツーエンドの待機時間を短縮するための省略可能な手法があります。 どちらも既定では有効になっていません。 個別に有効にする必要があります。

  • 非同期進行状況の追跡: オフセットへの書き込みとコミット ログを非同期スレッドに移動し、ステートレス クエリのバッチ間時間を短縮します。
  • 非同期状態チェックポイント処理: 状態のチェックポイント処理を待たずに、計算が完了するとすぐに次のマイクロバッチの処理を開始し、ステートフル クエリの待機時間を短縮します。

監視と可観測性

リアルタイム モードでは、従来のバッチ期間メトリックには、実際のエンドツーエンドの待機時間は反映されません。 待機時間を正確に測定し、クエリのボトルネックを特定するには、次の方法を使用します。

エンドツーエンドの待機時間はワークロード固有であり、ビジネス ロジックでのみ正確に測定できる場合があります。 たとえば、ソース タイムスタンプが Kafka で出力される場合、Kafka の出力タイムスタンプとソース タイムスタンプの差として待機時間を計算できます。

組み込みのメトリック StreamingQueryProgress

StreamingQueryProgress イベントは、ドライバー ログに自動的に記録され、StreamingQueryListenerのonQueryProgress()コールバック関数を介してアクセスできます。 これにより、たとえば外部監視システムにメトリックを発行する場合など、進行状況イベントにプログラムで対応できます。 これらのリアルタイムモードメトリックをQueryProgressEvent.json()またはtoString()に含めます。

  1. 待機時間の処理 (processingLatencyMs)。 リアルタイム モードのクエリがレコードを読み取ったときと、クエリが次のステージまたはダウンストリームにレコードを書き込むまでの経過時間。 システムは、タスクごとにこのメトリックを報告します。
  2. ソース キューの待機時間 (sourceQueuingLatencyMs)。 システムがメッセージ バスにレコードを書き込むまでの経過時間 (Kafka のログの追加時間など) と、リアルタイム モードクエリが最初にレコードを読み取るまでの時間。 システムは、タスクごとにこのメトリックを報告します。
  3. エンドツーエンドの待機時間 (e2eLatencyMs)。 システムがメッセージ バスにレコードを書き込み、リアルタイム モードクエリがレコードをダウンストリームに書き込むまでの時間。 システムは、すべてのタスクによって処理されたすべてのレコードにわたって、バッチごとにこのメトリックを集計します。

JSON の進捗イベントには、latencies 配下にこれらのメトリクスが含まれます。 例えば次が挙げられます。

{
  "latencies": {
    "processingLatencyMs": {
      "P0": 0,
      "P50": 0,
      "P90": 0,
      "P95": 0,
      "P99": 0
    },
    "sourceQueuingLatencyMs": {
      "P0": 0,
      "P50": 1,
      "P90": 1,
      "P95": 2,
      "P99": 3
    },
    "e2eLatencyMs": {
      "P0": 0,
      "P50": 1,
      "P90": 1,
      "P95": 2,
      "P99": 4
    }
  }
}

タスク利用率を監視する

busyTimeFractionを使って、Sparkタスクの利用率がスループットを制限しているかどうかを判断します。 値が1に近い場合は、タスクが十分に活用されており、クエリがより多くの計算を必要とすることを示します。 値が低いほど、タスクがアイドルやブロックされている時間が長くなっていることを示します。 この指標は0から1の範囲で、ステージとタスクごとに報告されます。

ストリーミングクエリを開始する前にデバッグメトリクスを有効にしてください:

spark.conf.set("spark.databricks.streaming.execution.enableDebugMetrics", "true")

タスク指標にアクセスするには、以下のいずれかの方法を用いてください。

生データ

クエリを開始または再開してください。 トリガーが完了した後、クエリセルの下にある生データを開き、latenciesの下にある_taskMetricsを見つけます。 以下の例はタスク指標を示しています:

{
  "latencies": {
    "_taskMetrics": {
      "stage_0_task_0": {
        "busyTimeFraction": 0.03
      }
    }
  }
}

Python

トリガーが完了した後、lastProgressを通じて_taskMetricsアクセスします:

import json

task_metrics = json.loads(query.lastProgress.json)["latencies"]["_taskMetrics"]
print(task_metrics)

Observe API を使用したカスタム待機時間の測定

Observe API を使用すると、別のジョブを起動せずに、インラインで待機時間を測定できます。 ソース データの到着時間に近いソース タイムスタンプがある場合は、シンクの前にタイムスタンプを記録し、その差を計算することで、バッチごとの待機時間を見積もることができます。 結果は進行状況レポートに表示され、リスナーが利用できるようになっています。

Python

from datetime import datetime

from pyspark.sql.functions import avg, col, lit, max, percentile_approx, udf, unix_millis
from pyspark.sql.types import TimestampType

@udf(returnType=TimestampType())
def current_timestamp():
  return datetime.now()

# Query before outputting
.withColumn("temp-timestamp", current_timestamp())
.withColumn(
  "latency",
  unix_millis(col("temp-timestamp")).cast("long") - unix_millis(col("timestamp")).cast("long"))
.observe(
  "observedLatency",
  avg(col("latency")).alias("avg"),
  max(col("latency")).alias("max"),
  percentile_approx(col("latency"), lit(0.99), lit(150)).alias("p99"),
  percentile_approx(col("latency"), lit(0.5), lit(150)).alias("p50"))
.drop(col("latency"))
.drop(col("temp-timestamp"))
# Output part of the query. For example, .WriteStream, etc.

Scala

import org.apache.spark.sql.functions.{avg, col, lit, max, percentile_approx, udf, unix_millis}

val currentTimestampUDF = udf(() => System.currentTimeMillis())

// Query before outputting
.withColumn("temp-timestamp", currentTimestampUDF())
.withColumn(
  "latency",
  col("temp-timestamp").cast("long") - unix_millis(col("timestamp")).cast("long"))
.observe(
  name = "observedLatency",
  avg(col("latency")).as("avg"),
  max(col("latency")).as("max"),
  percentile_approx(col("latency"), lit(0.99), lit(150)).as("p99"),
  percentile_approx(col("latency"), lit(0.5), lit(150)).as("p50"))
.drop(col("latency"))
.drop(col("temp-timestamp"))
// Output part of the query. For example, .WriteStream, etc.

サンプル出力:

"observedMetrics" : {
  "observedLatency" : {
    "avg" : 63.8369765176552,
    "max" : 219,
    "p99" : 154,
    "p50" : 49
  }
}

その他のリソース