バッチとストリーミング — 何を基準に選ぶのか
一言でいうと
バッチは境界がはっきりしたデータのまとまりを定期的に処理し、ストリーミングは終わりなく届くイベントを継続的に処理します。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;
この単純なパターンには、落とし穴がいくつも隠れています。
- 境界値を含めるかどうか。
>と>=を取り違えると、実行のたびに1件が重複したり、欠落したりします。 - 同一のタイムスタンプ。同じ時刻に複数の行が入ってくると、
>だけでは一部が永遠に欠落します。時刻と主キーを合わせて比較するか、ウォーターマークを少し後ろに戻しておき、重複をupsertで吸収するほうが安全です。 - 遅れて届くデータ。ソースがトランザクションを遅くコミットすると、ウォーターマークより過去の時刻の行が、あとから現れます。この場合、時間基準のウォーターマークは、その行を永遠に見逃します。
ストリーミングでは、この問題がより露骨です。イベントの時刻と処理の時刻が異なるため、「今このウィンドウを閉じてよいか」を判断する必要があり、そのため、ウォーターマークが許容遅延とともに定義されます。そして、遅れて届いたイベントを捨てるのか、ウィンドウを開き直すのかを、ポリシーとして決める必要があります。
配信保証も、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回実行しても結果が変わらないかを確認します。