How to use LakeFlow Pipelines
このページでは、データパイプラインのライフサイクル全体を通じた Lakeflow パイプラインの使用方法について、最初の設計上の決定から大規模な実行まで、および各段階におけるトレードオフを説明します。各セクションには、その方法を説明する記事へのLinkがあります。
このガイドは、コアとなるデータエンジニアリングの概念に精通していることを前提としています。パイプラインが初めてなら、 まずはApache Spark宣言型パイプライン から製品とその背後にある宣言モデルを読み、その後「 チェンジデータキャプチャを使ってETLパイプラインを構築する」チュートリアルを進めてください。
パイプラインのライフサイクルの概要
パイプラインは6つの段階を経て移動します:
- 計画と設計:何を作るかを決め、その合ったツールや言語、コンピュートを選びましょう。
- データを取り込む:ソースデータを確実かつ段階的にパイプラインに取り込む。
- 変換とモデリング: データをクリーニング、検証、結合、整形し、消費者が信頼できるテーブルを作成します。
- 運用化:パイプラインをバージョン管理下に置き、テストし、スケジュールを組み、環境全体でプロモーションします。
- 本番環境での実行:パイプラインが無人で稼働する際、監視、アラート、デバッグ、バックフィル、セキュリティ、そしてリネージの追跡を行います。
- 成熟と拡張: 本番運用への準備状況を確認し、データ量やチーム規模の拡大に合わせてパイプラインを健全に保ちます。
各ステージは厳密に逐次的ではありませんが、質問が発生する順序に対応しています。LakeFlow Pipelines はオーケストレーション、チェックポイント、再試行、インクリメンタル処理を処理するため、各ステージでの作業は実装というよりも設計上の決定が主となります。
計画と設計
最初の決定が、下流のすべてを形作ります。宣言型モデルと、手続き型ステップを自分で記述する場合の比較については、「Databricks における手続き型データ処理と宣言型データ処理」を参照してください。
いくつかの選択肢によって開始構成が設定されます:
- スタンドアロンのデータセットかパイプラインか。 単一のマテリアライズドビューやストリーミングテーブルは、SQLでスタンドアロンのデータセットとして定義でき、Databricksがその背後にあるリフレッシュパイプラインを管理します。Pythonのオーサリング、シンク、多段階オーケストレーションが必要なときに、LakeFlow Pipelinesをユニットとして作成・運用しましょう。単 独パイプラインとLakeFlow Pipelinesを参照してください。
- SQL または Python (あるいはその両方)。 SQL は、主にフィルター、結合、集計を行う変換に適しています。Python は、カスタムロジック、外部ライブラリ、またはプログラムによる多数の類似テーブルの生成に適しています。この選択はパイプライン全体ではなくファイルごとに行うため、両方を混在させることができ、事前に決定する必要はありません。
- Serverlessまたはクラシックコンピュート。 Serverlessが推奨されるdefaultであり、クラスター構成が不要になります。特定のインスタンスタイプ、カスタムクラスターポリシー、またはinitスクリプトが必要な場合は、クラシックを選択してください。Serverless パイプラインの設定およびパイプライン用のクラシック コンピュートの構成を参照してください。
- トリガー実行または継続実行。 Start Triggerされます。実行中のみコンピュートを消費するため。連続モードは新しいデータを最小限の遅延で処理するためにコンピュートを継続し、これが通常最大のコスト要因となるため、実証済みのレイテンシ要件に備えて活用してください。トリガー モードと連続パイプラインモードの比較を参照してください。
パイプラインは、コードが参照するデータセットから実行グラフを推論するため、設計作業の大部分はデータセットの命名と順序付けになります。重要な決定事項は、各出力がどのタイプであるべきかということです。追加処理が多い増分データの場合はストリーミングテーブル、再計算された集計や結合の場合はマテリアライズドビューを使用します。その選択はコストと正確性に大きく影響します。なぜなら、増分処理は新しいデータの速度に応じてスケーリングするのに対し、完全な再計算は履歴全体に応じてスケーリングするからです。どのタイプがどのジョブに適しているかについては、「パイプラインとは?」を参照してください。
パイプラインコードは通常のPythonやSQLなので、共有ワークスペースにデプロイする前に、自分のエディタで書き、リントし、検証できます。
このステージでは
このステージで考慮すべき質問:
- スタンドアロンデータセットとフルパイプラインのどちらを選ぶかはどうすればいいですか?
- データソースを特定し、それらに接続する方法を確認するにはどうすればよいですか?
- コードを書く前にパイプラインのアーキテクチャをどう設計すればいいですか?
- ファイル形式とストレージレイヤを選択するにはどうすればよいですか?
- ローカル開発環境を設定するにはどうすればよいですか?
- 構築を開始する前に、どのように拡張を計画し、コストを見積もればよいですか?
データを取り込む
中心となる設計上の問いは、ソースが追記専用か、それともインプレースで変更されるかという点です。それがターゲットのモデル化方法を決定します:
- クラウドストレージに格納されたファイルやメッセージバス上のイベントなどの 付加専用ソース はストリーミングテーブルに取り込み、進行状況をチェックポイントでリセットしても再処理やデータの削除を防ぎます。Auto Loaderはファイルを処理し、新しいファイルを発見し、届くたびにスキーマを推論・進化させます。Apache Kafka、Azure Event Hub、Amazon Kinesis、Google Pub/Subなどのメッセージバスは、ストリーミングテーブルに直接読み込みます。バスは同じイベントを複数回配信する可能性があるため、下流側で重複排除を行う。Azure Event Hubsについては、「 Use Azure Event Hubs as a パイプライン データソース」を参照してください。
- ほとんどのデータベースや多くの SaaS システムなど、 行の更新や削除を行うソース では、チェンジデータキャプチャ (CDC) を使用します。ランごとに完全コピーを行うのは無駄であり、ソースの増加に伴って低速化するため、CDC は前回のラン以降に変更された行のみを読み取ります。
AUTO CDCAPI は、手書きの Merge ロジックなしでそれらの変更を適用します。「AUTO CDC APIs : パイプラインによるチェンジデータキャプチャの簡素化」を参照してください。フローはストリーミングテーブルに CDC を適用します。複数のフローを 1 つのテーブルに供給できるため、複数のソースを単一のターゲットに集約(ファンイン)できます。
チェックポイントと再試行は自動的に行われるため、パイプラインはすべてを再処理するのではなく、最後に処理されたオフセットから再開されます。2つのセーフガードはオプトイン方式です:
- レスキューデータ列 は、期待されるスキーマに一致しないレコードをキャプチャします。
- Expectations は、定義した行レベルのアクションを適用します。
ストリーミングのチェックポイントが無効になった場合は、テーブルデータを保持する最も低コストな復旧方法を優先してください。
このステージでは
このステージで考慮すべき質問:
- データベースからデータを取り込むにはどうすればよいですか?また、フルロードとCDCのどちらを選択すればよいですか?
- APIからデータを取り込むにはどうすればいいですか?
- ストリーミングデータやイベントデータを取り込むには?
- ファイルを確実に取り込むにはどうすればよいですか?
- データを失わずにインジェストの失敗を処理するにはどうすればよいですか?
変換とモデル化
変換処理により、取り込まれたデータは、人やツールが信頼できるクリーンなテーブルに変わります。ここで、メダリオンパターン(ブロンズからシルバー、ゴールドへ)が具体的な形となります。
クリーニングと検証が最初に行われます。エクスペクテーションは、Lakeflowパイプラインの組み込み機能です。これは、パイプラインがすべてのランのすべての行に対して評価し、合格数と失敗数を報告するデータ品質制約であり、品質を一度限りのゲートではなく継続的なものにします。行が失敗した場合の処理(警告して保持する、破棄する、または更新を失敗させる)と、ゲートをどこに配置するかを決定します。ゲートは通常ブロンズからシルバーへの境界に配置されるため、ダウンストリームのすべてを再チェックなしで信頼できます。
結合と集計が、シルバーからゴールドへのステップを形成します。マテリアライズドビューは、既存のテーブルに対するバッチスタイルの結合や集計に適しています。これは、結果をソースと整合させるためです。クエリーとソースが許可する場合は増分的に更新され、それ以外の場合は完全に再計算されるため、どちらの方法でも同じ結果が得られます。ディメンションが変更されたときに結合を再計算するため、レイテンシーよりも正確性が重要である場合に適した選択肢となります。パイプラインはどのように更新されますか?を参照してください。ライブストリームの結合は無制限の状態を引き起こすため、ストリーミング結合と集計では、パイプラインがデータの到着遅延を待機する時間を制限するためにウォーターマークが必要です。
このステージには、2つの正確性に関する考え方が通底しています:
- べき等性 とは、パイプラインが同じ入力に対して何度実行されても同じ結果を生成することを意味します。LakeFlow Pipelines は、チェックポイント読み取りやキーベースの
AUTO CDCアップサートなど、管理対象の要素に対してべき等です。再計算されるビューで非決定論的な関数を使用しないようにすることで、独自のロジックをべき等に保つことができます。 - 「最低 1 回」と「厳密に 1 回」の処理。 管理された Delta-to-Delta テーブルは、各マイクロバッチの入力と出力をまとめて commit するため、default で「厳密に 1 回」の処理が実現されます。これは、カスタムシンク、非 Delta ターゲット、または未検証のカスタムソースなどの境界で停止します。これらの場合、書き込みを「最低 1 回」として扱い、キーに基づく Upsert などによってべき等にする必要があります。
ゆっくり変化する次元(SCD)もここに存在します。 AUTO CDC SCDタイプ1とタイプ2を直接実装しているため、履歴トラッキングロジックを書くのではなくタイプを設定するだけです。
このステージでは
このステージで考慮すべき質問:
- 受信データをクリーニングおよび検証するにはどうすればよいですか?
- 緩やかに変化するディメンション(SCD)を使用して時間の経過に伴う履歴を追跡するにはどうすればよいですか?SCD とは何ですか?
- ストリーミングデータと静的データをJOINする方法データを効率的に集計する方法
- 下流での使用に向けてデータをモデル化するにはどうすればよいですか?
- How do I ensure processing guarantees in LakeFlow Pipelines?
- 少なくとも1回処理と厳密に1回処理の違いは何ですか?どちらが必要ですか?
- 到着が遅れたデータや、順序が乱れたデータはどのように処理すればよいですか?
運用化する
運用化はパイプラインを単なる作業から、チームが繰り返し構築・テスト・出荷できるものへと移行させます。パイプラインはソースコードと設定で構成されているため、通常のソフトウェアエンジニアリングの実践が適用されます。
テストは、変換ロジックと、それを流れるデータの継続的な品質という 2 つの側面を同時にカバーします。エクスペクテーションは、データ側を継続的に処理します。ロジックについては、変換をプレーンな関数に分解してランタイム外で単体テストを行い、何かをマテリアライズする前にドライランでパイプライングラフを検証してください。パイプラインの単体テストを参照してください。
パイプラインコードをGitで管理し、デプロイ用にパッケージ化することで、レビュー、ロールバック、および環境間での一貫したデプロイが可能になります。このパッケージはLakeFlow Pipelinesの代替ではありません。これはパイプラインを囲むプロジェクトおよびCI/CDラッパーであり、データロジックは宣言型のまま維持されます。カタログ名やパスなどの環境固有の値をパラメーター化することで、同じコードを各環境で変更せずに実行できるようにします。パイプラインでのパラメーターの使用を参照してください。
スケジュール上でパイプラインを実行するには、 ワークフロー内のパイプライン実行でラップします。Databricksはジョブとのパイプラインのスケジューリングとオーケストレーションを推奨しており、これにより下流レポートや複数のパイプラインの連結など他の作業とパイプラインを調整することも可能です。実行内では パイプライン が独自のデータセットを注文・並列化するため、オーケストレーションはパイプライン外のタスクのみを調整します。
このステージでは
このステージで考慮すべき質問:
- データパイプラインをテストするにはどうすればよいですか。また、それは通常のソフトウェアのテストと何が違うのですか?
- チームとしてパイプラインコードのバージョン管理や共同作業はどうすればいいですか?
- パイプラインを自動で実行させるにはどうやってスケジュールやオーケストレーションを組めばいいですか?
- パイプラインを開発環境からステージング環境、本番運用環境へ安全に移行するにはどうすればよいですか?
- パイプラインの CI/CD をセットアップするにはどうすればよいですか?
本番運用
パイプラインが実際のデータに対して無人で稼働すると、その健全かどうかを見極め、そうでないときに修正することが課題となります。
モニタリングは3段階の深さで機能します。ジョブとパイプラインリストでは、最近の実行ステータスを一目で確認できます。パイプラインモニタリングUIには、すべてのテーブルとフローがステータス別に色分けされて表示され、行数、データ品質メトリクス、ストリーミングテーブルのバックログメトリクスが示されます。両者の基盤となるイベントlogは、プログラム的または履歴的なあらゆる事柄における信頼できる唯一のソースです。失敗通知を構成して、ステークホルダーから報告される前にランの破損を把握できるようにします。モニタリング画面の概要については、「パイプラインのモニタリング」を参照してください。
グラフ上で強調された失敗からイベントログのエラー詳細まで逆算してデバッグし、失敗した部分だけを再ランします。再試行の動作はTriggerによって異なります。手動でTriggerされた更新は自動再試行を無効にし、エラーが即座に現れますが、スケジュールされた更新は回復可能な失敗を再試行します。したがって、本番運用のアラートはリトライ時に自動的に解除されるかもしれませんが、インタラクティブ開発中は同じ失敗が起きないのです。開発中、 Genie Code は反復作業でコードレベルのエラーを診断・修正するのに役立ちますが、現在は本番運用実行の診断よりもオーサリングパイプラインをターゲットにしています。
バックフィルを、通常の増分フローと同じターゲットに供給する、明示的な1回限りのフローとしてモデル化します。履歴がいつ、どのように読み込まれたかを個別のレコードとして保持することで、定常状態のロジックをシンプルに保ちます。
パイプラインを保護するには、操作権限を制御し、個人アカウントではなく専用の Service Principal として実行し、認証情報をソースコードではなく Secret Scope に保持します。リネージは自動であり、列レベルまでキャプチャされます。パイプラインは、LakeFlow Pipelines のシンクを通じて外部システムに書き込みを行います。これは、前述の「少なくとも1回」という考え方が適用される境界です。
このステージでは
このステージで考慮すべき質問:
- パイプラインが正常に実行されたかどうかを監視するにはどうすればよいですか?
- 何かが壊れたときにどうやってアラートを受ければいいのですか?
- 失敗したパイプラインのランをデバッグするにはどうすればよいですか?
- ヒストリカルデータをバックフィルするにはどうすればよいですか?
- パイプラインの運営コストをどう管理し予測すればいいですか?
- 認証情報、アクセス制御、個人情報を含むパイプラインをどのように保護すればよいですか?
- パイプラインのドキュメント化とデータリネージの追跡を行うにはどうすればよいですか?
成熟してスケール
成熟したパイプラインは、書き換えなしで自動的に実行され、拡張されます。準備状況の確認とスケーリング方法の計画が、このステージを定義します。
本番運用の準備状況は、データ品質、信頼性、観測可能性、デプロイメント、コスト、ガバナンスにわたるチェックリストです。チェックされていない項目はすべて既知のギャップとして扱ってください。不正なデータを受け取る可能性のある各データセットに期待(expectation)が設定されているか、パイプラインは手動起動ではなくスケジュールされているか、失敗通知が構成されているか、Service Principalとして実行されているか、少なくとも開発環境と本番環境のターゲットに対してバージョン管理からデプロイされているかを確認します。データ品質と通知は、追加コストが最も低く、検出されない不正なランをキャッチする可能性が最も高い機能です。
パイプラインの健全性が低下しているという具体的なシグナルに応じてスケーリングします:
- 更新期間が上昇傾向にあります。
- オートスケールが繰り返し上限に達している。
- コストが基礎となるビジネスよりも速いペースで増加しています。
- マテリアライズドビューがフルリコンピュートにフォールバックしています。
まずはServerlessに移行したり、パフォーマンスモードを自分のレイテンシに合わせるなど、コンピュートレベルのレバーを試してみてください。それ以上に、パイプライン間のデータセットの整理方法が最も重要です:
- パイプラインには同時実行数の制限があります 。一度に更新できるデータセットの数は決まっています。パイプラインのデータセット数がその制限を超えると、超過分の更新はキューで待機するため、パイプラインの合計更新時間が増加します。
- 関連するデータセットをグループ化し、関連のないデータセットを分割します。 ドメイン、共有更新周期、依存関係でグループ化し、オーナシップ、レイヤー、レイテンシの境界で分割します。例えば、取り込みと変換を分離することで、低速な取り込みがダウンストリームのすべてを遅延させることを防ぎ、各パイプラインを同時実行制限内に収まる小ささに保つことができます。
すでに本番運用されている大きなパイプラインを分割するよりも、後から 2 つの小さなパイプラインをマージする方が簡単です。データセットのグループ化および分割方法については、Lakeflow Pipelines全体でのデータセットの整理を参照してください。
このステージでは
このステージで考慮すべき質問: