Google PDE: Dataflow と Apache Beam によるストリーム処理 — 学習ガイド

こちらの一部です: Google Professional Data Engineer — 学習ガイド. 検証済みの解答で練習: Google試験ハブ, または時間制限付き模擬試験に挑戦: ExamRoll.io.

概要

Google Cloudでのストリーム処理は、Dataflowランナーによって実行されるApache Beamの統一されたプログラミングモデルが中心となります。Beamは、PCollectionに対する変換のパイプラインという論理的な抽象化を提供し、これによりコードを並列処理、オートスケーリング、耐障害性といった実行の詳細から分離します。ストリーミングにおいて、正確性は時間セマンティクス(イベント時間 vs 処理時間)、ウィンドウ分割(固定、スライディング、セッション、グローバル)、ウォーターマーク、トリガー、そして遅延データの扱いに依存します。Dataflowでの運用の卓越性を実現するには、適切なワーカーサイジング、オートスケーリングポリシー、ストリーミングエンジン、シャッフルの選択、べき等なシンク設計、デッドレター処理、そして堅牢なオブザーバビリティが必要です。

Apache Beamモデルと時間セマンティクス

障害モードとトレードオフ:

ストリーミングワークロードのためのDataflowの運用

デプロイ、テンプレート、アップグレード戦略

実践的な問題シナリオ

NovaTrack社は、50,000個の温度センサーからグローバルなIoTテレメトリを取り込み、分単位の集計を提供し、生データを永続化し、リアルタイムダッシュボードを公開する必要があります。不正な形式のメッセージや順序が乱れた配信が時折発生することが想定されます。このソリューションは、自動スケーリングし、不良レコードを調査用に提示し、ゼロダウンタイムでのアップグレードをサポートする必要があります。

アプローチ:

  1. 取り込みと時間セマンティクス

    • リージョンごとのPub/Subトピックと、属性としてdeviceIdとeventTs(RFC3339)を持つリージョンごとのパブリッシャーを作成します。可能な場合は、deviceIdによる順序付けキーを有効にします。
    • 理由:Pub/Subは、at-least-once配信を備えた、耐久性があり伸縮自在なイングレスを提供します。エッジでイベントタイムスタンプを付与することで、真のイベント時刻を保持します。デバイスごとの順序付けは、中央のボトルネックなしにデバイス内の順序の乱れを低減します。
  2. イベントタイムウィンドウを使用したDataflowストリーミングパイプライン

    • 専用のサブスクリプションからPubSubIO経由で読み取り、eventTsをBeamタイムスタンプとして抽出します。欠落している場合はpublishTimeにフォールバックします。
    • 1分間のFixedWindowsを適用し、30秒での早期トリガーと、遅延要素ごとに遅延発火を設定します。許容される遅延を10分に設定し、ペインを蓄積します。
    • 理由:イベントタイムウィンドウは正確な分単位の集計を保証します。早期発火はサブ分単位の鮮度でダッシュボードにデータを提供し、遅延発火は遅延データが到着するにつれて集計を修正します。遅延の上限は、状態のサイズとコストを抑制します。
  3. 検証、エンリッチメント、デッドレタールーティング

    • JSONをパースし、スキーマと範囲を検証し、ジョブ開始時にBigQueryからロードされたサイドインプットを介して小さな静的参照データでエンリッチするParDoを実装します。
    • TupleTagsを使用して、有効なレコードをメイン出力に、失敗したレコードをペイロード、エラー、deviceId、パースタイムスタンプを含むデッドレターPCollectionに出力します。DLQはパーティション化されたBigQueryテーブルに書き込みます。
    • 理由:サイドインプットは参照データをメモリ内に保持し、低レイテンシを実現します。デッドレターのキャプチャにより、メインフローをブロックすることなく、不良行の調査と的を絞った再処理が可能になります。
  4. 集計とホットキーの緩和

    • deviceIdでキーを設定し、CombineFnsを使用して分単位のavg/min/maxを計算します。リージョンごとのトップNメトリクスについては、ホットキーを避けるためにregion#Nでシャーディングし、その後再集計します。
    • 理由:Combinerはシャッフル量とコストを最小限に抑えます。キーシャーディングは、リージョンごとのファンイン中に単一キーのボトルネックを防ぎます。
  5. シンクとexactly-once効果

    • 検証済みの生イベントと分単位の集計を、Storage Write APIを使用したBigQueryIOでBigQueryに書き込みます。カスタムリトライでのべき等性を確保するために、deviceId + eventTsに基づいて安定した挿入IDを設定します。
    • 理由:Storage Write APIは、ストリーム内でexactly-onceセマンティクスを備えた、高スループット、低レイテンシの取り込みを提供します。安定したIDは、リプレイが発生した場合にダウンストリームでの重複排除を保証します。
  6. ダッシュボードの一貫性戦略

    • ダッシュボードは、ウォーターマークに対して2分間のルックバック、またはストリーミングデータに対して観測された可用性レイテンシの2倍の固定遅延で、パーティション化された集計テーブルをクエリします。
    • 理由:BigQueryのストリーミング可視性は結果整合性です。読み取りをわずかに遅らせることで、ほぼリアルタイムの動作を維持しつつ、処理中の行が欠落するのを防ぎます。
  7. 運用:オートスケーリングとストリーミングエンジン

    • Streaming Engineを有効にします。予想されるピーク(例:平均の3倍)に基づいてmaxWorkersを設定し、CPUバウンドのパースと暗号化に合わせてマシンタイプを選択し、一時的なシャッフルに対応するためにブートディスクを増やします。
    • ウォーターマークの遅延、バックログ秒数、CPU、ステップごとのスループットを監視します。持続的な遅延やDLQレートの急上昇に対してアラートを設定します。
    • 理由:Streaming Engineは、状態/シャッフルを外部化することで、弾力性とより簡単なアップグレードを実現します。適切なサイジングと監視により、サイレントなSLO違反を防ぎます。
  8. Flex Templatesによるデプロイとアップグレード

    • パイプラインを、入力サブスクリプション、出力テーブル、DLQテーブル、maxWorkers、リージョンなどのパラメータを持つFlex Templateとしてパッケージ化します。互換性のない変更の場合は、新しいサブスクリプションで同じトピックをターゲットとする新しいパイプラインを開始し、出力を検証してから古いジョブをドレインします。オプションで、Pub/Subスナップショットを作成し、新しいサブスクリプションをそのスナップショットにシークさせることで、ギャップがないことを保証します。
    • 理由:Flex Templatesは、再現可能でパラメータ化されたデプロイを可能にします。ドレインを伴う検証済みのブルー/グリーン方式のカットオーバーは、ゼロデータ損失と最小限のダウンタイムを実現します。
  9. 再処理とバッチバックフィル

    • サイドアウトプットを介して、圧縮されたAvro形式の生イベントファイルをCloud Storageに保存します。モデルやスキーマが変更された際には、バッチDataflowパイプラインを実行してBigQueryへのバックフィルや再処理を行います。
    • 理由:耐久性のある生データのアーカイブは、ホットパスに影響を与えることなく、再現性とスキーマの進化をサポートします。

この設計は、グローバルスケールで順序が乱れたデータや遅延データを処理しながら、コストを抑制し、明確なエラー分離、強力な可観測性、安全なアップグレードパスを備えた、正確で低レイテンシの集計を実現します。


BigQuery による分析とウェアハウスエンジニアリング · すべてのドメイン · メッセージング、イベントの取り込み、リアルタイムサービス

これらの問題を練習する → · 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.

試験に合格する →

Googleを閲覧 →

Related guides

オールインワンアクセス

1つのサブスクリプション。すべての試験。

すべてのプランで、無制限の回答検索、模擬試験、AI解説、および完全なリソースライブラリが利用可能 — 20以上の言語に対応。

月額
24.87
Just €0.83/day
すべて含まれています:
  • 無制限の回答検索
  • 無制限の模擬試験
  • AIを活用した解説
  • 完全なリソースライブラリ
  • 20以上の言語
  • 毎週のコンテンツ更新
  • 特典 & 紹介
  • 優先サポート
無料トライアルを開始

クレジットカード不要*

ベストバリュー
12ヶ月
179.87
Just €0.49/daySave 40%
すべて含まれています:
  • 無制限の回答検索
  • 無制限の模擬試験
  • AIを活用した解説
  • 完全なリソースライブラリ
  • 20以上の言語
  • 毎週のコンテンツ更新
  • 特典 & 紹介
  • 優先サポート
無料トライアルを開始

クレジットカード不要*

✓ 無料プランが含まれています · ✓ いつでもキャンセル可能 · ✓ すべてのプランで製品の全機能が利用可能