Apache Flink — ストリームを本物のエンジンで動かす
同じ GROUP BY — バッチの1行とストリーミングの変更ログ229行
目標
同じ集計をバッチとストリーミングで動かして、最終状態は同じで、ストリーミングはチェンジログを出力することを確認します。オペレーターとシンクによって、変更の種類(+I・-U・+U・-D)がどう決まるかを、実行計画と拒否エラーから読み取ります。
なぜ重要なのか
ストリーミングの結果は、1度書いて終わりの値ではなく、書き換え続けられるテーブルです。その書き換えを受け取る側がリトラクション(-U)を処理できなければ数字が膨らみ、追加のみのストレージはそもそも受け取れません。どのクエリが更新を出し、どのシンクがそれを受け取れるかを、計画から事前に読めてはじめて、設計の段階で事故を防げます。このラボの採点ツールは、クラスターに問い合わせません。皆さんが保存したsql-clientの出力とシンクのファイルを読み、チェンジログを元のCSVから1行ずつ再現して照合します(入力が追加のみで並列度が1なら、ログの順序まで決まります)。
ステップ
flink-upでクラスターを起動し、/root/flink/dynamic/batch.sqlに、バッチモードで/opt/lab/fixtures/data/dynamic_clicks.csvを読み、user_id別にclicks(件数)・dwell(dwellの合計)を出すSQLを書いたあと、出力を、/root/flink/dynamic/batch.outに保存してください。- /root/flink/dynamic/stream.sqlに、同じクエリをストリーミングモードで動かすSQLを作成し、出力を、/root/flink/dynamic/stream.outに保存してください。
- /root/flink/dynamic/nested.sqlに、クリック数(
clicks)別のユーザー数(users)を出す、集計の上の集計をストリーミングで動かすSQLを作成し、出力を、/root/flink/dynamic/nested.outに保存してください。 - /root/flink/dynamic/explain.sqlに、
EXPLAIN CHANGELOG_MODEの文を2つ(ユーザーごとのCOUNT集計、dwell > 60の行のuser_id, urlだけを選ぶフィルター)書き、出力を、/root/flink/dynamic/explain.outに保存してください。 - /root/flink/dynamic/fs-sink.sqlに、filesystemコネクターのテーブル
user_clicks_fs(user_id STRING, clicks BIGINT)を作成し、ユーザーごとのCOUNTをINSERT INTOするように書いたあと、拒否エラーが入った出力を、/root/flink/dynamic/fs-sink.outに保存してください。 - /root/flink/dynamic/to-changelog.sqlで、ユーザーごとのCOUNTのビューを
TO_CHANGELOGで変換して、filesystemシンク(パスは/root/flink/dynamic/changelog、列はchange, user_id, clicks)に書き、出力を、/root/flink/dynamic/to-changelog.outに保存してください。 - /root/flink/dynamic/upsert.sqlに、printコネクターのテーブル
retract_sinkと、PRIMARY KEY (user_id) NOT ENFORCEDを置いたblackholeコネクターのテーブルupsert_sinkを作成し、2つのシンクに同じ集計を入れるEXPLAIN CHANGELOG_MODE INSERT INTO ...をそれぞれ実行して、出力を、/root/flink/dynamic/upsert.outに保存してください。 - /root/flink/dynamic/report.jsonに、
users・retract_messages・upsert_messages・nested_deletesを書いてください。
参考
- 元の列:
click_id BIGINT, user_id STRING, url STRING, dwell INT, ts TIMESTAMP(3)(ヘッダーなしのCSV、ファイルの順序 = 到着順序)。 - SQLファイルの実行:
sql-client.sh -f 파일.sql > 파일.out 2>&1(プレースホルダーはファイル名です)。エラーが起きた文で止まり、出力に[ERROR]が残ります。 - ストリーミングの結果テーブルは、先頭に
op列が付きます。終わりのあるファイルを読むので、ストリーミングジョブもファイルを読み終えると終了します。 - よくある間違い: ステップ6を2回実行すると、シンクのディレクトリにファイルがたまります。もう一度実行するときは、ディレクトリを先に空にしてください。
INSERTはデフォルトが非同期の提出なので、ジョブが終わる前にsql-clientが終了します。SET 'table.dml-sync' = 'true';を指定すると、終わるまで待ちます。 - よくある間違い:
SET 'table.exec.mini-batch.enabled'のような設定を有効にすると、途中のログがまとめられて、行数が変わります。このラボは、デフォルト設定のままにします。 - 公式ドキュメント: Dynamic Tables・EXPLAIN・Changelog Conversion・FileSystem・Print・BlackHole
バッチで一度に数える
flink-upでクラスターを起動し、/root/flink/dynamic/batch.sqlにSET 'execution.runtime-mode' = 'batch';・元のCSVを読むCREATE TABLE・user_id別にclicks(COUNT(*))とdwell(SUM(dwell))を出すSELECTを書き、出力を、/root/flink/dynamic/batch.outに保存してください。
filesystemコネクターで、pathはfile:///opt/lab/fixtures/data/dynamic_clicks.csv、formatはcsvです。結果の列にはAS clicks・AS dwellで別名を付けます。バッチの結果にはop列がなく、ユーザーごとに1行ずつ出ます。
同じクエリをストリーミングで動かす
/root/flink/dynamic/stream.sqlに、ステップ1と同じクエリをSET 'execution.runtime-mode' = 'streaming';で動かすSQLを作成し、出力を、/root/flink/dynamic/stream.outに保存してください。ログを最後まで適用した最終状態が、バッチの結果と同じである必要があります。
行が1つ来るたびに、結果テーブルが変わります。初めて見るユーザーは+Iの1行、すでにいたユーザーは、古い値を消す-Uと新しい値+Uの2行です。ログの行数は、ユーザー数 + 2 × (クリック数 − ユーザー数)になるはずです。設定はデフォルトのままにします。
集計の上の集計で-Dを見る
/root/flink/dynamic/nested.sqlに、ストリーミングモードでSELECT clicks, COUNT(*) AS users FROM (사용자별 COUNT(*) AS clicks) GROUP BY clicks(プレースホルダーは、ユーザーごとのCOUNTのサブクエリです)を実行するSQLを作成し、出力を、/root/flink/dynamic/nested.outに保存してください。
ユーザーがクリック数1のグループからクリック数2のグループへ移るとき、内側の集計は-U(1)と+U(2)を出します。外側の集計は-Uを受け取って、クリック数1のグループのユーザー数を1つ減らし、そのグループが0になると結果の行が消えなければならないので、-Dを出します。内側のクエリの列名をclicksに揃えてください。
計画から変更の種類を読む
/root/flink/dynamic/explain.sqlに、EXPLAIN CHANGELOG_MODEで始まる文を2つ、すなわち、ユーザーごとのCOUNT集計と、dwell > 60の行のuser_id, urlだけを選ぶフィルターを書き、出力を、/root/flink/dynamic/explain.outに保存してください。
EXPLAINは、ジョブを動かさず、計画だけを出力します。CHANGELOG_MODEを付けると、Optimized Physical Planのノードごとにチェンジログモード(changelogMode=[...])が付きます。Iは追加、UBは更新前、UAは更新後、Dは削除です。フィルターだけのクエリのノードは、何だけを出すでしょうか。
更新をファイルシンクに入れて拒否される
/root/flink/dynamic/fs-sink.sqlに、filesystemコネクターのテーブルuser_clicks_fs(user_id STRING, clicks BIGINT)(パスは/root/flink/dynamic/fs-out、formatはcsv)を作成し、INSERT INTO user_clicks_fs SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id;を書いたあと、出力を、/root/flink/dynamic/fs-sink.outに保存してください。エラーで終わるのが正常です。
ファイルは、1度書いた行を直せません。オプティマイザーは、シンクが受け取れる変更の種類を先に確認し、集計が出すUB/UAを受け取れなければ、ジョブを提出する前に拒否します。エラーの文に、どのシンクがどのノードの何を受け取れないかが書かれています。
変更の種類を列として取り出してファイルに書く
/root/flink/dynamic/to-changelog.sqlで、ユーザーごとのCOUNTをビュー(user_id、clicksの2列)として作成し、TO_CHANGELOG(input => TABLE 뷰, op => DESCRIPTOR(change))(プレースホルダーはビューの名前です)の結果を、filesystemシンク(パスは/root/flink/dynamic/changelog、列はchange STRING, user_id STRING, clicks BIGINT、formatはcsv)にINSERTしてください。出力は、/root/flink/dynamic/to-changelog.outに保存します。
TO_CHANGELOGは、更新ログの各行を追加(INSERT)に変え、元の変更の種類を文字列の列に入れます。そのため、追加しかしないファイルシンクでも受け取れます。INSERTは非同期なので、table.dml-syncを有効にしないと、ジョブが終わったあとにファイルが確定した状態で残りません。もう一度実行するときは、シンクのディレクトリを先に空にしてください。
キーのあるシンクの前で-Uが消える
/root/flink/dynamic/upsert.sqlに、printコネクターのテーブルretract_sink(user_id STRING, clicks BIGINT)と、同じ列にPRIMARY KEY (user_id) NOT ENFORCEDを置いたblackholeコネクターのテーブルupsert_sinkを作成してください。そして、ユーザーごとのCOUNTをそれぞれのシンクに入れるEXPLAIN CHANGELOG_MODE INSERT INTO ...の2つの文を実行し、出力を、/root/flink/dynamic/upsert.outに保存してください。
printシンクは、受け取ったものをそのまま出力するだけなので、古い行を消すにはリトラクション(UB)が必要です。キーが宣言されたシンクは、キーで探して上書きすればよいので、UBは不要です。2つの計画で、GroupAggregateのchangelogModeを比べてみてください。NOT ENFORCEDは、Flinkがキーを検査しないという意味です。
報告書: リトラクトとアップサートのメッセージ数
/root/flink/dynamic/report.jsonに、users(ユーザー数 = stream.outの+Iの数)、retract_messages(stream.outのログの全行数)、upsert_messages(同じ結果をアップサートシンクに送っていたなら送ったはずのメッセージ数)、nested_deletes(nested.outの-Dの数)を、整数で書いてください。
アップサートシンクは、キーで上書きするので、更新前の値(-U)を受け取る必要がありません。同じログからどのopだけが残るかを考えてみてください。数字は、保存した出力のop列を数えれば出ます。grepで、行頭の「| +I |」のような形を数えられます。