Google PDE: データの取り込み、統合、移行 — 学習ガイド
こちらの一部です: Google Professional Data Engineer — 学習ガイド. 検証済みの解答で練習: Google試験ハブ, または時間制限付き模擬試験に挑戦: ExamRoll.io.
概要
Google Cloudにおけるデータの取り込み、統合、移行は、多様なソースシステムを信頼性が高くクエリ可能なデータセットに変換するための、再現可能なパターン、マネージドサービス、運用管理に及びます。効果的な設計では、転送と変換を分離し、プロデューサーとコンシューマーを分離(デカップリング)し、明確なリネージと検証機能を備えた、べき等でチェックポイントが設定されたパイプラインが推奨されます。このセクションでは、取り込みパターン、データ移動とCDCのためのGoogle Cloudサービス、スキーマとデータ品質の管理、接続性とハイブリッド統合、そしてカットオーバーストラテジーについて、設計上のトレードオフと障害モードを随所で指摘しながら解説します。
取り込みパターンとワークロード
- バッチ取り込み: 定義された間隔での定期的なプルまたはファイルドロップ。予測可能なコストとバックフィルに適しています。障害モード: 大規模で頻度の低いバッチは、リソースのスパイク、長いキャッチアップ期間、SLA違反を引き起こします。緩和策: バッチウィンドウのサイズを適正化し、時間またはキーでシャーディングし、並列処理を利用します。
- 一括ロード: 1回限りまたは大規模なロード(例:初回の履歴バックフィル)。カラムナフォーマットまたは自己記述型フォーマット(Parquet、Avro)が望ましく、分析ストレージ(BigQuery)に直接ロードするか、Cloud Storageにステージングします。トレードオフ: 外部テーブルクエリはロードステップを回避しますが、コストをクエリ時のスキャンにシフトします。
- 増分ロード: タイムスタンプまたはハイウォーターマークによる定期的な差分ロード。堅牢な重複排除とべき等なアップサートが必要です。障害モード: クロックのずれや遅延到着レコード。サーバーサイドのコミットタイムスタンプとウォーターマークを使用します。
- 変更データキャプチャ(CDC): オペレーショナルデータベースからの挿入、更新、削除の継続的なレプリケーション。ニアリアルタイム分析や低ダウンタイムでの移行に最適です。トレードオフ:
- 順序性: ほとんどのCDCツールはトランザクション内、そして通常はシャード内の順序を保持しますが、グローバルなクロスシャードでの順序は保証しません。トランザクションのコミットタイムスタンプと主キーを使用してシーケンスを再構築します。
- 配信セマンティクス: at-least-once(少なくとも1回)が一般的です。べき等なシンクを構築するか、一意の変更IDを使用して重複排除します。
- スナップショット + CDC: 一貫性のあるスナップショットから開始し、その後、正確なログシーケンスからの変更を適用して、ダウンタイムなしでパリティに到達します。
リレーショナル、SaaS、オンプレミス、ファイルソース:
- リレーショナルソース: ネイティブCDCまたはタイムスタンプカラムを使用します。一括ロードの場合は、Avro/Parquetにエクスポートし、Cloud Storageにステージングします。
- SaaSソース: 増分トークンを持つベンダーAPIが望ましいです。マネージドコネクタ(例:Data Fusion内)を介して統合します。レート制限のためにスロットリングを行い、スキーマドリフトに対応します。
- オンプレミスソース: エージェントベースの転送、VPN/Interconnect + Private Google Access、またはTransfer Applianceによるオフラインシーディングから選択します。
- ファイル取り込み: 多数の小さなファイルがある場合は、バンドル(例:tar)してRPCオーバーヘッドを削減します。
gsutil -mまたは並列化されたクライアントを使用し、分析用に、より大きなカラムナファイルに合成または変換します。
取り込み、統合、移行のためのGoogle Cloudサービス
- Datastream(サーバーレスCDC): MySQL、PostgreSQL、Oracleからの変更をキャプチャし、Cloud Storage、BigQuery(テンプレート経由)、またはPub/Subに取り込みます。トランザクション境界とコミットメタデータを保持しますが、グローバルな順序は保証されません。ダウンストリームでキーとコミットタイムスタンプによって順序を適用します。at-least-once配信を想定し、べき等なコンシューマー(例:変更IDを使用したBigQuery MERGE)を設計します。
- Database Migration Service (DMS): ネイティブレプリケーションを使用した、最小限のダウンタイムでのデータベース移行のためのサービスです。DMSは一貫性のあるスナップショットを作成し、その後GTID/LSN/SCNを使用して変更を継続的にレプリケートします。任意の変換ではなく、リフトアンドシフト専用に構築されています。分析目的の場合、必要に応じてDMSをDataflowまたはData Fusionで補強します。
- Cloud Data Fusion: リレーショナル、SaaS、ファイル、メッセージングシステムへのコネクタを備えたマネージド統合サービスです。変換ステージ(結合、集計、フォーマット変換、カスタムWranglerレシピ)を持つパイプラインを構築し、ソースとフィールドを横断してリネージをキャプチャします。運用面では、スケジューリング、リトライ、メトリクスの発行を行います。ノーコード/ローコードのELT/ETLやコネクタ管理の一元化にData Fusionを使用します。
- Storage Transfer Service (STS): AWS S3、Azure Blob、オンプレミス(エージェント使用)、SFTP、URLリストからCloud Storageへの転送を管理・スケジューリングします。マニフェスト、増分同期、帯域幅制御、チェックサムによる整合性をサポートします。障害モードには、小さなファイルの非効率性やAPIスロットリングが含まれます。バッチ処理や調整可能な同時実行数で緩和します。
- Transfer Appliance: ネットワーク帯域幅が限られている場合や、データが機密性が高く長時間の転送に適さない場合に、数テラバイトからペタバイト規模の初期シーディングを行うための、オフラインの暗号化アプライアンスです。管理履歴(Chain-of-custody)と暗号化が組み込まれています。シーディング後、差分データについてはSTSまたはCDCでフォローアップします。
- Cloud Pub/Sub + Dataflow: Pub/Subは、ストリーミングまたはマイクロバッチパターンのためにプロデューサーとコンシューマーを分離します。Dataflowは、チェックポイントとウォーターマークを備えた、自動スケーリングされるステートフルなストリーム/バッチ処理を提供します。デフォルトストリームごとにexactly-once保証を提供する低レイテンシのストリーミングにはBigQuery Storage Write APIを使用します。それ以外の場合は、insertIdによる重複排除セマンティクスに依存します。
HadoopからDataprocへの移行では、GCSコネクタを使用してデータをCloud Storageに保存し、エフェメラルクラスタまたは自動スケーリングクラスタを使用することで、Persistent Diskの使用を最小限に抑えます。これにより、処理のためにHDFS互換のセマンティクスを維持しつつ、大規模なブロックストレージのコストを回避できます。
境界におけるスキーマ、検証、データ品質
- スキーマのマッピングと型変換:早い段階で厳密に型付けされたスキーマに標準化します。AvroやParquetはスキーマを保持し、クリーンに進化させることができます。BigQueryでは、スキャンコストを削減するために、パーティション化テーブルとクラスタ化テーブルを優先的に使用します。 例:日次分析用のパーティション化テーブルを作成する
undefined
- 不正な形式のレコードの処理:リジェクトされたレコードは、デッドレターキュー(Pub/Sub)またはCloud Storageの隔離バケットにルーティングします。Dataflowではサイドアウトプットを、Data Fusionではエラーコレクタを使用します。トリアージのために、サンプルペイロードとスキーマバージョンとともに解析エラーをログに記録します。
- 検証:永続化する前に境界チェックを実行します:
- 構造的:スキーマへの準拠、必須フィールド、データ型、enumドメイン。
- 参照:キャッシュされたディメンションルックアップによる外部キーの存在確認。
- 妥当性:タイムスタンプの範囲、ジオフェンス、負でない金額。
- 一意性:主キーまたは複合キーの衝突。
- べき等な読み込み:決定論的なキーとupsert操作を使用します。BigQueryでは、自然キーまたは代理変更キーを使用してMERGEを実装します。 例:
undefined
- ウォーターマークと遅延:ストリーミングパイプラインでは、イベントタイムのウォーターマークと許容される遅延を設定し、完全性とレイテンシのバランスをとります。遅延データは修正パスにルーティングされるか、バックフィルをトリガーします。
- 照合:ソースからシンクまで、パーティション/ウィンドウごとの行数とチェックサムを追跡します。CDCのログポジション(LSN/SCN)とコミットタイムスタンプをキャプチャし、制御テーブルに保存して継続性を証明し、ギャップを特定します。
接続性、信頼性、および運用
ネットワーク接続とプライベートアクセス:
- ハイブリッド:プライベート接続にはCloud VPNまたはDedicated/Partner Interconnectを使用します。Cloud StorageなどのGoogle APIへのプライベートアクセスには、Private Google AccessまたはPrivate Service Connectを有効にします。
- セキュリティ:ワークロードIDにはサービスアカウント、最小権限のIAM、データ漏洩防止にはVPC Service Controls、必要に応じてCMEKを使用します。
- スループット:クライアント側で並列処理をスケールさせますが、最終的には帯域幅がスループットを決定します。大規模な転送には、最初のバルク転送にはTransfer Applianceを、その後の増分更新にはSTSまたはCDCを優先します。
チェックポイントとバックプレッシャー:Dataflowはチェックポイントと自動スケーリングを管理します。バーストを吸収できるシンク(Cloud Storageへのバッファリング、BigQueryへのバッチ書き込み)を設計します。Pub/Subでは、メッセージの再配信ストームを防ぐために、フロー制御とackデッドラインを調整します。
CDCによる順序性と一貫性:
- Datastreamはトランザクション内の順序を保持し、コミットメタデータを生成します。コンシューマーはコミットタイムスタンプを使用してキーごとの順序を再構築します。at-least-onceを想定し、べき等性を構築します。
- DMSは、ネイティブログを使用してスナップショットとレプリケーションのカットオーバー全体でデータベースの一貫性を保証します。段階的なカットオーバーには、リードレプリカまたはデュアルライト戦略を使用します。
分析のためのファイル戦略:大規模で複数のエンジンからアクセスする場合、正規データをCloud Storageに保存し、コスト効率が良い場合はアドホッククエリ用に永続的な外部テーブルを公開します。本番分析では、クエリごとのスキャンコストを最小限に抑えるために、パーティション化されたBigQueryテーブルにロードします。
小規模ファイルの最適化:転送前に小さなファイル(例:tarあたり約1,000個)をバンドルし、クラウドで展開します。並列gsutilとライフサイクルルールを使用して、ステージングアーティファクトを階層化し、有効期限切れにします。
運用上の落とし穴と緩和策:
- SaaSからのスキーマドリフト:Data Fusionでスキーマの進化を有効にし、互換性を強制します。互換性のない変更についてアラートを出します。
- タイムゾーンとエンコーディング:入力時にUTC、UTF-8に正規化します。
- CDCのギャップ:ソースのログ保持期間を監視し、レプリカラグが保持期間の制限に近づいたらアラートを出します。
- クォータ:BigQueryストリーミングインサート、APIレート制限。制限に近づいたらバッチ処理を行います。
カットオーバー、バックフィル、および検証
- カットオーバー計画:
- ビッグバン方式: 短時間の凍結、一度の切り替え。運用上の複雑さは最も低いが、ロールバックが必要な場合のリスクは最も高い。
- 段階的またはブルー/グリーン方式: ミラーリングされた書き込みによる並行稼働、段階的なトラフィックシフト、シャドウリード。コストは高くなるが、ロールバックはより安全。
- バックフィル:
- 初期の一括ロード(Transfer ApplianceまたはSTS)をAvro/Parquetを使用して実行し、スキーマを保持する。ロード中にパーティション化とクラスタリングを行い、手戻りを避ける。
- 一括転送中の差分をキャプチャするため、スナップショットと同時に既知のログポジションからCDCを開始する。本番稼働前に共通のウォーターマークで整合性を取る。
- ロールバック:
- 検証中はレガシーシステムを読み取り専用で維持する。デュアルライトのシナリオでは、書き込みを機能フラグの背後に配置して迅速に元に戻せるようにする。必要に応じてCDCの変更をリプレイまたは巻き戻すために、一貫性のあるチェックポイントを保持する。
- 移行の検証:
- 構造的: 行数とパーティションごとのチェックサムが一致し、スキーマと制約が同等であること。
- 時間的: スナップショットの境界からカットオーバーまでにギャップがなく、CDCのポジションが連続していること。
- ビジネスパリティ: ウィンドウ期間にわたる集計値とKPIを比較し、受け入れクエリを実行する。
- パフォーマンス: 取り込みスループット、クエリのレイテンシ、コストを予算と照らし合わせて検証する。
実践的な問題シナリオ
Northstar Retail社は、ニアリアルタイム分析と機械学習を強化するため、オンプレミスのOracleおよびMySQLトランザクションシステム、SaaS CRMイベント、日次のCSVドロップといったグローバルに混在するデータをGoogle Cloudに統合する必要があります。また、レガシーなHadoopクラスターを、高額なブロックストレージ費用を発生させることなく移行し、ゼロから低ダウンタイムでのカットオーバーを達成する必要もあります。
- 安全なハイブリッド接続の確立
- 主要な帯域幅にはPartner Interconnectを使用し、フォールバックとしてCloud VPNを利用する。オンプレミスのワークロードがCloud StorageとPub/Subにプライベートにアクセスできるよう、Private Google Accessを有効にする。 理由: プライベートパスは下り(egress)の露出とレイテンシを最小限に抑え、Private Google Accessはセキュリティポリシーを満たしつつパブリックIPの要件を回避します。
- 履歴データの効率的な初期投入(シーディング)
- 800 TBの履歴HDFSデータについては、Transfer Applianceを使用してCloud Storageにコピーします(初期一括ロード)。初期投入後、カットオーバーまでオンプレミスのNFSエクスポートからStorage Transfer Serviceを毎日実行し、変更を取得します。 理由: Transfer Applianceは長時間のネットワーク飽和を回避し、STSはスケジュール化されたチェックサム付きの増分同期を提供します。GCSコネクタを使用してCloud Storageに保存することで、ノードごとに50 TBのPersistent Diskを必要とせずにDataprocでの処理が可能になります。
- CDCによる業務データベースの移行
- DMSを使用して、最小限のダウンタイムでMySQLとPostgreSQLを移行します。Oracleから分析システムへのCDCには、Datastreamを使用してCloud Storageにランディングし、その後Google提供のDataflowテンプレートを使ってBigQueryにロードします。 理由: DMSは信頼性の高いスナップショット+継続的同期のためにネイティブレプリケーションを活用します。Datastreamはコミットメタデータ付きのサーバーレスCDCを提供し、Dataflowテンプレートは順序付けされたべき等のBigQuery書き込みを保証します。
- SaaSおよびファイルベースのフィードの取り込み
- Cloud Data Fusionパイプラインを構築し、SaaSコネクタを使用して増分トークン付きのCRMイベントを取り込み、ファイルパイプラインでベンダーのSFTPからSTS経由で日次CSVを取り込みます。キュレートされたCloud StorageバケットでAvroに正規化し、その後パーティション化されたBigQueryテーブルにロードします。 理由: Data Fusionはコネクタ、変換、リネージを一元管理します。Avroに標準化することでスキーマが保持され、進化が容易になります。パーティション化されたBigQueryテーブルはクエリコストを削減します。
- リアルタイムイベントのストリーミング
- ウェブおよび店舗のイベントをPub/Subにパブリッシュします。Dataflowで処理し、パース、検証、エンリッチメント、ウォーターマーキングを行います。Storage Write API経由でBigQueryに書き込み、生のAvroをCloud Storageにアーカイブします。 理由: Pub/Subはプロデューサーとコンシューマーを分離します。Dataflowはオートスケーリング、ステートフル処理、チェックポイント、遅延データ処理を提供します。デュアルライトにより、低レイテンシの分析と耐久性のある生データの保持の両方が保証されます。
- 境界でのデータ品質とスキーマの制御の徹底
- Dataflow/Data Fusionにスキーマレジストリと検証を実装します。不正な形式のレコードはGCSの隔離バケットとPub/Subのデッドレタートピックにルーティングします。ドメインチェック(例: 通貨コード、UTCタイムスタンプ)を適用し、複合キーを使用して重複排除を行います。 理由: 早期の拒否と隔離により、不正なデータの伝播を防ぎます。べき等性と重複排除は、CDCおよびストリーミングソースからのat-least-once配信から保護します。
- 分析ストレージとアクセスの最適化
- キュレートされたデータセットを、パーティション化およびクラスタリングされたBigQueryテーブルにロードします。生のアーカイブは、低頻度の探索用に永続的な外部テーブルとして公開します。トランザクション処理が残るOLTPワークロードについては、リードレプリカを持つCloud SQLを維持します。 理由: パーティショニングとクラスタリングはスキャンコストを最小限に抑えます。外部テーブルは、時折のアクセスのための不要なロードを回避します。Cloud SQLはトランザクションアプリのACIDセマンティクスを保持します。
- カットオーバー、バックフィル、ロールバックの計画
- 各RDBMSに対してスナップショット+CDCを実行し、行数とチェックサムが一致する調整ポイントに到達させます。48時間、デュアルライトによるブルー/グリーン方式で稼働させ、読み取りを徐々にBigQueryにシフトします。不一致が検出された場合に書き込みを元に戻すための機能フラグを維持します。 理由: ブルー/グリーン方式はリスクを低減します。既知のウォーターマークでの検証は完全性を保証し、フラグは迅速なロールバックを可能にします。
- 検証と可観測性
- ソースのLSN/SCN、コミットタイムスタンプ、行数、パーティションごとのチェックサムをキャプチャするコントロールテーブルを構築します。Datastreamのラグ、DMSのレプリケーション状態、Dataflowのウォーターマーク、Pub/Subのバックログ、STSのジョブステータス、BigQueryのストリーミング挿入メトリクスを監視します。 理由: エンドツーエンドのリネージと定量的な管理は、正確性の監査可能な証明と、ギャップやラグに関するタイムリーなアラートを提供します。
ランディング、キュレーション、サービングの各レイヤーを分離し、Cloud Storageを耐久性のある低コストのステージングおよびアーカイブとして使用し、べき等なコンシューマーを持つCDCのためにDMS/Datastreamを活用し、イングレスでスキーマと品質を徹底することにより、Northstar Retail社は、安全でスケーラブルな取り込みと、予測可能なコストで低リスクかつ検証可能な移行を実現します。
← Spark、Dataproc、分散データ処理 · すべてのドメイン · ワークフローのオーケストレーションとパイプラインの自動化 →
これらの問題を練習する → · 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.
試験に合格する →