代表的な面接トピック

データエンジニアリング面接:ストリーミングパイプラインをSLOごとに分割すべきタイミングとは?

データ難しい
Offer.cc 編集チーム公開日 更新日

質問

1つのイベントストリームが、60秒以内に更新する必要がある運用ダッシュボードと、翌朝7:00までに完了する必要がある財務レポートにデータを供給しています。1つの共有パイプライン、1つのパイプライン内の分岐、別々のパイプラインのどれを採用しますか?SLO、リソース、公平性(fairness)、リプレイ、障害の観点から説明してください。

プロンプトと適用範囲

これは、データプラットフォーム、リアルタイム分析、シニアデータエンジニア職でよく出題される設計問題です。ピーク時に毎秒50,000イベント、ダッシュボードの鮮度目標が60秒、前日分のレポート締め切りが07:00であると仮定します。面接では、共有実行グラフがすべてのコンシューマーに最も厳しいSLOを強制してしまうタイミングや、独立した処理進行とリソースプールによってどのように結合度を下げられるかを評価します。

面接官が評価しているポイント

  • ツール名を挙げる前に、「リアルタイム」や「定刻通り」をエンドツーエンドのSLOに落とし込めているか。
  • 単一のBeamグラフ内の分岐、同一ソースへの独立したサブスクリプション、完全に独立したコンピュートリソースの違いを区別できているか。
  • 重複読み取り、コスト、バックプレッシャー、リプレイ、遅延イベント、障害ドメインのトレードオフを説明できるか。
  • パイプラインの分割によってユーザー体験が向上することを証明するメトリクス、カナリアリリース、ロールバック、リコンシリエーション(整合性検証)を提示できるか。

最初に確認すべき明確化のための質問

60秒というのがイベント発生からクエリ可能になるまでの鮮度(event-time-to-queryable freshness)を指すのか、それともプロセッサのレイテンシのみを指すのかを確認します。また、財務レポートで遅延パーティションのバックフィルが許容されるか、両方のコンシューマーがイミュータブルな未加工ストレージを共有できるか、イベントにテナントや優先度の公平性が必要か、重複・欠損・順序付けに関してどのような保証が求められるか、予算はコストとダッシュボードのテールレイテンシのどちらを優先するかを尋ねます。レポートがT+1(翌日)でありながらダッシュボードに厳格な低レイテンシ目標がある場合、分離はすでに評価に値する有力なデフォルト案となります。

30秒での回答構成

まず、2つのコンシューマーのエンドツーエンドSLO、正確性の保証、復旧目標を書き出します。1つの共有パイプラインで60秒という厳格な目標を満たす必要がある場合は、共有のベースラインを確立した上で、独立したサブスクリプションを使用してリアルタイム処理とバッチ処理の進行を分離します。負荷の高い共通のパース処理は共有しますが、リソースの逼迫や障害の波及が重大な問題となる場合はコンピュートを分割します。各パスに鮮度、バックログ、遅延、重複、レポート完了のメトリクスを個別に持たせ、リプレイやフォールトインジェクション(障害注入)によって追加の分離がコストに見合うかを判断します。

詳細な回答

1. SLOから構造を導き出す

リアルタイムのSLOを、p99のイベント発生からクエリ可能までの鮮度60秒未満と定義します。バッチのSLOは、前日の有効なイベントの処理を07:00までに完了し、遅延データは制限された修復ウィンドウ内で処理することと定義します。GoogleのDataflowのガイダンスでは、複数の異なるSLOに対応する単一のパイプラインは、より厳しい目標を満たす必要があり、その結果、優先度の低い処理がリアルタイム処理のキャパシティを消費してしまう可能性があると指摘されています。エラーバジェット、アラートルート、オートスケーリングポリシーが異なることは、処理パスを分離する強力な理由になります。

2. 3つのトポロジを比較する

グラフ内での分岐は、共通のデコードや軽量なルーティングに役立ちます。Apache Beamのドキュメントによると、複数の変換(transform)が同じPCollectionを読み取ることができますが、各変換が再度入力を処理します。単一のマルチ出力変換を使用すれば、共通の処理について各要素を1回だけ処理できます。

同一トピックへの独立したサブスクリプションを使用すると、リアルタイムコンシューマーとバッチコンシューマーが個別の確認応答(ack)、バックログ、リプレイ位置を保持できます。GoogleのDataflowのガイダンスでは、各ジョブが独立してプルおよび確認応答を行えるように、別々のサブスクリプションを使用する複数のパイプラインが説明されています。完全に独立したジョブにすると、重複読み取り、シリアライズ、運用の所有権といったコストが発生するものの、CPU、メモリ、リリース頻度、障害ドメインを完全に追加で分離できます。

3. 推奨パスの設計

イミュータブルな未加工イベント層を1つ維持し、同じソースからリアルタイム用サブスクリプションとバッチ用サブスクリプションを作成します。リアルタイムジョブは、低レイテンシの配信ストアに向けて軽量な集計を実行します。バッチジョブは、イベント日付ごとに保持データを読み取り、パーティションテーブルに書き込みます。共通のパース処理が合計CPUの20%以上を占める場合は、取り込み(ingress)時に1度正規化して未加工層にバージョニングされたイベントを書き込みます。リアルタイムコンシューマーの進行を、バッチパスがコミットされたことの証明として扱ってはなりません。

text
raw-events
  -> realtime-subscription -> stream-aggregate -> serving-store
  -> batch-subscription or retained-raw -> daily-transform -> partitioned-lake

4. バックプレッシャー、優先度、コストの管理

リアルタイムパスには独自の並行性上限とバックログアラートを設定します。リアルタイムのリソースが逼迫したときはバッチ側の並行性を下げられるようにしますが、両方を1つの無制限キューの後ろに配置してはいけません。2つの完全なコンピュートを実行することがコスト的に高すぎる場合は、未加工データのデコードと配置を共有し、ダウンストリームのステージを分離します。リアルタイムのバックログが60秒を超えた場合は、バッチキャパシティの拡張を停止し、リアルタイムSLOの復旧を最優先します。ジョブ数を比較するのではなく、入力バイト数、CPU、バックログの滞留時間、100万イベントあたりのコストを各パスに帰属させて評価します。

5. 遅延データ、リプレイ、復旧の処理

リアルタイムウィンドウにはウォーターマークと制限された許容遅延(allowed lateness)を使用します。ウィンドウを超えたイベントは遅延データキューまたは未加工層にルーティングし、バッチ側でバージョニングされたべき等な書き込みによって影響を受けたパーティションを修復させます。リプレイは、保存されたソース位置から新しいサブスクリプションを作成して実行します。本番コンシューマーの確認応答位置を決して巻き戻してはいけません。バッチパスは、リアルタイム側で障害が発生した後でも未加工データから復旧できるようにすべきです。未加工ストレージが利用できない場合、両方のパスで明示的な縮退運転とアラートが必要になります。

6. 受け入れ実験による検証と判断

まず共有ベースラインを実行し、次に小規模なテナントスライスに対して分離されたリソースを有効にします。ダッシュボードのp50/p95/p99鮮度、レポート完了、バックログ滞留時間、重複率、リプレイ時間、CPU、ストレージ、100万イベントあたりのコストを比較します。バッチのバースト、プロセッサの再起動、重複メッセージ、遅延パーティション、サブスクリプションの一時停止などを注入(テスト)します。もし分割によってプロセッサのレイテンシのみが改善し、データの重複が増加したりコスト予算を超過したりする場合は、共有構成のまま共通ステージを最適化します。

優れた回答例

私ならまず2つのエンドツーエンドSLOを定義します。ダッシュボードのp99イベント鮮度(クエリ可能まで)60秒未満と、制限された遅延データ修復ウィンドウを持つ07:00までの前日財務レポート完了です。共有ベースラインを構築しますが、リアルタイムとバッチの間で確認応答の位置は決して共有しません。私のデフォルト設計は、1つのイミュータブルな未加工イベント層、2つの独立したサブスクリプション、そして2つのダウンストリームジョブです。リアルタイムジョブは軽量な集計を実行し、バッチジョブはイベント日付ごとに保持データを読み取ります。共通のパース処理は取り込み時に一度バージョニングできますが、完全なジョブの分離はリソースや障害ドメインで必要とされる場合にのみ行います。各パスが鮮度、バックログ、遅延、重複、レポート完了、ユニットコストのメトリクスを保持し、リプレイには新しいサブスクリプションとべき等性キーを使用します。バースト、再起動、遅延イベントを注入してSLOを検証します。分離によってユーザー成果が向上しない場合は、コンピュートを共有したまま共通ステージを最適化します。

よくある間違い

  • 兆候: コンシューマーが2つあるという理由だけで、2つの完全なパイプラインを複製する。失敗する理由: 分離性が向上しないまま、共通のパースと配置のコストが2倍になる。対策: まずイミュータブルな未加工データを共有し、その後SLOに応じてダウンストリームのコンピュートを分離する。
  • 兆候: 1つのグローバルオフセットで両方のコンシューマーを駆動する。失敗する理由: 遅いコンシューマーが速いコンシューマーをブロックし、リプレイを独立して行えなくなる。対策: 独立したサブスクリプション、または独立して検証可能な進行管理を使用する。
  • 兆候: プロセッサレイテンシのみを測定する。失敗する理由: ストレージ、クエリ、更新時間によってユーザー目標に違反する可能性がある。対策: エンドツーエンドの鮮度とパーセンタイルを測定する。
  • 兆候: 遅延イベントをライブ結果に直接書き込む。失敗する理由: カウントが重複したり、公開済みのレポートが意図せず変更されたりする。対策: ウォーターマーク、修復ウィンドウ、バージョン、べき等性キーを使用する。
  • 兆候: 履歴をリプレイするために本番コンシューマーを巻き戻す。失敗する理由: オンラインの処理進行が妨げられ、トラフィックが増幅する可能性がある。対策: 保持データからレート制限付きのリプレイスクリプションを作成する。

フォローアップの質問

両方のパスで同じ高負荷な特徴量計算が必要な場合はどうしますか?

特徴量計算ステージをバージョニングしてリプレイ可能にし、その出力を一度配置して、両方のパスから読み取れるようにします。状態をオンラインに維持する必要があり共有できない場合にのみ、重複計算を許容します。CPU、レイテンシ、一貫性の実験を使用して、共有中間層と重複計算を比較評価します。

独立したサブスクリプションによって入力コストは2倍になりますか?

読み取りと確認応答のオーバーヘッドは増加しますが、イミュータブルな配置、圧縮、保持ウィンドウ、オンデマンドリプレイによって抑制できます。ストレージコスト単体ではなく、リアルタイムSLOおよび障害分離と合わせて100万イベントあたりのコストを評価します。

バッチ処理が遅延した際、リアルタイム側のキャパシティを借りることはできますか?

リース機能と自動回収を備えた、制限付きのプリエンプティブルな低優先度キャパシティを使用します。リアルタイム側は厳格な上限と独立したバックログメトリクスを維持し、借りたキャパシティによってそのp99が60秒を超えないようにします。

2つのパスが最終的に一致することをどのように証明しますか?

同じイベントバージョン、ビジネスキー、時間境界からリコンシリエーションセットを構築します。件数、金額、欠損レコード、重複、遅延修正を比較します。リアルタイムの出力が近似値である場合は、単一の合計値の一致を証明とするのではなく、収束ウィンドウと説明可能な差異の許容量を定義します。

公開情報ソース

関連する質問