Apache Flink — ストリームを本物のエンジンで動かす
ウォーターマーク — エンジンが「このウィンドウはもう閉じる」と決める方法
一言でいうと
イベント時間ウィンドウは、ウォーターマークがウィンドウの終わりを過ぎてはじめて閉じます。Flink SQLのウォーターマークは、WATERMARK FOR ts AS ts - INTERVAL '5' SECONDのように宣言し、値は「これまでに見た最大のts − 遅延」です。ウィンドウが閉じたあとに届いた行は、黙って捨てられます。遅延は、精度とレイテンシのトレードオフを調整するつまみです。
なぜ必要なのか
前のモジュールのGROUP BY user_idは、結果を書き換え続ければ済みました。ところが、「1分ごとにセンサー値を合計して、1回だけ出力せよ」という要求は違います。1回だけ出力するには、その1分が終わったことを知らなければなりません。壁時計で判断すれば(処理時間)簡単ですが、公式ドキュメント(Timely Stream Processing)が指摘するとおり、処理時間は到着速度・障害・再処理に左右されて、決定的ではありません。昨日のデータをもう一度動かすと、結果が変わります。
そこで、レコードに書かれたイベント時間でウィンドウを分けます。問題は、イベントが順序どおりに来ないことです。12:00:59のセンサー値が、12:01:10の値より遅れて届くことがあります。12:00のウィンドウをいつ閉じればよいのでしょうか。永遠に待つことはできません。ウォーターマークは、この問いに対する約束です。Watermark(t)は、「これからはts ≤ tのイベントはもう来ないものとみなす」という宣言です。
どう動くのか
宣言は、次のとおりです。ドキュメント(CREATE Statements)によると、WATERMARK句は、TIMESTAMP(3)の列1つをイベント時間属性にします。DESCRIBEで見ると、その列の型に*ROWTIME*が付きます。よく使う戦略は3つです。
| 式 | 意味 |
|---|---|
ts |
厳密な昇順。見た最大のtsが、そのままウォーターマーク |
ts - INTERVAL '0.001' SECOND |
昇順。最大のtsと同じ時刻の行は、遅延行ではない |
ts - INTERVAL '5' SECOND |
順序が乱れた入力。5秒まで待つ |
ウォーターマークが出るタイミングは、次のとおりです。式はレコードごとに評価されますが、ドキュメントは、ウォーターマークがpipeline.auto-watermark-interval(デフォルト200ms)の周期で出ると記しています。ところが、ウォーターマークを付けるオペレーターのソースを見ると、もう1つあります。新しいウォーターマークが、最後に出力した値より間隔を超えて先に進んでいれば、周期を待たずにそのレコードですぐ出力します。そして、順序が重要です。レコードを先に下流へ流したあとで、ウォーターマークを上げます。そのため、ある行が見るウォーターマークは、「その前の行までの最大ts − 遅延」です。秒単位のタイムスタンプなら、前進幅が常に200msより大きいので、レコードごとにウォーターマークが出て、結果が実行速度と関係なく決まります。ラボの採点ツールが、到着順だけで結果を再現できる理由です。
ウィンドウが閉じる条件は、次のとおりです。ウィンドウ集計は、window_end ≤ 워터마크(プレースホルダーはウォーターマークです)になった瞬間に、そのウィンドウの結果を1回出力して、状態を空にします。そのあとで、そのウィンドウに入る行が来ると、受け取る場所がないので捨てます。エラーも警告もありません。終わりのあるファイルを読み終えると、エンジンが最大のウォーターマークを送って、残りのウィンドウをすべて閉じます。
行が遅れていることと、ウィンドウが閉じていることは別です。CURRENT_WATERMARK(ts)は、その行が通過するオペレーターの現在のウォーターマークを返します(ドキュメントのBuilt-in Functions、まだなければNULL)。ts <= CURRENT_WATERMARK(ts)の行は、ウォーターマーク基準で遅れた行です。しかし、その行のウィンドウが、まだ開いていることがあります。12:00:58の行がウォーターマーク12:00:59のあとに来ても、12:00ウィンドウの終わり(12:01:00)は、まだウォーターマークの右側にあります。そのため、ドキュメントのフィルター式(CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts))でウィンドウの前でフィルタリングすると、ウィンドウに任せるときより多くの行を捨てます。ラボのデータで測ると、5秒の遅延でウィンドウが捨てる行は15件ですが、ウォーターマーク基準で遅れた行は98件です。
並列の場合は、複数の入力を受けるオペレーターのイベント時間が、入力のウォーターマークのうち最小値になります(ドキュメントのWatermarks in Parallel Streams)。1つのパーティションだけが静かでも、全体のウォーターマークが止まります。これを解決するために、table.exec.source.idle-timeoutで静かなソースを一時的に外す設定があります。このラボは並列度1なので、この効果は見られません。
現場での姿
「集計の数字が元データより少し少ない」という報告の多くは、遅延行です。捨てられた行はログにも残らないので、同じデータをバッチで1回動かして差を測るのが、最初の診断です。差が、遅延を増やすと減るかを見れば、原因を絞り込めます。遅延を増やすと、そのぶんウィンドウの結果が遅く出ます。ダッシュボードが5秒遅れてもよいのか、30秒まで大丈夫なのかは、ビジネスが決める問題です。
2つ目は、「結果がまったく出ない」場合です。ウォーターマークが進まなければ、ウィンドウは永遠に閉じません。パーティションの1つが空であるか、ソースのtsがNULLであるか、テストデータの時刻が1点に集中している場合です。ウォーターマークの間隔をとても大きくしても、似たことが起きます。ラボで間隔を1時間にすると、最初のウォーターマークのあとは進まなくなり、終わりのあるファイルなので最後にまとめて閉じて、捨てられた行が1つもない結果になります。無限ストリームなら、ウィンドウが1時間閉じなかったでしょう。
3つ目は、遅延行を捨てずに別に集めたいという要求です。SQLでは、CURRENT_WATERMARKで遅延行にマークを付けて、別のシンクに送るのが一般的な方法です。ただし、上で見たとおり、「ウォーターマークより早い行」と「ウィンドウがすでに閉じて捨てられる行」は、別の集合です。どちらを集めるのかを、先に決める必要があります。
次のラボですること
センサーイベント600件(ファイルの順序が到着順)を、5秒遅延のウォーターマークで宣言し、DESCRIBEで確認します。1分のTUMBLE集計をバッチで動かしてベースラインを作り、遅延5秒・0秒・30秒のストリーミング結果と比べて、捨てられた行を数えます。CURRENT_WATERMARKで遅延行を取り出し、その行をウィンドウの前でフィルタリングした結果がどう違うかを確認し、ウォーターマークの間隔を1時間に変えて、前進が止まると何が変わるかを確認したあと、数字を報告書にまとめます。