Google PDE: Dataflow と Apache Beam によるストリーム処理 — 学習ガイド
こちらの一部です: Google Professional Data Engineer — 学習ガイド. 検証済みの解答で練習: Google試験ハブ, または時間制限付き模擬試験に挑戦: ExamRoll.io.
概要
Google Cloudでのストリーム処理は、Dataflowランナーによって実行されるApache Beamの統一されたプログラミングモデルが中心となります。Beamは、PCollectionに対する変換のパイプラインという論理的な抽象化を提供し、これによりコードを並列処理、オートスケーリング、耐障害性といった実行の詳細から分離します。ストリーミングにおいて、正確性は時間セマンティクス(イベント時間 vs 処理時間)、ウィンドウ分割(固定、スライディング、セッション、グローバル)、ウォーターマーク、トリガー、そして遅延データの扱いに依存します。Dataflowでの運用の卓越性を実現するには、適切なワーカーサイジング、オートスケーリングポリシー、ストリーミングエンジン、シャッフルの選択、べき等なシンク設計、デッドレター処理、そして堅牢なオブザーバビリティが必要です。
Apache Beamモデルと時間セマンティクス
パイプライン、変換、PCollection、ランナー:
- Beamパイプラインは、PTransformの有向非巡回グラフをPCollection(有界または非有界)に適用します。
- ランナー(Dataflow、Spark、Flink、Direct)がパイプラインを実行します。Dataflowは、マネージドのオートスケーリング、チェックポイント作成、運用上の可視性を提供します。
- 変換には、要素ごとの処理(ParDo)、グループ化と結合(GroupByKey、Combine)、ジョイン(CoGroupByKey)、IO(PubSubIO、BigQueryIO、FileIO)などがあります。
ウィンドウ:
- 固定ウィンドウ:重複しないスライス(例:1分間のタンブリングウィンドウ)で、定期的な集計に使用します。
- スライディングウィンドウ:重複するウィンドウで、スムーズな移動メトリクス(例:1分ごとにスライドする5分間のウィンドウ)に使用します。
- セッションウィンドウ:非アクティブな期間が一定時間続くと閉じる動的なウィンドウで、ユーザーセッションやデバイスのバースト処理に最適です。
- グローバルウィンドウ:非有界ストリーム全体に対する、ウィンドウ分割されていないデフォルトのビューです。多くの場合、定期的な実体化のためにトリガーと組み合わせて使用されます。
イベント時間 vs 処理時間:
- イベント時間:イベントがソースで発生した時刻。転送レイテンシが変動しても、論理的に一貫した集計を可能にします。
- 処理時間:イベントがパイプラインによって観測された時刻。運用上のトリガーには役立ちますが、セマンティックな正確性の担保には適していません。
ウォーターマーク:
- ウォーターマークは、イベント時間の完全性(ランナーが時刻Tまでのすべてのイベントを観測したという推測)を推定するものです。
- ウォーターマークは、バックプレッシャーやソースの遅延により不規則に進んだり停滞したりすることがあります。遅延データとは、ウォーターマークより前のタイムスタンプで到着するあらゆるデータのことです。
トリガーと遅延:
- デフォルト:ウォーターマークがウィンドウの終端を通過したときに発火するAfterWatermarkトリガー。許容される遅延が0の場合、遅延データは破棄されます。
- 早期発火(処理時間ベースまたはカウントベース)により、低レイテンシの予備的な結果が得られます。
- 遅延発火により、遅延データが到着した際の修正が可能になります。蓄積モードは、ペインが結果を蓄積するか、以前の出力を破棄するかを制御します。
- 許容される遅延は、ビジネス上の許容度とストレージ/コンピューティングのトレードオフに基づいて選択します。遅延を許容するほど、ステートの保持期間が長くなり、コストが増加します。
ステートフル処理、タイマー、セッション化、重複排除:
- ステートフルなDoFnは、キーごとのステート(例:最後に観測されたイベント、実行中の集計値)を保持し、タイマーを設定してステートを出力またはクリアします。
- セッション化は、SessionWindowsによって自然に表現されます。カスタムロジックのためには、キー付きステートと処理時間/イベント時間タイマーを使用します。
- 重複排除:イベントごとに安定したIDを使用し、ウィンドウごとのDistinct/Combine、またはキーごとのステート(例:ブルームフィルタやTTL付きのセット)のいずれかを用います。メモリと偽陽性のトレードオフと、厳密な正確性との間でバランスを取ります。
障害モードとトレードオフ:
- ビジネスメトリクスに処理時間ウィンドウを使用すると、スパイクやリトライ時にドリフト(ずれ)が発生します。イベント時間ウィンドウの使用が推奨されます。
- ウィンドウが小さすぎ、早期トリガーが頻繁に発生すると、過剰なペインの出力とシンクへの書き込み増幅を引き起こします。
- 許容される遅延を無制限にすると、ステートが肥大化する可能性があります。常にステートのTTLを制限し、休眠中のキーをクリアするためのタイマーを設定してください。
ストリーミングワークロードのためのDataflowの運用
ワーカーのサイジングとオートスケーリング:
- 水平オートスケーリングは、バックログ、ウォーターマークの遅延、CPU、スループットに基づいてワーカーを追加/削除します。スパイクを吸収するために、適切なmaxWorkersを設定します。
- ボトルネックに応じてマシンタイプを選択します:CPUバウンド(より多くのvCPU)、メモリバウンド(高メモリタイプ)、ネットワークバウンド(より大きなVMでシャッフルのオーバーヘッドを削減)。
- 大量のシャッフルやファイルベースのシンクには、ブートディスクを増やします。システムラグとバックログ秒数を監視します。
Streaming Engineとシャッフル:
- Streaming Engineは、状態とシャッフルをサービスバックエンドに外部化し、弾力性を向上させ、ワーカーのメモリプレッシャーを軽減し、より高速な更新を可能にします。
- バッチ処理の多いステージや大規模なキーグルーピングには、Dataflow Shuffleを使用してワーカーからシャッフルI/Oをオフロードします。どちらもホットワーカーの障害やディスクスラッシングを削減します。
バックプレッシャー、ホットキー、スキュー:
- Dataflowは動的な作業の再分散によってバックプレッシャーを管理しますが、それでも、該当する場合はソースのフロー制御(例:Pub/Subの未処理のメッセージ/バイト数)を調整します。
- ホットキー(例:人気のID)はストラグラー(処理の遅延する要素)を生み出します。キーシャーディング(key#N)、部分的な事前集約とその後の再キーイング、またはスケッチベースの近似によって緩和します。
- 外れ値レコード(巨大なペイロード)やバースト的なパブリッシャーからのスキューには、パブリッシャーごとのパーティション分割、バッチ処理、または圧縮が必要になる場合があります。
Pub/Subとの統合:
- 取り込みにはPub/Subトピックを使用します。メタデータ(例:deviceId、イベントタイムスタンプ)のためにメッセージ属性を有効にします。
- PubSubIOで取り込みます。イベントタイムスタンプを属性またはペイロードから抽出し、それ以外の場合は発行時刻にフォールバックします。
- 順序付けキーはキーごとの順序を提供します。at-least-once(少なくとも1回)配信のため、Dataflowは依然として下流でのべき等な動作を必要とします。
BigQueryへのストリーミングパターン:
- 高スループット、低レイテンシー、そしてストリームオフセットと自動リトライによるストリーム内での「exactly-once」セマンティクスを実現するために、BigQueryIOとStorage Write APIの使用を推奨します。
- 低レートの単純なパイプラインでは、ストリーミング挿入も許容されます。クライアントのリトライを重複排除するためにinsertIdを設定します。
- ストリーミングバッファに対するクエリは結果整合性を持ちます。時間に制約のある分析では、バッファの遅延(例:観測された可用性レイテンシーの約2倍待機)の後でクエリを実行するか、マイクロバッチウィンドウとStorage Write APIのコミットモードを介してマテリアライズします。
Exactly-onceの効果、べき等性、リプレイ、シンク:
- Beamはat-least-once(少なくとも1回)の処理を保証します。「exactly-once」は、べき等な書き込み、トランザクション、または重複排除キーを使用してシンクで達成する必要があります。
- BigQuery:ストリーム内でexactly-onceを実現するには、Storage Write APIのデフォルトストリームまたはコミット済みストリームを使用します。ストリーミング挿入では、安定したinsertIdを設定します。
- ファイル:一意の名前で一時ファイルを書き込み、ウィンドウの完了時にファイナライズし、アトミックな名前変更を保証します。部分的な重複を防ぐために上書きを避けます。
- 外部データベース:安定したIDをキーにしたupsertを使用するか、重複排除ウィンドウを実装します。
- リプレイを考慮した設計:決定論的な変換を維持します。シンクがリトライ時に重複排除することを保証します。
デッドレター処理、エラールーティング、オブザーバビリティ:
- リスクのあるパース/エンリッチメント処理をParDo内でtry/catchでラップし、TupleTagを介して失敗をデッドレターPCollectionに出力します。ペイロード、エラーコード、コンテキストを含めます。
- DLQ(デッドレターキュー)を分析のためにBigQueryまたはCloud Storageにルーティングします。再処理のために別のPub/Subトピックを検討します。
- オブザーバビリティ:Dataflowジョブメトリクス(ウォーターマークの遅延、システムラグ、スループット)、カスタムカウンター、分布メトリクス、およびCloud Loggingのステップごとのログを使用します。Cloud Monitoringで遅延とエラー率に関するアラートを作成します。Error Reportingを使用して例外を集約します。
パフォーマンスチューニングのパターン:
- 効率的な読み取り:BigQueryソースの場合、Storage Read API、または必要なフィールドとフィルタのみを選択するクエリベースの読み取りを推奨します。
- コンバインのリフティング:GroupByKeyの前にcombinerを使用してシャッフル量を削減します。
- サイドインプット:小さな参照データをメモリにキャッシュします。ファンアウトと更新頻度に注意します。
- シリアライゼーション:コンパクトなスキーマ(Avro/Proto)を使用し、ホットパスでの過剰なJSONパースを避けます。
デプロイ、テンプレート、アップグレード戦略
Flex Templates:
- パイプラインをコンテナ化され、パラメータ化されたテンプレートにパッケージ化し、再現可能なデプロイを実現します。Flex Templatesは、カスタム依存関係、GPUイメージ、環境の分離をサポートします。
- 実行時パラメータ(例:入力サブスクリプション、出力テーブル、デッドレターシンク、maxWorkers)を外部化し、環境固有のデプロイを可能にします。
パイプラインの更新と互換性:
- Dataflowは、変換名、状態仕様、出力タイプに互換性がある場合、多くのストリーミングパイプラインでインプレース更新をサポートします。安定したPTransform名を使用してください。
- 互換性のないグラフや状態の変更には、制御されたカットオーバーを実行します。新しいジョブを開始し、その後古いジョブをドレインして処理中の作業を完了させ、新しい要素の読み取りを停止します。
ドレインとスナップショット:
- ドレインは、処理を正常に完了させ、残りの出力を書き込み、終了します。データのギャップを避けるために、Pub/Subの保持期間やスナップショットと連携させます。
- 継続性を確保するために、Pub/Subスナップショットを作成し、新しいパイプラインをそのスナップショットまたは適切なタイムスタンプから開始し、出力を検証してから古いジョブをドレインすることができます。
設定例:
早期/遅延トリガーと蓄積を使用したウィンドウ処理の例: window .into(FixedWindows.of(Duration.standardMinutes(1))) .triggering( AfterWatermark.pastEndOfWindow() .withEarlyFirings(AfterProcessingTime .pastFirstElementInPane().plusDelayOf(Duration.standardSeconds(30))) .withLateFirings(AfterPane.elementCountAtLeast(1)) ) .withAllowedLateness(Duration.standardMinutes(10)) .accumulatingFiredPanes();
Storage Write APIを使用したBigQueryIOの例: BigQueryIO.writeTableRows() .to(“project:dataset.table”) .withMethod(BigQueryIO.Write.Method.STORAGE_WRITE_API) .withWriteDisposition(WriteDisposition.WRITE_APPEND);
よくある落とし穴:
- ストリーミングでウィンドウ処理された書き込みを行わずにファイルベースのシンクに書き込むと、ファイナライズが停止する可能性があります。ウィンドウ処理された書き込みとトリガーを有効にしてください。
- 無制限の増大:状態や許容される遅延に上限を設けるのを忘れると、メモリリークやスケーリングの失敗を引き起こす可能性があります。
- タイムスタンプの欠落:イベントタイムスタンプを割り当てないと、パイプラインはデフォルトで処理時間を使用するようになり、変動する遅延下で正確性が失われます。
実践的な問題シナリオ
NovaTrack社は、50,000個の温度センサーからグローバルなIoTテレメトリを取り込み、分単位の集計を提供し、生データを永続化し、リアルタイムダッシュボードを公開する必要があります。不正な形式のメッセージや順序が乱れた配信が時折発生することが想定されます。このソリューションは、自動スケーリングし、不良レコードを調査用に提示し、ゼロダウンタイムでのアップグレードをサポートする必要があります。
アプローチ:
取り込みと時間セマンティクス
- リージョンごとのPub/Subトピックと、属性としてdeviceIdとeventTs(RFC3339)を持つリージョンごとのパブリッシャーを作成します。可能な場合は、deviceIdによる順序付けキーを有効にします。
- 理由:Pub/Subは、at-least-once配信を備えた、耐久性があり伸縮自在なイングレスを提供します。エッジでイベントタイムスタンプを付与することで、真のイベント時刻を保持します。デバイスごとの順序付けは、中央のボトルネックなしにデバイス内の順序の乱れを低減します。
イベントタイムウィンドウを使用したDataflowストリーミングパイプライン
- 専用のサブスクリプションからPubSubIO経由で読み取り、eventTsをBeamタイムスタンプとして抽出します。欠落している場合はpublishTimeにフォールバックします。
- 1分間のFixedWindowsを適用し、30秒での早期トリガーと、遅延要素ごとに遅延発火を設定します。許容される遅延を10分に設定し、ペインを蓄積します。
- 理由:イベントタイムウィンドウは正確な分単位の集計を保証します。早期発火はサブ分単位の鮮度でダッシュボードにデータを提供し、遅延発火は遅延データが到着するにつれて集計を修正します。遅延の上限は、状態のサイズとコストを抑制します。
検証、エンリッチメント、デッドレタールーティング
- JSONをパースし、スキーマと範囲を検証し、ジョブ開始時にBigQueryからロードされたサイドインプットを介して小さな静的参照データでエンリッチするParDoを実装します。
- TupleTagsを使用して、有効なレコードをメイン出力に、失敗したレコードをペイロード、エラー、deviceId、パースタイムスタンプを含むデッドレターPCollectionに出力します。DLQはパーティション化されたBigQueryテーブルに書き込みます。
- 理由:サイドインプットは参照データをメモリ内に保持し、低レイテンシを実現します。デッドレターのキャプチャにより、メインフローをブロックすることなく、不良行の調査と的を絞った再処理が可能になります。
集計とホットキーの緩和
- deviceIdでキーを設定し、CombineFnsを使用して分単位のavg/min/maxを計算します。リージョンごとのトップNメトリクスについては、ホットキーを避けるためにregion#Nでシャーディングし、その後再集計します。
- 理由:Combinerはシャッフル量とコストを最小限に抑えます。キーシャーディングは、リージョンごとのファンイン中に単一キーのボトルネックを防ぎます。
シンクとexactly-once効果
- 検証済みの生イベントと分単位の集計を、Storage Write APIを使用したBigQueryIOでBigQueryに書き込みます。カスタムリトライでのべき等性を確保するために、deviceId + eventTsに基づいて安定した挿入IDを設定します。
- 理由:Storage Write APIは、ストリーム内でexactly-onceセマンティクスを備えた、高スループット、低レイテンシの取り込みを提供します。安定したIDは、リプレイが発生した場合にダウンストリームでの重複排除を保証します。
ダッシュボードの一貫性戦略
- ダッシュボードは、ウォーターマークに対して2分間のルックバック、またはストリーミングデータに対して観測された可用性レイテンシの2倍の固定遅延で、パーティション化された集計テーブルをクエリします。
- 理由:BigQueryのストリーミング可視性は結果整合性です。読み取りをわずかに遅らせることで、ほぼリアルタイムの動作を維持しつつ、処理中の行が欠落するのを防ぎます。
運用:オートスケーリングとストリーミングエンジン
- Streaming Engineを有効にします。予想されるピーク(例:平均の3倍)に基づいてmaxWorkersを設定し、CPUバウンドのパースと暗号化に合わせてマシンタイプを選択し、一時的なシャッフルに対応するためにブートディスクを増やします。
- ウォーターマークの遅延、バックログ秒数、CPU、ステップごとのスループットを監視します。持続的な遅延やDLQレートの急上昇に対してアラートを設定します。
- 理由:Streaming Engineは、状態/シャッフルを外部化することで、弾力性とより簡単なアップグレードを実現します。適切なサイジングと監視により、サイレントなSLO違反を防ぎます。
Flex Templatesによるデプロイとアップグレード
- パイプラインを、入力サブスクリプション、出力テーブル、DLQテーブル、maxWorkers、リージョンなどのパラメータを持つFlex Templateとしてパッケージ化します。互換性のない変更の場合は、新しいサブスクリプションで同じトピックをターゲットとする新しいパイプラインを開始し、出力を検証してから古いジョブをドレインします。オプションで、Pub/Subスナップショットを作成し、新しいサブスクリプションをそのスナップショットにシークさせることで、ギャップがないことを保証します。
- 理由:Flex Templatesは、再現可能でパラメータ化されたデプロイを可能にします。ドレインを伴う検証済みのブルー/グリーン方式のカットオーバーは、ゼロデータ損失と最小限のダウンタイムを実現します。
再処理とバッチバックフィル
- サイドアウトプットを介して、圧縮された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.
試験に合格する →