ウォーターマークは事実ではなく約束だ
一言でいうと
ウォーターマークは、「この時刻までのデータはすべて届いた」という観測ではなく約束であり、その約束を破るデータに出会ったときに何をするかを前もって決めておかなければ、昨日出した数字を説明できなくなります。
なぜ必要なのか
このコースの前のラボで、ウォーターマークをすでに一度扱いました。ロードしたデータの最大時刻を表に書いておき、次の実行がその後だけを読むようにするものでした。そのウォーターマークが答える問いは、1つです。どこまで読んだかです。
しかし、その問いだけでは解けないことがあります。モバイルアプリは、地下鉄で電波を失って数分後に、機内モードだった端末は数時間後に、イベントを送ります。そのイベントの発生した時刻は、すでに過ぎ去った時刻です。時間基準のウォーターマークは、そのデータを永遠に見られません。そして、私たちはすでにその時間帯の集計を送り出しています。
そのため、朝の会議で、こんな話が出ます。「昨日ダッシュボードで見た数字と、今日の数字が違います」。この問いに答えられなければ、そのダッシュボードは、その日から誰にも信じられなくなります。
どう動くのか
まず、時間が3種類あることから分ける必要があります。イベント時間は実際に起きたとき、取り込み時間は私たちのシステムに入ったとき、処理時間は私たちが計算したときです。分析の基準は常にイベント時間です。処理時間を基準にすると、再処理のたびに答えが変わります。
イベント時間でウィンドウ(window)を分けると、すぐに1つ問題が生じます。このウィンドウをいつ閉じるかです。永遠に開けておくと結果が出ず、早く閉じすぎると、まだ届いていないデータを取りこぼします。
ウォーターマークが、その判断を代わりに行います。Flinkのドキュメントは、ウォーターマークを「イベント時間での進行状況をシステムに知らせるもの」と説明しています。よくある計算式は次のとおりです。
워터마크 = 지금까지 본 최대 이벤트 시간 − 허용 지연
창이 닫힌다 = 워터마크가 그 창의 끝을 지났다
ここで、2つを見落としやすいです。
1つ目は、ウォーターマークは後ろに戻らないことです。最大イベント時間は単調増加なので、ウォーターマークも単調増加です。遅れて届いたイベントが、ウォーターマークを戻すことはありません。戻せるようにすると、すでに閉じたウィンドウが、再び開いて閉じることを繰り返し、どんな値も確定しません。
2つ目は、遅れはデータの性質ではなく、到着順の性質だということです。同じイベントでも、いつ到着したかによって、遅いものになることも、ならないこともあります。そのため、判定は、そのイベントより先に到着したもので計算したウォーターマークで行います。ここで、よく起きる事故があります。データを扱いやすくするために、イベント時間で一度ソートしてしまうことです。その瞬間に、到着順が失われ、遅延データが0件で出ます。実測してみると、前後が入れ替わった到着が20件以上あったデータでも、ちょうど0が出ます。問題がないからではなく、見られる目をなくしたのです。
許容遅延は、正確さとレイテンシを引き換えるノブです。大きくとれば、遅延データをより多く受け入れますが、ウィンドウがその分遅く閉じ、小さくとれば、早く出せますが、より多く取りこぼします。この値は、勘で決めるのではなく、実際の遅延の分布から読み取ります。遅延の中央値と95パーセンタイルと最大値を測ってみると、たいてい95パーセンタイルの近くで折れ曲がり、その上は、尾が長く伸びています。最大値に合わせると、その1件のために、全員が待つことになります。
遅延データをどうするか
ウォーターマークを越えて届いたデータに対する選択肢は、2つです。
捨てます。確定した数字は、絶対に変わりません。その代わり、その分が静かに消えるので、捨てた件数と金額を必ず別に数えます。数えずに捨てると、あとで「なぜ私たちの合計が合わないのですか」を説明する根拠がありません。
遅延更新として反映します。同じFlinkのドキュメントのウィンドウの説明は、許容遅延を置いたあとで、遅れて到着した要素が、ウィンドウをもう一度発火させることがあり、そのとき出る値は、前の計算の更新された結果として扱わなければならないと書いています。下流がそれを更新として受け入れられず、追記だけをすると、重複が生じます。そのため、この選択は、こちら側だけの決定ではありません。
どちらにしても、閉じた瞬間の値と、そのあとの訂正を別に残すことが核心です。1つにまとめておくと、「昨日の数字が今日変わった」に答えられませんが、別に残せば、変わった量がそのまま答えになります。
現場での姿
1つ目は、ウィンドウが閉じないことです。データが途絶えると、最大イベント時間が止まり、ウォーターマークも止まり、ウィンドウは永遠に開いたままになります。同じFlinkのドキュメントが扱うアイドル(idleness)の問題が、これです。静かなパーティション1つが、全体を引き止めます。
2つ目は、許容遅延を増やして忘れることです。事故が起きると、「とりあえず余裕を持たせて」と増やされ、そのまま固まります。6時間後に出るダッシュボードは、リアルタイムではありません。増やした値は、戻す日付と一緒に書いておきます。
3つ目は、再処理と遅延データが混ざることです。過去の区間をもう一度回すことと、遅れて届いたデータを反映することは、どちらも「古い数字が変わる」ように見えます。原因が違うため、記録も別に残す必要があります。
4つ目は、遅延データが一方に偏ることです。地域や端末の機種で分けて見ると、特定のグループだけが大きく遅れています。全体の平均だけで見ると、そのグループのデータは常に捨てられ、その事実を誰も知りません。
実務で本当に大切なこと
- 遅延を分布で測ります。平均1つで許容遅延を決めません。
- ウォーターマークは単調増加です。遅延データがウォーターマークを戻すようにしません。
- 到着順を失いません。イベント時間でもう一度ソートすると、遅延データが0件で出ます。
- 捨てたら数えます。捨てた件数と金額がなければ、説明できません。
- 閉じた瞬間の値と訂正を別に残します。その差が、そのまま答えです。
次のラボですること
地下鉄と機内モードを通ってきた注文イベントを作り、ツールwm.pyを1ステップずつ育てます。遅延を分布で測り、到着順にウォーターマークを動かし、イベント時間でウィンドウを分け、ウィンドウが閉じたあとに届いたデータを分けます。そのあと、同じデータを、捨てるポリシーと反映するポリシーでそれぞれ集計して、どのウィンドウがどれだけ変わるかを見て、閉じた瞬間の値と訂正を別に出します。最後に、その数字でレポートを書きます。採点ツールは、毎回異なるウィンドウサイズと許容遅延で、自分で作ったツールを実際に動かして、答えを突き合わせます。