Apache Flink — ストリームを本物のエンジンで動かす
重複排除と Top-N — 同じ ROW_NUMBER、異なる変更ログ
一言でいうと
Flink SQLの重複排除とTop-Nは、どちらもROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...)に順位番号の条件を付けた同じパターンです。ただし、先頭行を残すと結果がappend-onlyになり、末尾行を残すか順位を付けると、すでに出した行をリトラクションして直す更新ログになります。どちらなのかが、下流のシンクと集計を決めます。
なぜ必要なのか
公式ドキュメントが挙げる例は、現実そのままです。上流のETLがエンドツーエンドで正確に1回を保証できないと、障害復旧のときに同じレコードがシンクに2回入ります。その状態でSUM・COUNTをすると、数字が膨らみます。そのため、分析の前に重複を取り除く必要があります。
ところが、「重複」には2つの意味が混ざっています。1つは同じ出来事が2回来た場合(再送)で、最初の1つだけを残せばよいです。もう1つは同じ対象の状態が何度も変わった場合(注文が作成 → 支払い → 配送)で、現在の状態、つまり最後のものを残す必要があります。2つの要求は、SQLではASCとDESCの1語の違いですが、エンジンの中ではまったく違うことをします。先頭行は1回出せば終わりですが、末尾行は「これまでの最後」が変わり続けるからです。
どう動くのか
パターンをそのまま守る必要があります。ドキュメントは、重複排除をROW_NUMBER()・PARTITION BY 키・ORDER BY 시간 속성・外側のWHERE rownum = 1で定義し(プレースホルダーは順に、キーと時間属性です)、この形を正確に守ってはじめて、オプティマイザーが認識すると記しています。ORDER BYは必ず時間属性(処理時間またはイベント時間)でなければならず、ASCは先頭行を、DESCは末尾行を残します。理論的には、重複排除は、Nが1で、時間で並べたTop-Nの特殊なケースです。
計画で見ると、次のように分かれます。EXPLAIN CHANGELOG_MODEの物理計画は、このオペレーターをRank(strategy=[AppendFastStrategy], rankRange=[rankStart=1, rankEnd=1], ...)と書き、行末に出力する変更の種類を付けます。実行計画では、同じ場所がDeduplicateノードに変わります。
물리 계획 Rank(... orderBy=[ROWTIME ts ASC] ...) changelogMode=[I]
실행 계획 Deduplicate(keep=[FirstRow], key=[order_id], order=[ROWTIME], outputInsertOnly=[true])
물리 계획 Rank(... orderBy=[ROWTIME ts DESC] ...) changelogMode=[I,UA,D]
실행 계획 Deduplicate(keep=[LastRow], key=[order_id], order=[ROWTIME], outputInsertOnly=[false])
このコードブロックの韓国語の語句は、順に、物理計画と実行計画(Optimized Physical PlanとOptimized Execution Plan)を表しています。
先頭行を残す方式は、キーごとに「すでに見た」という印だけを状態に置き、2番目からは捨てます。結果はinsert-onlyなので、ファイルのようなappend-onlyシンクにもそのまま入ります。末尾行を残す方式は、キーごとにこれまでの最後の行を状態に置き、新しい行が来ると、先に出した行をリトラクションして、新しい行を出します。実測で、イベント622行・注文240件を入れたところ、+Iが240個、-U/+Uが382組出ました。382 = 622 − 240、つまり同じキーの2番目以降の行ごとに1組です。まったく同じ行が再送されても、組が出ました。処理時間(PROCTIME())で並べると、同じ入力で再送41件のうち18件の組が出ませんでした。到着時刻に頼る並べ方なので、このような件数は実行環境によって変わることがあり、ラボはイベント時間で判定します。
Top-Nは、順位番号の条件が<= Nである同じパターンです。ドキュメントは、Top-Nが結果更新型で、上位Nが変わると、変わった行をリトラクションと更新で送ると記しています。入力自体が更新される場合(販売量SUMの上の順位)なら、計画にRank(strategy=[RetractStrategy], ...)が出ます。ここで重要な選択が、順位番号を出力するかです。順位番号の列は結果の一意キーの一部になるので、ドキュメントの例のように、9位が1位に上がると、1–9位の行がすべて送り直されます。外側のSELECTから順位番号を除くと、変わった商品1つだけを送れば済みます。実測で、同じ販売ファイルのカテゴリ別Top-3は、順位番号を出力したときはログが3,419行、除いたときは1,523行で、最終結果は同じでした。
末尾行を残す方式の上にGROUP BY statusを載せると、2つのオペレーターが噛み合います。注文1つがcreatedからpaidに変わると、重複排除が-U created・+U paidを出し、集計はそれを受けて、createdのグループから1つ引き、paidのグループに1つ足します。重複排除なしでイベントをそのまま数えると、1つの注文が複数のグループに同時に数えられます。
前のモジュールのテンポラル結合でバージョン付きビューを作ったのも、まさにこの末尾行を残す方式です。ウィンドウ単位のTop-Nは、ウィンドウのモジュールで扱いました。そこではウィンドウが閉じるときに1回だけ出すので、リトラクションがありません。
現場での姿
最もよくある失敗は、「重複排除の結果をファイルやappend-onlyのトピックに書く」ことです。先頭行を残す方式なら問題ありませんが、末尾行を残す方式に変えた瞬間に、結果が更新ログになり、更新を受け取れないシンクには入れられません(動的テーブルのモジュールで見た拒否です)。どちらが必要かを、要求から先に分けるのが先です。再送を取り除くのなら先頭行、現在の状態が必要なら末尾行です。
2つ目は、Top-Nのログの急増です。リアルタイムのランキング表をキーバリューストアに書いていて、書き込み量が予想の数倍なら、順位番号を出力していないかをまず確認します。画面が自分で順位を付けられるなら、順位番号を除くだけで、書き込みが大きく減ります。
3つ目は、並べ替えの基準です。処理時間で重複排除すると、結果が到着順に左右されて、再実行のたびに違うことがあります。再処理しても同じ答えが出なければならないパイプラインなら、イベント時間を使います。
次のラボですること
注文状態のイベントと販売記録を定義して、行数と注文数を数えます。先頭行を残す方式と末尾行を残す方式を動かして、ログの件数がどう違うかを確認し、EXPLAINで2つの計画のチェンジログモードを比べます。末尾行を残す方式の上で状態別の現況を数え、カテゴリ別のTop-3を、順位番号ありと順位番号なしで動かして、ログ量を比べたあと、報告書にまとめます。