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

統合CDCパイプラインのスマートなクローズ

統合された CDC パイプラインは Trigger モードでランします。各更新では、ソースデータベースから変更が抽出され、送信先テーブルに適用された後、停止します。 Smart closure は、更新が実行される時間を決定するポリシーです。各更新に変更を処理する時間を付与し、ソースに追いついた後に更新を停止し、更新の実行時間を制限します。

注記

スマートクローズは、変更抽出機能と適用機能を個別のコンポーネントとしてではなく、単一のトリガーされたパイプラインで一緒に実行する統合CDCパイプライン(ダイレクトCDCとも呼ばれます)に適用されます。これは、取り込みゲートウェイが継続的に実行される標準のゲートウェイベースのアーキテクチャには適用されません。SQL Server向けの統合CDCパイプラインを作成するを参照してください。

更新が停止したとき

Databricksは、運用上の経験に基づいて、最小ランタイム、最大ランタイム、およびラグのthresholdを構成および調整します。

条件

意味

最小ランタイム

各更新は、ソースに追いついたために停止できるようになるまで、最低限の時間実行されます。

ソースは最新の状態です。

最小ランタイムの後、保留中の変更のバックログが適用され、ソースにほぼ追いついた時点で更新が停止します。パイプラインイベントログには、完了理由 lag-converged が記録されます。

ランタイムの制限に達しました。

アップデートが追いつかない場合、最大ランタイムで停止します。次のアップデートは、このアップデートが停止した場所から再開されます。パイプラインイベントログには、完了理由 max-runtime-cap-hit が記録されます。

ソーススキーマの変更

ソーススキーマの変更により、現在のパイプラインの更新が停止します。その後、パイプラインは新しいスキーマを使用する新しい更新を起動します。

条件

意味

最小ランタイム

各更新は、ソースに追いついたために停止できるようになるまで、最低限の時間実行されます。

ソースは最新の状態です。

最小ランタイムの後、保留中の変更のバックログが適用され、ソースにほぼ追いついた時点で更新が停止します。パイプラインイベントログには、完了理由 lag-converged が記録されます。

ランタイムの制限に達しました。

アップデートが追いつかない場合、最大ランタイムで停止します。次のアップデートは、このアップデートが停止した場所から再開されます。パイプラインイベントログには、完了理由 max-runtime-cap-hit が記録されます。

ソーススキーマの変更

ソーススキーマの変更により、現在のパイプラインの更新が停止します。その後、パイプラインは新しいスキーマを使用する新しい更新を起動します。

スマートクロージャーの活用法

  • 変更が少ない場合のコスト削減: 最小ランタイム経過後、処理が追いついた時点で更新が終了します。この動作により、更新の間で変更を蓄積させることができ、コンピュートの継続的な実行にかかるコストを削減できます。
  • 境界のある、予測可能なランタイム: 大量のバックログがある場合でも、1回の更新がいつまでも実行され続けることはありません。各更新には上限があり、大規模なワークロードは後続のスケジュールされた更新に分散されます。
  • 完了状況の可視性:各更新は、それが終了した理由を記録するため、ソースに追いついたのか、それともランタイム制限で停止したのかを判別できます。

更新完了を監視

完了理由は、パイプライン イベント ログの COMPLETED イベントのメッセージに表示されます。ソースに追いついた更新は理由lag-convergedで完了し、ランタイムの制限で停止した更新は理由max-runtime-cap-hitで完了します。

完了理由を確認するには、エクストラクタの COMPLETED イベントについてパイプラインイベントログをクエリします。<pipeline-id> をパイプライン ID に置き換えます:

SQL
SELECT timestamp, message
FROM event_log('<pipeline-id>')
WHERE message LIKE '%Direct Cdc Extraction has COMPLETED%'
ORDER BY timestamp DESC

イベントメッセージには理由が埋め込まれています。たとえば Direct Cdc Extraction has COMPLETED (reason=lag-converged)

定期的な更新をスケジュール

更新の継続時間はソースの変更データの量によって異なるため、大規模なバックログは1回の更新では完了しない場合があります。定期的なスケジュールでデータを取り込むには、パイプラインを実行する Lakeflow Jobs タスクを作成します。後続の更新が追いつくように十分な頻度でスケジュールします。60分という開始点は、ほとんどのワークロードでうまく機能します。更新が停止した後、次の更新は、設定されたスケジュールまたは手動の起動のいずれか早い方で開始されます。前の更新がまだ実行されている間にスケジュールされたTriggerが起動した場合、Databricksはその更新をスキップし、次のスケジュールされたランを使用します。

関連リソース