Apache Flink — ストリームを本物のエンジンで動かす
遅延を変えながら捨てられた行を数える
目標
イベント時間ウィンドウが、ウォーターマークによっていつ閉じ、どの行を捨てるかを、同じデータをバッチと複数の遅延のストリーミングで動かして、数字で確認します。CURRENT_WATERMARKで遅延行を自分で見つけ、ウィンドウが捨てる行と比較します。
なぜ重要なのか
ウォーターマークがウィンドウを閉じたあとに届いた行は、エラーなしに消えます。そのため、「数字が少し足りない」という問題はログで探せず、ウォーターマークがどう動くかを知っていてはじめて説明できます。このラボの入力は秒単位のタイムスタンプなので、ウォーターマークがレコードごとにすぐ進み、結果が実行速度と関係なく、到着順だけで決まります。採点ツールは、クラスターに問い合わせません。皆さんが保存したsql-clientの出力を読み、元のCSVを到着順に流して、エンジンと同じルール(行を先に出力してからウォーターマークを上げる・window_end ≤ 워터마크のウィンドウの行を捨てる。プレースホルダーはウォーターマークです)で計算した値と照合します。
ステップ
flink-upでクラスターを起動し、/root/flink/watermark/ddl.sqlに、WATERMARK FOR ts AS ts - INTERVAL '5' SECONDを置いたeventsテーブルとDESCRIBE events;を書いて、出力を、/root/flink/watermark/ddl.outに保存してください。- /root/flink/watermark/batch.sqlに、バッチモードで1分の
TUMBLEウィンドウごとにcnt(件数)・total(readingの合計)を出すクエリを書いて実行し、出力を、/root/flink/watermark/batch.outに保存してください。 - /root/flink/watermark/w5.sqlに、同じ集計をストリーミングモード(遅延5秒)で動かすクエリを書いて実行し、出力を、/root/flink/watermark/w5.outに保存してください。
- /root/flink/watermark/sweep.sqlに、ウォーターマークが
tsのテーブルevents_0と、ts - INTERVAL '30' SECONDのテーブルevents_30を作成し、同じ集計をevents_0からevents_30の順に実行して、出力を、/root/flink/watermark/sweep.outに保存してください。 - /root/flink/watermark/late.sqlに、5秒遅延の
eventsで、CURRENT_WATERMARK(ts)がNULLではなくts <= CURRENT_WATERMARK(ts)である行のevent_id, ts, wmを取り出すクエリを書いて実行し、出力を、/root/flink/watermark/late.outに保存してください。 - /root/flink/watermark/filtered.sqlに、遅延行をウィンドウの前でフィルタリング(
CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts))したあと、同じ1分の集計を出すクエリを書いて実行し、出力を、/root/flink/watermark/filtered.outに保存してください。 - /root/flink/watermark/slow.sqlに、ステップ3と同じ集計に
SET 'pipeline.auto-watermark-interval' = '1 h';だけを加えたクエリを書いて実行し、出力を、/root/flink/watermark/slow.outに保存してください。 - /root/flink/watermark/report.jsonに、
total_rows・dropped_0・dropped_5・dropped_30・late_rows_5・dropped_slowを書いてください。
参考
- 元の列:
event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3)(ヘッダーなしのCSV、ファイルは/opt/lab/fixtures/data/watermark_events.csv、ファイルの順序 = 到着順)。 - ウィンドウ集計の形:
SELECT window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total FROM TUMBLE(TABLE 표, DESCRIPTOR(ts), INTERVAL '1' MINUTE) GROUP BY window_start, window_end;(プレースホルダーはテーブル名です) - 1つのSQLファイルにSELECTが複数あれば、順にジョブが動き、出力に結果のテーブルが順番に出力されます。
- よくある間違い: バッチモードはウォーターマークを使いません。遅延行を見るには、ストリーミングで動かす必要があります。並列度はデフォルトの1のままにしてください。入力が複数あると、ウォーターマークはそのうちの最小値に従います。
- よくある間違い: 最初の行では、ウォーターマークがまだなく、
CURRENT_WATERMARK(ts)がNULLです。NULLとの比較は真ではなく、WHEREを通過できないので、ステップ6のフィルター式からIS NULLの条件を外すと、最初の行まで捨てられます。 - 公式ドキュメント: Timely Stream Processing・CREATE — WATERMARK・Time Attributes・Windowing TVF・Built-in Functions・Configuration
tsをイベント時間として宣言する
flink-upでクラスターを起動し、/root/flink/watermark/ddl.sqlに、元のCSVを読むeventsテーブル(列はevent_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3)、WATERMARK FOR ts AS ts - INTERVAL '5' SECOND)とDESCRIBE events;を書いて、出力を、/root/flink/watermark/ddl.outに保存してください。
WATERMARK句は、列の一覧の中、最後の列の後ろに置きます。DESCRIBEの結果で、tsの型の横にROWTIMEが付き、watermarkの列に式が見えれば、イベント時間属性になっています。filesystemコネクターのpathは、file:///opt/lab/fixtures/data/watermark_events.csvです。
バッチでベースラインを作る
/root/flink/watermark/batch.sqlに、SET 'execution.runtime-mode' = 'batch';・ステップ1のeventsテーブル・1分のTUMBLEウィンドウごとにwindow_start, window_end, COUNT(*) AS cnt, SUM(reading) AS totalを出す集計を書き、出力を、/root/flink/watermark/batch.outに保存してください。
バッチは、入力をすべて集めてから計算するので、ウォーターマークで行を捨てません。そのため、この結果が「遅延行が1つもなかったなら」のベースラインになります。バッチでも、TUMBLEはTIMESTAMP列に使えます。cntをすべて足すと、元の行数と同じになるはずです。
遅延5秒のストリーミング: 捨てられる行
/root/flink/watermark/w5.sqlに、ステップ2と同じ集計をSET 'execution.runtime-mode' = 'streaming';で動かすクエリを書き(ウォーターマークは5秒遅延のまま)、出力を、/root/flink/watermark/w5.outに保存してください。
ウィンドウは、window_endがウォーターマーク以下になった瞬間に結果を1回出力して、状態を空にします。そのあとでそのウィンドウに入る行が来ると、捨てます。結果テーブルのcntの合計を、バッチと比べてみてください。ウィンドウの結果は+Iだけで出ます。
遅延0秒と30秒を一度に比べる
/root/flink/watermark/sweep.sqlに、ウォーターマークがtsのテーブルevents_0と、ts - INTERVAL '30' SECONDのテーブルevents_30(列・ソースはeventsと同じ)を作成し、ストリーミングモードで同じ1分の集計を、events_0を先に、events_30をあとに実行して、出力を、/root/flink/watermark/sweep.outに保存してください。
ウォーターマークの式はテーブルごとに付くので、遅延を変えるにはテーブルを別に作ります。1つのファイルのSELECT 2つは、ジョブ2つとして順に動き、出力に結果のテーブルがその順序で出力されます。遅延が0なら、少し遅れただけでも捨てられ、30秒ならほとんど待ってくれます。
CURRENT_WATERMARKで遅延行を取り出す
/root/flink/watermark/late.sqlに、ストリーミングモードで5秒遅延のeventsに対して、SELECT event_id, ts, CURRENT_WATERMARK(ts) AS wm ... WHERE CURRENT_WATERMARK(ts) IS NOT NULL AND ts <= CURRENT_WATERMARK(ts)を実行するクエリを作成し、出力を、/root/flink/watermark/late.outに保存してください。
CURRENT_WATERMARKは、その行が通過するオペレーターの現在のウォーターマークです。ウォーターマーク生成器は、行を先に出力してからウォーターマークを上げるので、ある行が見る値は、その前の行までの最大ts − 5秒です。遅延行の数を、ステップ3で捨てられた行数と比べてみてください。遅れたからといって、ウィンドウが閉じているとは限りません。
ウィンドウの前でフィルタリングすると、より多く捨てる
/root/flink/watermark/filtered.sqlに、5秒遅延のeventsから、CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts)の行だけを残した(ビューやサブクエリ)あと、同じ1分のTUMBLE集計をストリーミングで出すクエリを作成し、出力を、/root/flink/watermark/filtered.outに保存してください。
ドキュメントが、遅延行をフィルタリングするときに使うよう勧めている式です。行単位で「ウォーターマークより早いか」を見るので、ウィンドウがまだ開いていて受け取ってくれたはずの行まで捨ててしまいます。ウィンドウTVFの入力にビューを入れるには、CREATE VIEWのあと、TUMBLE(TABLE view_name, ...)の形で書きます。
ウォーターマークの間隔を1時間に: 前進が止まる
/root/flink/watermark/slow.sqlに、ステップ3のw5.sqlの先頭にSET 'pipeline.auto-watermark-interval' = '1 h';の1行だけを加えたクエリを書いて実行し、出力を、/root/flink/watermark/slow.outに保存してください。
ウォーターマークは周期(この設定)ごとに出力され、その間は、新しい値が最後に出した値より間隔を超えて先に進んだときだけ、レコードからすぐ出力されます。間隔が1時間なら、最初のウォーターマーク(それ以前に出したものがないのですぐ出ます)のあとは、ファイルが終わるまで進みません。そうすると、閉じるウィンドウはあるでしょうか。
報告書: 遅延と捨てられた行
/root/flink/watermark/report.jsonに、total_rows(batch.outのcntの合計)、dropped_0・dropped_5・dropped_30(バッチの合計から、各遅延のcntの合計を引いた値)、late_rows_5(late.outの行数)、dropped_slow(バッチの合計 − slow.outのcntの合計)を、整数で書いてください。
すべて、保存した出力から数えられます。結果の行は「| +I |」で始まり、バッチの結果の行は日付で始まります。awk -F'|'でcntの列を足せばよいです。sweep.outには結果のテーブルが2つあるので、ヘッダー行(opが入った行)を基準に分けて数えてください。