Google PDE: ワークフローのオーケストレーションとパイプラインの自動化 — 学習ガイド
こちらの一部です: Google Professional Data Engineer — 学習ガイド. 検証済みの解答で練習: Google試験ハブ, または時間制限付き模擬試験に挑戦: ExamRoll.io.
概要
ワークフローのオーケストレーションとパイプラインの自動化は、複数のサービスにまたがるデータタスクを調整し、取り込み、変換、品質チェック、公開が、信頼性、安全性、コスト効率に優れた方法で実行されるようにします。Google Cloud では、オーケストレーションは、スケジュールされたバッチ、イベント駆動のストリーム、アドホック、または長時間実行ジョブといった、各ワークロードの実行モデルに合わせる必要があります。設計目標は、再現性、べき等性、可観測性、最小権限、そして環境間での安全なプロモーションです。
主な選択肢:
- DAG、タスク依存関係、高度なスケジューリングのための、Cloud Composer (Apache Airflow) を使用したコード中心のバッチオーケストレーション。
- 軽量でイベント駆動の、サービスをまたいだシーケンスのための、Cloud Workflows を使用したサーバーレス API コレオグラフィ。
- cron の場合は Cloud Scheduler、イベントの場合は Eventarc によってトリガーされる、Cloud Run ジョブや Dataproc ジョブなどの実行エンドポイント。
- BigQuery の変換、アサーション、リリース管理のための、Dataform を使用した SQL ネイティブのオーケストレーション。
運用モデルでは、上限付き指数バックオフによるリトライ、タイムアウト、SLA、キャッチアップとバックフィル、安全な再実行のためのべき等なタスク設計、そしてデッドレターのキャプチャによる堅牢な障害処理を重視します。セキュリティは、パイプラインごとのサービスアカウント、シークレットの分離、パラメータ化、最小権限の IAM を通じて強化されます。CI/CD、Infrastructure as Code、そして包括的なテレメトリーが、本番環境に対応したアプローチを完成させます。
Google Cloud 上のオーケストレーション: ツールとパターン
Cloud Composer (Airflow)
- DAG は、明示的な依存関係を持つ有向非巡回実行グラフを定義します。タスクを表現するには、TaskFlow API またはオペレーター(例: BigQuery、Dataflow、Dataproc、Cloud Run)を使用します。センサーと遅延可能オペレーターは、待機条件(例: Cloud Storage でのオブジェクトのファイナライズや BigQuery でのパーティションの出現)に対するスケジューラの負荷を軽減します。
- スケジューリング: cron 式、start_date、end_date、および catchup が過去の実行を制御します。バックフィルには catchup を使用し、ストリーミングに隣接するターゲットやべき等でないターゲットでは無効にします。下流システムを保護するために、max_active_runs とプールで同時実行性を制限します。
- 依存関係: set_upstream/set_downstream またはタスクフローの依存関係。メタデータ駆動型のオーケストレーションでは、BigQuery のコントロールテーブル(例: クライアント/パーティションのリスト)から動的タスクマッピングを使用してタスクを動的に生成し、DAG のパース時間を安定させ、タスクをデータ駆動にします。
- DAG の(簡潔な)コード例:
undefined
undefined
undefined
undefined
undefined
undefined
undefined
undefined
undefined
Cloud Workflows、Cloud Scheduler、Cloud Run ジョブ、およびイベント駆動実行
- Cloud Workflows は、組み込みのリトライ、ループ、並列分岐、補償ロジックを使用して、Google API と HTTP エンドポイントをオーケストレーションします。BigQuery、Dataflow、Batch、Cloud Run ジョブなどのサービスをまたいだ軽量な制御フローに最適です。
- Cloud Scheduler は、cron スタイルの自動化のために、Workflows、Pub/Sub トピック、または HTTP サービスをトリガーします。毎日 02:00 のバッチ処理には、Dataflow ジョブまたは Dataproc ジョブを起動する Workflow をスケジュールします。
- Cloud Run ジョブは、自動リトライと最小限の運用で、コンテナ化されたバッチステップを実行します。Dataflow や BigQuery の周辺での多段階のデータタスクや前処理/後処理のために、Workflows との相性が良いです。
- イベント駆動: Eventarc を使用して、Cloud Storage オブジェクトのファイナライズ、Pub/Sub メッセージ、または監査ログを Cloud Run または Workflows にルーティングします。単一テーブルへの BigQuery の挿入ジョブ通知には、高度なフィルタを持つ Cloud Logging シンクを Pub/Sub に作成し、そのトピックからコンシューマをトリガーします。
Dataform: BigQuery のための SQL ワークフロー
- ref() で依存関係グラフをモデル化し、テーブル/ビュー/増分テーブルを定義し、タグやスケジュールによってビルドをオーケストレーションします。Dataform は SQLX を順序付けられた実行計画にコンパイルし、宣言的な定義からメタデータ駆動型のオーケストレーションを可能にします。
- アサーションはデータ品質を保証します。アサーションとは、パスするために 0 行を返す必要があるクエリです。 アサーションの例:
undefined
undefined
undefined
- リリースとリポジトリの管理: コードをリポジトリに保存し、ブランチとレビューを使用し、タグ付けされたリリースを環境固有の変数とともに環境(例: dev、test、prod)にプロモートします。CI/CD チェックとアサーションの結果を介してデプロイをゲートします。
Dataproc、Dataflow、およびストレージパターン
- 最小限の運用で Hadoop/Spark を再利用するには、GCS コネクタを備えた Dataproc を使用して、クラスタのライフタイムを超えてデータを永続化し、永続ディスクのコストを最小限に抑えます。分離とコスト管理のためにジョブごとにエフェメラルクラスタを作成し、Composer または Workflows でオーケストレーションします。
- 不正な形式の行を含むバッチ取り込みには、Dataflow を実行して有効なレコードを BigQuery に書き込み、パース/検証エラーを調査用のデッドレター BigQuery テーブルにルーティングします。
信頼性、障害処理、べき等性
リトライ、タイムアウト、バックオフ
- 一時的な障害には上限付き指数バックオフを使用し、合計リトライ期間をジョブのSLA内に収めます。例えば、15分ごとにデータベースをポーリングするフロントエンドやタスクは、最大15分まで指数バックオフでリトライし、その後は制御された障害として表面化させるべきです。
- Airflowではタスクごとの
execution_timeoutとグローバルなDAGのSLAを設定します。Workflowsでは、ステップごとのタイムアウトとmax_doublingsおよびmax_retry_durationを持つリトライポリシーを設定します。Cloud Runジョブでは、リトライ回数とバックオフを設定します。
バックフィル、キャッチアップ、障害処理
- タスクがべき等で、ソースが日付でパーティション分割されている場合は、履歴の再計算のためにキャッチアップを有効にします。非決定的な出力や外部の副作用がある場合は、バックフィル専用のDAGや、生成されたものを追跡するための書き込み監査テーブルを検討します。
- ストリーミング/バッチ変換におけるレコードレベルの障害には、デッドレター トピック/テーブルを使用します。バッチDataflowでは、エラータグで不正な行を捕捉し、エラーメトリクスを集計します。ストリーミングでは、Pub/SubのDLQを使用します。
べき等なタスク設計と再実行
- BigQuery:
MERGEまたは重複排除キーを持つINSERTを推奨します。ストリーミング挿入の重複排除にはinsertIdを使用します。バッチの場合は、ステージングテーブルに書き込み、その後トランザクション的に安全なステップ内でターゲットにMERGEすることで、完全な再実行を可能にします。 - Cloud Storage: generation事前条件と決定的なオブジェクト名(例:
prefix/date/hash)を使用して、再実行時に予期された場合にのみ安全に上書きされるようにします。 - Pub/SubとDataflow: at-least-once(少なくとも1回)配信を前提に設計します。メッセージ識別子(例: パッケージID、論理イベントタイムスタンプ)を含めることで、後続のシステムが重複排除や遅延の判断を行えるようにします。ビジネスルールが「最初に処理されたイベントが優先される」セマンティクスを許容する場合は、そのトレードオフを文書化し、偏りを監視します。そうでなければ、イベント時間とタイブレーカーによって勝者を決定します。
- 部分的な障害からの回復: 出力を
run_idや日付でパーティション分割し、完了マーカーを書き込み、後続タスクがマーカーに依存するようにします。未完了とマークされたパーティションのみを再処理します。
トラブルシューティングとスケーラビリティ
- ストリーミングダッシュボードでイベントが欠落しているが、Pub/Subには存在することが示されている場合、既知の固定データセットをDataflowパイプラインに流して、変換の欠陥を特定します。ウィンドウ、トリガー、許容される遅延を検証します。
- 一般的な障害モード: 非境界ソースに対して適切なウィンドウ/トリガーなしでストリーミングパイプラインを作成したり、シャーディングされたウィンドウを誤って使用したりすると、パイプラインの作成に失敗したり、状態が肥大化したりする可能性があります。
- Dataflowは、最大ワーカー数と自動スケーリングアルゴリズムによってスケーリングします。スパイク(例: 50,000インストール)に対しては、ピーク時に水平スケーリングを可能にするために最大ワーカー数を増やします。
セキュリティ、パラメータ化、環境、CI/CD
パラメータ化と構成管理
- 環境ごとに構成を外部化します。Composerでは、変数、接続、環境変数を使用し、実行日やパーティションによってDAGパラメータをテンプレート化します。Workflowsでは、ランタイム引数を使用し、環境ごとにワークフローを分離するか、Secret Managerから構成を読み取ります。
- クライアント、ソース、またはパーティションをリストした制御テーブル(例: BigQueryの構成データセット)を読み取ることで、メタデータ駆動のオーケストレーションを使用します。タスクを動的に生成することで、コードの変更をデータ駆動の変更から分離します。
シークレット、サービスアカウント、最小権限
- 認証情報はSecret Managerに保存し、実行時に参照します。コードやAirflowの変数にシークレットを埋め込むことは避けます。
- パイプラインごとに個別のサービスアカウントを割り当て、必要最小限のIAMロールを付与します。規制対象のBigQueryアクセスについては、クライアントデータを個別のデータセットに分離し、承認されたユーザーにのみデータセット固有のロールを付与し、BigQuery APIへのアクセスを承認されたプリンシパルに制限します。マルチテナンシーの場合は、クライアントごとにデータセットを作成し、適切なロールのみをバインドします。
CI/CDとInfrastructure as Code
- インフラストラクチャ(Composer環境、Workflows、Schedulerジョブ、Pub/Subトピック、ログシンク)はTerraformで管理します。モジュールを使用して、プロジェクト/環境、シークレット、サービスアカウントを標準化します。
- パイプラインのコードはCloud BuildやGitHub Actionsでビルドおよびテストします。単体テスト、SQLリンティング、Dataformのドライラン、Airflow DAGの検証を自動化します。アーティファクトはタグを介して昇格させます。Composerの場合は、DAGをデプロイ可能なバンドルとしてパッケージ化します。Dataformの場合は、アサーションがパスした後に昇格するリリースブランチを使用します。
- デプロイの昇格: dev → test → prodは、個別のプロジェクトとパラメータ化された構成を介して行います。リスクの高い昇格には、手動承認ゲートと変更ウィンドウを備えた継続的デリバリーを使用します。
可観測性、アラート、ランブック
テレメトリとアラート
- すべてのオーケストレーションログを、構造化フィールド(pipeline, dag_id, run_id, task_id, partition)と共にCloud Loggingにルーティングします。エラーログはログベースの指標を介してMonitoringにエクスポートします。以下についてアラートを設定します:
- スケジュールの失敗やSLA違反
- 連続したタスクの失敗
- バックログの増加(例:Pub/Subの未確認応答メッセージ、Dataflowのシステムラグ)
- データ品質アサーションの失敗
- Cloud Composer:DAG/タスクの実行時間、成功率、キューの深さ、スケジューラの健全性を監視します。ページング(担当者呼び出し)と修復ランブックのために
on_failure_callbackを設定します。 - Cloud Workflows:実行ログとステップのレイテンシを調査し、明示的なリトライとエラーハンドラを追加し、相関ID付きのカスタムログを出力します。
- BigQueryテーブル変更通知:特定のテーブルを対象とする挿入ジョブ用の高度なフィルタを持つプロジェクトレベルのLoggingシンクを作成し、Pub/Subにエクスポートします。監視ツールはそのトピックをサブスクライブし、他のテーブルからのノイズなしに即時アラートを受け取ります。
ランブックの設計
- 各パイプラインについて、トリガー、依存関係、SLA、ロールバック/リトライ手順、安全なバックフィル手順を文書化します。Dataflowの「固定データセットのリプレイ」、ストリーミングジョブのドレイン方法、失敗したパーティションの再処理方法、DLQメッセージの修復方法を含めます。
- 一般的な障害シグネチャ(例:権限拒否、クォータ超過、スキーマの不一致)を、デシジョンツリーとエスカレーションパスと共にキャプチャします。
実践的な問題シナリオ
Acme Retail Analytics社は、時折不正な形式の行を含むパートナーからの日次CSV提供を取り込み、有効なデータを変換してBigQueryにロードし、調査のために不良な行を表面化させる必要があります。また、ほぼリアルタイムの価格更新のためのイベント駆動型エンリッチメントと、開発環境から本番環境への安全なプロモーションも望んでいます。
アプローチ:
ストレージとイベントトリガー
- オブジェクトのバージョニングと均一なバケットレベルのアクセスが有効な専用のCloud Storageバケットを作成します。Eventarcを介してPub/Subへのオブジェクトファイナライズ通知を有効にします。
- 理由:オブジェクトのファイナライズは、下流の取り込みをトリガーするための信頼性の高いイベントです。バージョニングは再実行と監査をサポートします。
バッチ取り込みとデッドレター処理
- Cloud Composerを使用して、catchupを有効にした日次のAirflow DAGを02:00にスケジュールします。DAGはDataflowバッチジョブを起動し、CSVを解析し、スキーマを検証し、有効なレコードを決定論的なステージングテーブルに書き込み、その後パーティション化されたターゲットテーブルに
MERGEします。不正な形式/失敗したレコードは、BigQueryのデッドレターテーブルにルーティングします。 - 理由:Dataflowは解析/検証をスケーリングします。
MERGEはべき等性を保証します。デッドレターのキャプチャは、パイプラインをブロックすることなく調査をサポートし、不正な形式の行に対する推奨パターンに一致します。
- Cloud Composerを使用して、catchupを有効にした日次のAirflow DAGを02:00にスケジュールします。DAGはDataflowバッチジョブを起動し、CSVを解析し、スキーマを検証し、有効なレコードを決定論的なステージングテーブルに書き込み、その後パーティション化されたターゲットテーブルに
イベント駆動型のエンリッチメント
- 増分的な価格更新のための軽量なエンリッチメントを実行するために、Cloud Runジョブをデプロイします。日中に小さな更新ファイルが到着した際に、EventarcからのPub/SubメッセージをリッスンしているCloud Workflowsを介してトリガーします。
- 理由:Workflowsを備えたサーバーレスコンテナは、重い変換処理をバッチに残しつつ、小さなイベントに対して低レイテンシ、低運用負荷のオーケストレーションを提供します。
信頼性の制御
- DataflowおよびCloud Runジョブの一時的な障害に対して、指数バックオフ付きのリトライを設定し、合計リトライ時間をDAGのSLAに制限します。Airflowではタスクごとの実行タイムアウトと
on_failureコールバックを設定し、Workflowsではmax_doublingsとmax_retry_durationを設定します。 - 理由:上限のあるバックオフはSLAを維持し、リトライの暴走を防ぎます。
- DataflowおよびCloud Runジョブの一時的な障害に対して、指数バックオフ付きのリトライを設定し、合計リトライ時間をDAGのSLAに制限します。Airflowではタスクごとの実行タイムアウトと
セキュリティと最小権限
- 各コンポーネントを専用のサービスアカウント(ComposerオーケストレータSA、DataflowワーカーSA、Cloud RunジョブSA)で実行します。必要なロールのみを付与します:Dataflowに取り込みバケットへのGCS読み取り権限、ターゲットデータセットへのBigQuery
dataEditor権限、ログへのViewer権限。シークレットはSecret Managerに保存し、実行時に参照します。 - 理由:最小権限の原則を強制し、影響範囲(ブラスト半径)を隔離します。
- 各コンポーネントを専用のサービスアカウント(ComposerオーケストレータSA、DataflowワーカーSA、Cloud RunジョブSA)で実行します。必要なロールのみを付与します:Dataflowに取り込みバケットへのGCS読み取り権限、ターゲットデータセットへのBigQuery
メタデータ駆動のオーケストレーション
- パートナーソース、ファイルパターン、ターゲットデータセットをリストしたBigQueryコントロールテーブルを維持します。DAGの実行時に、Airflowはこのテーブルをクエリし、動的タスクマッピングを使用してパートナーごとのタスクを生成します。
- 理由:パートナーの追加がコードの変更ではなくデータの変更となり、デプロイのリスクを低減します。
可観測性とアラート
run_idとpartner_idを含む構造化ログを出力します。DAGのSLA違反、Dataflowのシステムラグ、空でないデッドレター数に対するアラートポリシーを作成します。ターゲットテーブルへのBigQuery挿入については、そのテーブル用の高度なフィルタを持つCloud Loggingシンクを、Acme社の監視ツールが利用するPub/Subトピックに設定します。- 理由:きめ細かいアラートにより、ノイズのない迅速なトリアージが可能になります。
CI/CDとプロモーション
- インフラ(バケット、Pub/Sub、Eventarc、Composer、Workflows、BigQueryデータセット)をTerraformで管理します。Cloud Buildを使用してAirflow DAGの構文を検証し、単体テストを実行し、開発用のComposer環境にデプロイします。Dataformアサーションと統合テストが合格した後、パラメータ化された設定と手動の承認ゲートを使用して、テスト環境および本番環境にプロモートします。
- 理由:宣言的で再現可能なデプロイと、環境間での安全なプロモーションを実現します。
ランブックとリカバリ
- 特定の日付をリプレイする手順を文書化します:オブジェクトのバージョニングからCSVを復元し、そのパーティションに対してDataflowジョブを再実行し、結果を
MERGEし、DLQレコードを確認します。不一致が発生した場合に変換処理のバグを特定するための「固定データセットのリプレイ」手順を含めます。 - 理由:べき等な設計と文書化されたリカバリ手順により、部分的な障害からの修復が効率化されます。
- 特定の日付をリプレイする手順を文書化します:オブジェクトのバージョニングからCSVを復元し、そのパーティションに対してDataflowジョブを再実行し、結果を
← データの取り込み、統合、移行 · すべてのドメイン · 機械学習、AI、データサービング →
これらの問題を練習する → · ExamRoll.ioで時間制限付き練習 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
試験に合格する →