TT Lab
はじめる
学ぶ 学習パス コース

データパイプライン

バッチとストリーミング — 何を基準に選ぶのか

TT Labで続きを見る

一言でいうと

バッチは境界がはっきりしたデータのまとまりを定期的に処理し、ストリーミングは終わりなく届くイベントを継続的に処理します。2つの本当の違いは速度ではなく、境界を誰が決めるかです。

なぜ必要なのか

「リアルタイムがいいのでストリーミングにしよう」という決定はよく下され、その代償は、たいてい6か月後に請求されます。バッチは失敗したらもう一度実行すればよいですが、ストリーミングは状態を持って動き続けるシステムなので、再処理の設計がはるかに難しいです。

どう動くのか

バッチの中核の仕組みはウォーターマークです。どこまで処理したかを保存しておき、次の実行ではそれ以降だけを取得します。

SELECT * FROM orders
WHERE ordered_at > (SELECT last_ordered_at FROM etl_watermark WHERE job_name = 'orders_archive')
  AND ordered_at < :batch_end;

この単純なパターンには、落とし穴がいくつも隠れています。

ストリーミングでは、この問題がより露骨です。イベントの時刻と処理の時刻が異なるため、「今このウィンドウを閉じてよいか」を判断する必要があり、そのため、ウォーターマークが許容遅延とともに定義されます。そして、遅れて届いたイベントを捨てるのか、ウィンドウを開き直すのかを、ポリシーとして決める必要があります。

配信保証も、3種類に分かれます。最大1回(失われることがある)、最小1回(重複することがある)、正確に1回です。実際のシステムの大半は、最小1回を提供し、重複をコンシューマー側で吸収して、結果として1回のようにする方式を選びます。そのため、次のモジュールの冪等性が重要になります。

何を基準に選ぶのか

「リアルタイムがいい」は基準ではありません。4つを検討すれば、たいてい答えが決まります。

問い バッチが合う ストリーミングが合う
結果がいつ必要か 時間単位・日単位 秒単位・分単位
遅れて届いたデータをどうするか 次のバッチに自然に含まれる ウォーターマーク・再処理の設計が必要
全体を再計算できるか 簡単 難しい(状態の復元)
運用の負担 失敗したらもう一度実行すればよい 常に稼働している必要がある

元に戻せるかが最も重要です。バッチは、ロジックを直して昨日の分をもう一度実行すれば終わりです。ストリーミングで同じことをするには、オフセットを巻き戻し、状態を初期化し、下流の重複に対処しなければなりません。

そのため、実務でよくある答えは両方です。ストリーミングで素早い近似値を出し、バッチで正確な値を、あとから上書きします(ラムダアーキテクチャ)。最近は、同じコードで両方を処理する方式(カッパ、Flink・Beam)が増えましたが、運用の複雑さは、依然としてストリーミングのほうが大きいです。

時間に関する3つのこと

ストリーミングでずれるものの半分は、時間の定義です。

モバイルアプリは、機内モードだったあと、数時間後にイベントを送ります。イベント時刻で集計すると、すでに締め切ったウィンドウ(window)に、遅れたデータが入ってきます。ウォーターマークは、「この時刻より前のデータは、もう来ないとみなす」という宣言で、その線を越えて到着したものは、捨てるか、別の経路で処理します。

워터마크 = 지금까지 본 최대 이벤트 시각 − 허용 지연(예: 10분)

許容遅延を増やせば正確になりますが、結果がその分遅く出ます。正確さとレイテンシを引き換えるノブは、これ1つだけです。

バッチも増分で回す

バッチだからといって、毎回全体を読む必要はありません。最後に処理した地点を記録して、その後だけを読みます。ただし、境界を重ねてとります。

-- 워터마크를 그대로 쓰면 경계에 걸친 것을 놓친다
where updated_at >= :last_watermark - interval '10 minutes'
  and updated_at <  :now

重なる区間はもう一度読みますが、ロードが冪等なら、何の害もありません。冪等性があれば、重ねて読むのはただ同然という原理が、ここでも同じです。

現場での姿

バッチがストリーミングより優れている状況は、思ったより多いです。ソースが1日1回更新されるのに、パイプラインだけがリアルタイムでは、何の価値もありません。レポートを見る人が、朝に1回見るなら、早朝のバッチで十分です。コンシューマーが実際にどれくらいの頻度で見るかが、最初の問いになる必要があります。

反対に、バッチの代償も明らかです。周期が長いほど、一度失敗したときの遅延が大きくなり、一度に処理する量が大きいため、リソースの使用が尖ります。そのため、大きなバッチを細かく分けて、チャンク単位で処理する方式が、よく使われます。失敗したチャンクだけをもう一度実行でき、ロックの時間も短くなります。

次のラボですること

注文データをCSVに抽出して、アーカイブ表にロードし、ウォーターマークを記録して増分ロードを実行したあと、同じ作業を2回実行しても結果が変わらないかを確認します。