マージ前の重複を確かめ、FINAL・argMax・CLEANUP で扱う
目標
ReplacingMergeTreeに変更履歴(更新・削除)を入れ、マージ前には古いバージョンが見えること、それをFINALやargMaxで隠す方法、マージとCLEANUPが実際に何を消すのか、ソートキーと再送が結果をどう変えるかを、数値で確認します。
なぜ重要なのか
分析DBに移した本番データは、変わり続けます。ClickHouseは更新を「新しいバージョンのINSERT+あとのマージ」で処理しますが、マージはいつ起きるかわかりません。そのため、同じクエリが、正しいときも間違うときもあります。このラボでは、SYSTEM STOP MERGESで「マージ前」を固定してその間違った答えを目で確かめ、クエリ実行時に正しい答えを出す2つの方法を身に付けます。採点ツールは、期待値を元データの生成式から直接計算し、マージが実際に起きたかはsystem.part_logの記録で確認します。
ステップ
- データベース
rmtとテーブルrmt.usersを作成してください。列はuser_id UInt64, email String, plan LowCardinality(String), score UInt32, ver UInt32, is_deleted UInt8(この順序)、エンジンはReplacingMergeTree(ver, is_deleted)、ORDER BY user_idです。 SYSTEM STOP MERGES rmt.usersでこのテーブルのマージを止めたあと、/opt/lab/fixtures/replacing/changes.sqlを1回だけ実行してください(INSERT 3回=パート3つ)。- FINALなしで全行を数えるクエリ(/root/ch/replacing/q_raw.sql)と、
FROM rmt.users FINALで数えるクエリ(/root/ch/replacing/q_final.sql)を作成し、2つの結果をraw・finalとして書き込んでください(保存先: /root/ch/replacing/counts.json)。 - FINALを使わずに、料金プラン(plan)ごとの現在のユーザー数を求めるクエリを作成してください(/root/ch/replacing/q_plan.sql)。結果は、planとユーザー数の2つの列です。
rmt.mergedをrmt.usersと同じ列・エンジンで作成し、rmt.usersの行をver = 1、ver = 2、ver = 3の順に3回のINSERTで移して、OPTIMIZE TABLE rmt.merged FINALを実行したあと、rows_after_merge(count())・deleted_rows_kept(is_deleted = 1の行数)・final_count(FINALで数えた数)を書き込んでください(保存先: /root/ch/replacing/merged.json)。rmt.cleanedを同じ列・エンジンにテーブル設定allow_experimental_replacing_merge_with_cleanup = 1を付けて作成し、同じ方法で3回入れたあと、OPTIMIZE TABLE rmt.cleaned FINAL CLEANUPを実行して、rows_after_cleanup・deleted_rows_leftを書き込んでください(保存先: /root/ch/replacing/cleanup.json)。- 同じ列・エンジンでソートキーが
ORDER BY (user_id, plan)のrmt.badを作成し、同じ方法で3回入れてOPTIMIZE TABLE rmt.bad FINALを実行したあと、good_final(rmt.mergedをFINALで数えた数)・bad_final(rmt.badをFINALで数えた数)・bad_users(rmt.bad FINALの異なるuser_idの数)を書き込んでください(保存先: /root/ch/replacing/key.json)。 - 遅れた再送を真似ます。
rmt.usersからver = 1 AND user_id <= 1000の行を、rmt.mergedとrmt.cleanedにそれぞれもう一度入れ、2つのテーブルをFINALで数えた値とその差を、merged_final・cleaned_final・resurrectedとして書き込んでください(保存先: /root/ch/replacing/replay.json)。
参考
- サーバーはPodの起動時にすでに立ち上がっています。
clickhouse-clientと入力するだけで接続できます。止まっていた場合はch-upを実行してください(サーバーを立ち上げ直すとSTOP MERGESが解除されるので、ステップ2からやり直します)。 SYSTEM STOP MERGESは、そのテーブルにだけかかります。ステップ5–7の新しいテーブルはマージが有効なので、OPTIMIZEが実行できます。止めたテーブルでOPTIMIZEを実行すると、「Cancelled merging parts」で拒否されます。- マージの記録は
system.part_logに残ります(event_type = 'MergeParts'、merged_fromにまとめられたパート名)。今起きたことが見えない場合はSYSTEM FLUSH LOGSを実行します。 - よくある間違い: ステップ5–7で
INSERT ... SELECT * FROM rmt.usersの1回で移すことです。1回のINSERTの中の重複は、入れた瞬間に減る(optimize_on_insert)ので、マージするものがなくなります。バージョンごとに3回入れてください。 - 数値をJSONの中で引用符なしで受け取るには、
--output_format_json_quote_64bit_integers 0を使います。 - 公式ドキュメント: ReplacingMergeTree・Working with the ReplacingMergeTree engine・OPTIMIZE・system.part_log
バージョンと削除マーカーのあるテーブル
データベースrmtとテーブルrmt.usersを作成してください。列はuser_id UInt64, email String, plan LowCardinality(String), score UInt32, ver UInt32, is_deleted UInt8の順序、エンジンはReplacingMergeTree(ver, is_deleted)、ソートキーはORDER BY user_idです。
エンジンの第1引数は、どの行が勝つかを決めるバージョン列、第2引数は、勝った行が削除かどうかを示す列です。削除マーカーの列は、バージョン列なしでは使えません。ソートキーがそのまま「同じ行」の定義なので、変わらない識別子だけを置きます。
マージを止めて、変更履歴の3つのまとまりを入れる
SYSTEM STOP MERGES rmt.usersでこのテーブルのマージを止めたあと、/opt/lab/fixtures/replacing/changes.sqlを1回だけ実行してください。最初のロード・更新・削除の3つのINSERTが、パート3つとして残る必要があります。
マージはバックグラウンドでいつでも起きるので、放っておくと「マージ前」の状態が数秒で消えることがあります。止めるのは、入れる前に行う必要があります。すでに入れてしまった場合は、止めてからTRUNCATEして、もう一度入れてください。
FINALなしで数える場合とFINALで数える場合
FINALなしでrmt.usersの全行を数えるクエリ(/root/ch/replacing/q_raw.sql)と、FROM rmt.users FINALで数えるクエリ(/root/ch/replacing/q_final.sql)を作成し、2つのクエリの結果をraw・finalとして書き込んでください(保存先: /root/ch/replacing/counts.json)。
FINALは、クエリ実行中にマージのルール(ソートキーが同じならverが大きい行だけを残し、その行が削除マーカーなら除く)を適用します。2つの数値の差は、古いバージョンの行と削除の行を合わせた分です。q_final.sqlにWHEREを付けないでください。
FINALなしで同じ答えを出す: argMax
FINALを使わずに、料金プラン(plan)ごとの現在のユーザー数を求めるクエリを作成してください(/root/ch/replacing/q_plan.sql)。結果はplanとユーザー数の2つの列で、退会したユーザーは除く必要があります。
単にGROUP BY planで数えると、古いバージョンと削除の行まで数えてしまいます。まず、ユーザーごとに1行に減らす必要があります。user_idでグループ化し、argMax(値, ver)でバージョンが最も大きい行のplanとis_deletedを選んでから、外側で退会者を除き、planでもう一度グループ化します。
マージしても削除の行は残る
rmt.mergedをrmt.usersと同じ列・エンジンで作成し、rmt.usersの行をver = 1・ver = 2・ver = 3の順に3回のINSERTで移して、OPTIMIZE TABLE rmt.merged FINALを実行したあと、rows_after_merge(count())・deleted_rows_kept(is_deleted = 1の行数)・final_count(FINALで数えた数)を書き込んでください(保存先: /root/ch/replacing/merged.json)。
CREATE TABLE ... ASに別のテーブルを指定すると、列とエンジンをコピーしますが、STOP MERGESの状態はコピーしません。マージのあとの行数がユーザー数と同じになるか、それなのにFINALなしで数えた値がなぜまだ間違っているのかを、見てください。マージが実際に起きたかどうかは、system.part_logに残ります。
CLEANUPマージは削除の行まで消す
rmt.cleanedをrmt.usersと同じ列・エンジンにテーブル設定allow_experimental_replacing_merge_with_cleanup = 1を付けて作成し、ステップ5と同じように3回移したあと、OPTIMIZE TABLE rmt.cleaned FINAL CLEANUPを実行して、rows_after_cleanup(count())・deleted_rows_left(is_deleted = 1の行数)を書き込んでください(保存先: /root/ch/replacing/cleanup.json)。
CREATE TABLE ... ASで別のテーブルを指定した文の後ろに、SETTINGSを付けられます。設定なしでCLEANUPを実行すると、サーバーが拒否します。削除の行を消すと、あとから古いバージョンが入ってきたときに防げなくなるので、わざと有効にする必要がある機能です。
ソートキーが「同じ行」を決める
同じ列・エンジンでORDER BY (user_id, plan)のrmt.badを作成し、同じ方法で3回移してOPTIMIZE TABLE rmt.bad FINALを実行したあと、good_final(rmt.mergedをFINALで数えた数)・bad_final(rmt.badをFINALで数えた数)・bad_users(rmt.bad FINALの異なるuser_idの数)を書き込んでください(保存先: /root/ch/replacing/key.json)。
ソートキーにplanが入ると、料金プランが変わったユーザーの古い行と新しい行は、キーが違うので1つのグループになれません。削除の行も、退会直前の料金プランのグループだけを隠します。bad_finalがユーザー数より大きいか、bad_usersが現在のユーザー数より大きいかを見てください。
CLEANUPのあとに古い行がもう一度来ると
rmt.usersからver = 1 AND user_id <= 1000の行を、rmt.mergedとrmt.cleanedにそれぞれもう一度入れたあと、2つのテーブルをFINALで数えた値とその差を、merged_final・cleaned_final・resurrected(cleaned_final − merged_final)として書き込んでください(保存先: /root/ch/replacing/replay.json)。
パイプラインが古いバッチをもう一度送る状況です。削除マーカーが残っているテーブルでは何が古い行に勝つのか、削除の行を消したテーブルでは古い行が誰と競うのかを、考えてみてください。復活したのは、user_idが1–1000のユーザーのうち、退会していたユーザーです。