構造化ストリーミングの運用に関する考慮事項

Azure Databricksで、スケジュールされた Lakeflow ジョブとして運用構造化ストリーミング ワークロードを実行します。 「Lakeflow ジョブ」を参照してください。

Databricks では、常に次を構成することをお勧めします。

  • displaycountなど、結果を返すノートブックから不要なコードを削除します。
  • 汎用コンピューティングを使用して構造化ストリーミング ワークロードを実行しないでください。 ストリームは常に jobs compute を使用して Lakeflow Jobs としてスケジュールしてください。
  • Continuous モードを使用して Lakeflow ジョブをスケジュールします。 これは、Azure Databricks ジョブのスケジューリング機能を指し、構造化ストリーミングのトリガー間隔ではありません。
  • 構造化ストリーミング ジョブのコンピューティングに対して自動スケールを有効にしないでください。

一部のワークロードには、次の利点があります。

Databricks では、構造化ストリーミング ワークロードの運用インフラストラクチャの管理の複雑さを軽減するために、Lakeflow パイプラインが導入されました。 Databricks では、新しい構造化ストリーミング パイプラインに Lakeflow パイプラインを使用することをお勧めします。 Spark 宣言型パイプラインに関するページを参照してください。

コンピューティングの自動スケールには、構造化ストリーミング ワークロードのクラスター サイズのスケールダウンに制限があります。 Databricks では、ストリーミング ワークロード用に拡張された自動スケーリングを使用して、Lakeflow で Spark 宣言パイプラインを使用することをお勧めします。 自動スケーリングを使用した Lakeflow パイプライン クラスターの使用率の最適化に関する説明を参照してください。

:::note サーバーレス コンピューティング

サーバーレス コンピューティングでは、 Trigger.AvailableNow()Trigger.Once() のみがサポートされます。 Databricks では Trigger.AvailableNow() を推奨しています。

サーバーレス コンピューティングでの継続的ストリーミングには、継続的パイプライン モードを使用します。トリガーされたパイプライン モードではなく、連続モードで行ってください。

ストリーミングの制限事項を参照してください。

:::

運用ストリーミングの遅延削減

運用ストリーミングワークロードは、データをほぼリアルタイムで取り込み、変換し、対応します。 一般的な例としては、不正検出、異常検出、パーソナライズ、リアルタイム監視やアラートなどがあり、処理遅延がビジネスの成果に直接影響します。 これらのワークロードの低遅延は通常、数十から数百ミリ秒に及ぶことが多いですが、多くのチームは高いパーセンタイルでの変動を考慮して、サービスレベルアグリーメント(SLA)を秒単位で設定しています。

エンドツーエンドの遅延を最小限に抑えるにはリアルタイムモードを用いましょう。これはテールで1秒未満、一般的なケースで約300ミリ秒のエンドツーエンド遅延を実現します。 リアルタイム モードの概念を参照してください。

リアルタイムモードがあなたのワークロードに合わない場合は、以下のベストプラクティスがマイクロバッチ構造化ストリーミングの遅延を削減します:

  • 出力モード:クエリオペレーターやシンクが対応している更新モードを使いましょう。 アップデートモードでは、各トリガーの後に更新された行を出力し、ウォーターマークの有効期限が切れるまでそれらを更新し続けるため、更新された結果を処理できるよう、下流のシンクを冪等にしてください。 更新モードが対応しないワークロード、例えばストリームストリームの結合や遅れて到着したデータをドロップできる場合は追加モードを使いましょう。 低遅延のためにコンプリートモードは使わないでください。 「構造化ストリーミングの出力モードを選択する」を参照してください。
  • トリガー:0の間隔を持つprocessingTimeトリガーを使用し、前のマイクロバッチが終了し新しいデータが利用可能になったらすぐに次のマイクロバッチを開始します。 これによりマイクロバッチ遅延が最も低くなりますが、クラウドストレージAPIコストが増加します。 運用作業には AvailableNowOnceContinuous は使わないでください。 「構造化ストリーミングのトリガー間隔を構成する」を参照してください。
  • ウォーターマーク:作業量が減らないように、遅れて到着したデータを含めるためにウォーターマークを十分に長く設定してください。 ウォーターマークはクエリが順序が乱れたイベントタイムデータを受け入れる期間を制御し、その後に削除して状態を追い出すため、短すぎるウォーターマークは有効な遅延レコードを静かに破棄します。 その制約の中で、ウォーターマークが短いほどレイテンシが低く、状態の保持も少なくなります。ウォーターマークが長いほど遅延データを許容しますが、その代わりにレイテンシと状態が犠牲になります。 遅延SLAのわずかな倍数、例えば2倍の調整が妥当な出発点です。 「透かしを適用してデータ処理のしきい値を制御する」を参照してください。
  • ソースとシンク: Apache Kafka、Amazon Kinesis、Apache Pulsar、Google Cloud Pub/Sub などのメッセージバスといった低遅延ソース、または Delta Lake や Apache Iceberg テーブルの変更データフィードから読み取ります。 メッセージバス、運用データベース、 foreach シンクなどの低遅延・高スループットのシンクに書き込みを行います。 シンク操作は冪等性を持つよう設計し、下流の消費者が重複データや遅れて到達するデータを処理できるようにします。
  • 状態とチェックポイント:ステートフルクエリには、変更ログチェックポイントと非同期状態チェックポイントの両方に必要となるRocksDBステートストアを使用します。 チェンジログチェックポイントを有効にして、段階的な状態変更のみを永続化させてください。 状態チェックポイントがバッチ時間のボトルネックの場合は、障害・回復やクラスタサイズ変更の注意事項を確認した後、非同期状態チェックポイントを有効にして次のマイクロバッチとチェックポイント書き込みを重ねます。 各クエリに独自のチェックポイントディレクトリを、持続可能なクラウドストレージに設置しましょう。 Azure Databricks で RocksDB ステートストアを構成するステートフル クエリの非同期状態チェックポイント、およびStructured Streaming のチェックポイントを参照してください。
  • オフセット管理:連続ストリームでのオフセットチェックポイントによる遅延を減らすために、非同期進行状況追跡を有効にし、データ処理をブロックせずにオフセットとコミットログを更新します。 AvailableNowOnceトリガーとは互換性がありません。 非同期進行状況の追跡を参照してください。
  • ストレージホップ:可能な限り計算を単一のストリーミングパイプライン内に収めます。 複数のジョブやパイプラインにロジックを分散すると、ストレージホップが発生し遅延が増加します。

失敗が予想されるストリーミング ワークロードを設計する

Databricks では、障害発生時に自動的に再起動するようにストリーミング ジョブを常に構成することをお勧めします。 スキーマの進化を含む一部の機能では、構造化ストリーミング ワークロードが自動的に再試行される必要があります。 「障害時にストリーミング クエリを再起動するように、構造化ストリーミング ジョブを構成する」を参照してください。

foreachBatch のような一部の操作では、1 回限りの保証ではなく、少なくとも 1 回の保証が提供されます。 これらの操作を行う際には、処理パイプラインが冪等性を持つことを確認してください。 foreachBatch を使用した任意のデータ シンクへの書き込みに関するページを参照してください。

クエリが再起動すると、前回の実行中に予定されたマイクロバッチが処理されます。 メモリ不足エラーが原因でジョブが失敗した場合、またはマイクロバッチのサイズが大きいためにジョブを手動で取り消した場合は、マイクロバッチを正常に処理するためにコンピューティングのスケールアップが必要になる場合があります。

実行間で構成を変更した場合、これらの構成は計画された最初の新しいバッチに適用されます。 「構造化ストリーミング クエリの変更後に復旧する」を参照してください。

ジョブの再試行時

Azure Databricks ジョブの一部として複数のタスクをスケジュールできます。 継続的トリガーを使用してジョブを構成する場合、タスク間の依存関係を設定することはできません。

次のいずれかの方法を使用して、1 つのジョブで複数のストリームをスケジュールすることができます。

  • 複数のタスク: 継続的トリガーを使用してストリーミング ワークロードを実行する複数のタスクを含むジョブを定義します。
  • 複数のクエリ: 1 つのタスクのソース コードで複数のストリーミング クエリを定義します。

これらの戦略を組み合わせることもできます。 次の表では、これらの方法を比較します。

戦略 複数のタスク 複数のクエリ
コンピューティングはどのように共有されますか? Databricks では、各ストリーミング タスクに適したサイズの Jobs Compute をデプロイすることをお勧めします。 必要に応じて、タスク間でコンピューティングを共有できます。 すべてのクエリで同じコンピューティングが共有されます。 必要に応じて、 スケジューラ プールにクエリを割り当てることができます。
再試行はどのように処理されますか? ジョブの再試行前に、すべてのタスクが失敗します。 クエリが失敗すると、タスクは再試行します。

複数のタスクまたはクエリの操作の詳細については、「 同じクラスターで複数の構造化ストリーミング クエリを実行する」を参照してください。

障害時にストリーミング クエリを再起動するように、構造化ストリーミング ジョブを構成する

Databricks では、継続的トリガーを使用して、すべてのストリーミング ワークロードを構成することをお勧めします。 「ジョブを継続的に実行する」を参照してください。

継続的トリガーの既定の動作は次のとおりです。

  • ジョブの複数の同時実行を阻止します。
  • 前の実行が失敗したときに新しい実行を開始します。
  • 再試行にエクスポネンシャル バックオフを使用します。

Databricks では、ワークフローをスケジュールするときに、All-Purpose Compute ではなく Jobs Compute を常に使用することをお勧めします。 ジョブが失敗して再試行すると、新しいコンピューティング リソースがデプロイされます。

Databricks では、 streamingQuery.awaitTermination()spark.streams.awaitAnyTermination()は使用しないことをお勧めします。 awaitTermination()を使用するタイミングを参照してください。

使用するタイミング awaitTermination()

ストリーミング クエリが終了するまで streamingQuery.awaitTermination() および spark.streams.awaitAnyTermination() は、現在のスレッドをブロックします。 これらの関数を使用するかどうかは、実行環境によって異なります。

Lakeflow ジョブの場合は、 streamingQuery.awaitTermination() または spark.streams.awaitAnyTermination()を使用しないでください。 これらの関数は、ストリーミング クエリがアクティブな場合にジョブ サービスによって実行の完了が自動的に防止されるため、必要ありません。 どちらの関数もノートブック のセルが完了するのをブロックし、ジョブ サービスがストリーミング クエリを追跡するのを防ぎます。これにより、バックログ メトリックとジョブ通知が中断されます。

次の場合は、 awaitTermination() を使用します。

利用シーン 行動
万能コンピューティング上の対話型ノートブック awaitTermination() は、セルの実行を維持し、クエリの状態を観察し、ノートブックの出力でエラーが発生することを確認します。
ローカル環境と開発環境 Spark プログラムをローカルで実行すると、メイン スレッドが完了するとプロセスが終了します。 awaitTermination()を呼び出して、ストリーミング クエリが完了または失敗するまでプログラムを維持します。
ドライバーへのエラー伝達 awaitTermination()しないと、ジョブ以外のコンテキストでのストリーミング クエリエラーが呼び出し元のスレッドに伝達されない可能性があります。 クエリは警告なしに失敗する可能性があり、エラーの検出と診断が困難になります。 awaitTermination()を呼び出すと、ドライバーのクエリ例外が再度発生します。