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

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

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

Databricksは、常に以下の設定を行うことを推奨します。

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

ワークロードによっては、次のメリットがあります。

Databricksは、Structured Streamingワークロードの本番運用インフラストラクチャの管理の複雑さを軽減するために、LakeFlow パイプラインを導入しました。Databricksでは、新しいStructured StreamingパイプラインにLakeFlow Pipelinesの使用を推奨しています。See Spark宣言型パイプライン.

注記

コンピュートのオートスケーリングには、Structured Streaming ワークロードのクラスターサイズをスケールダウンする際の制限があります。Databricks は、ストリーミングワークロード向けに、拡張オートスケールを備えたLakeflow 上でのSpark宣言型パイプラインの使用を推奨しています。オートスケールを使用したLakeFlow Pipelinesクラスター使用率の最適化を参照してください。

:::note サーバレスコンピュート

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

サーバーレス コンピュートでの連続ストリーミングの場合は、連続モードでトリガー モードと連続パイプライン モードを使用します。

ストリーミングの制限事項をご覧ください。

:::

運用ストリーミングのレイテンシー削減

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

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

リアルタイムモードがワークロードに適さない場合、以下のベストプラクティスに従うことで、マイクロバッチ Structured Streaming のレイテンシーを削減できます。

  • 出力モード : クエリー演算子とシンクがサポートしている場合は、更新モードを使用してください。更新モードは各Triggerの後に更新された行を出力し、ウォーターマークが期限切れになるまで更新を続けます。そのため、更新された結果を処理できるように、ダウンストリームのシンクを冪等(べきとう)にする必要があります。ストリーム・ストリーム結合など、更新モードがサポートされていないワークロードの場合、または到着が遅れたデータを破棄できる場合は、アペンドモードを使用してください。低レイテンシが必要な場合は、コンプリートモードを使用しないでください。「構造化ストリーミングの出力モードの選択」を参照してください。
  • Trigger : processingTimeTriggerと0間隔を使用します。これにより、前のマイクロバッチが終了し、新しいデータが利用可能になるとすぐに次のマイクロバッチが起動します。これによりマイクロバッチのレイテンシは最小になりますが、クラウドストレージのAPIコストが増加します。運用ワークロードにはAvailableNowOnce、またはContinuousを使用しないでください。「Configure Structured Streaming Trigger intervals」を参照してください。
  • ウォーターマーク : ワークロードで破棄してはならない遅延データを含められるよう、ウォーターマークを十分に長く設定してください。ウォーターマークは、クエリーが順不同のイベント時間データを受け入れてから、それを破棄して状態を削除するまでの時間を制御します。そのため、ウォーターマークが短すぎると、有効な遅延レコードが警告なしに破棄されます。その制約内で、ウォーターマークを短くするとレイテンシーが低下し、保持する状態が少なくなります。一方、ウォーターマークを長くすると、レイテンシーと状態のコストはかかりますが、より多くの遅延データを許容できます。2倍などのレイテンシSLAの小さな倍数は、チューニングを開始する際の妥当な基準となります。「ウォーターマークを適用してデータ処理のthresholdを制御する」を参照してください。
  • ソースとスミック :メッセージバス(Apache Kafka、Amazon Kinesis、Apache Pulsar、Google Cloud Pub/Sub)などの低遅延ソースから読み込み、Delta LakeやApache Icebergテーブルからチェンジデータフィードを変更します。メッセージバス、運用データベース、 foreach シンクなどの低遅延・高throughputのシンクに書き込みを行います。下流のコンシューマーが重複データや遅延データを処理できるように、シンク操作は冪等性を持つように設計してください。
  • 状態とチェックポイント処理 :ステートフルクエリーには RocksDB 状態ストアを使用します。これは、変更ログのチェックポイント処理と非同期状態チェックポイント処理の両方に必要です。変更ログのチェックポイント設定を有効にして、増分状態の変更のみを永続化します。状態チェックポイント処理がバッチ期間のボトルネックになっている場合は、非同期状態チェックポイント処理を有効にして、チェックポイントの書き込みを次のマイクロバッチとオーバーラップさせます。ただし、その前に障害復旧およびクラスターサイズ変更に関する注意事項を確認してください。各クエリーに、耐久性のあるクラウドストレージ内の独自のチェックポイントディレクトリを割り当てます。Databricks での RocksDB 状態ストアの構成ステートフルクエリーの非同期状態チェックポイント、およびStructured Streaming チェックポイントを参照してください。
  • オフセット管理 : 連続ストリームにおけるオフセットのチェックポインティングによるレイテンシを削減するには、非同期進行状況の追跡を有効にします。これにより、データ処理をブロックすることなくオフセットとcommit Logsが更新されます。これは AvailableNow または Once のTriggerとは互換性がありません。非同期進行状況の追跡を参照してください。
  • ストレージホップ : 可能な限り、単一のストリーミングパイプライン内でコンピュートを維持してください。複数のジョブやパイプラインにロジックを分割すると、ストレージホップが増加し、レイテンシが長くなります。

障害を想定するようにストリーミング ワークロードを設計する

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

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

注記

クエリが再開されると、前回の実行中に計画されたマイクロバッチが処理されます。 メモリ不足エラーが原因でジョブが失敗した場合、またはマイクロバッチが大きすぎるためにジョブを手動でキャンセルした場合は、マイクロバッチを正常に処理するためにコンピュートをスケールアップする必要があります。

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

ジョブが再試行されるとき

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

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

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

これらの戦略を組み合わせることもできます。 次の表では、これらのアプローチを比較しています。

戦略

複数のタスク

複数のクエリ

コンピュートはどのように共有されますか?

Databricks では、各ストリーミング タスクに適切なサイズでジョブ コンピュートをデプロイすることをお勧めします。 必要に応じて、タスク間でコンピュートを共有できます。

すべてのクエリは同じコンピュートを共有します。 クエリをスケジューラプールに任意で割り当てることができます。

再試行はどのように処理されますか?

すべてのタスクは、ジョブが再試行される前に失敗する必要があります。

クエリが失敗した場合、タスクは再試行します。

戦略

複数のタスク

複数のクエリ

コンピュートはどのように共有されますか?

Databricks では、各ストリーミング タスクに適切なサイズでジョブ コンピュートをデプロイすることをお勧めします。 必要に応じて、タスク間でコンピュートを共有できます。

すべてのクエリは同じコンピュートを共有します。 クエリをスケジューラプールに任意で割り当てることができます。

再試行はどのように処理されますか?

すべてのタスクは、ジョブが再試行される前に失敗する必要があります。

クエリが失敗した場合、タスクは再試行します。

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

構造化ストリーミング ジョブを構成して、失敗時にストリーミング クエリを再開する

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

継続トリガーは、デフォルトでは以下の動作をします。

  • ジョブの複数の並列実行を防止します。
  • 前の実行が失敗したときに、新しい実行を開始します。
  • 再試行には指数バックオフを使用します。

Databricks 、ワークフローをスケジュールする際には、常に 汎用 コンピュートではなく、ジョブ コンピュートを使用することをお勧めします。 ジョブの失敗と再試行時に、新しいコンピュート リソースがデプロイされます。

注記

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()を呼び出すと、ドライバー上でクエリ例外が再度発生します。

ユースケース

挙動

汎用コンピュートに関するインタラクティブなノートブック

awaitTermination() セルを常に実行状態に保ち、クエリの状態を監視できるようにし、ノートブックの出力に障害が確実に表示されるようにします。

地域および開発環境

Sparkプログラムをローカルで実行する場合、メインスレッドが完了するとプロセスは終了します。ストリーミングクエリが完了または失敗するまでプログラムを継続させるには、 awaitTermination()を呼び出してください。

ドライバーへの障害伝播

awaitTermination()がない場合、ジョブ以外のコンテキストでのストリーミングクエリの失敗は、呼び出し元のスレッドに伝播しない可能性があります。クエリーはサイレントに失敗する可能性があり、その結果、失敗の検出と診断が困難になります。awaitTermination()を呼び出すと、ドライバー上でクエリ例外が再度発生します。