TT Lab
はじめる
学ぶ 学習パス コース

Apache Flink — ストリームを本物のエンジンで動かす

同じ GROUP BY — バッチの1行とストリーミングの変更ログ229行

TT Labで続きを見る

目標

同じ集計をバッチとストリーミングで動かして、最終状態は同じで、ストリーミングはチェンジログを出力することを確認します。オペレーターとシンクによって、変更の種類(+I・-U・+U・-D)がどう決まるかを、実行計画と拒否エラーから読み取ります。

なぜ重要なのか

ストリーミングの結果は、1度書いて終わりの値ではなく、書き換え続けられるテーブルです。その書き換えを受け取る側がリトラクション(-U)を処理できなければ数字が膨らみ、追加のみのストレージはそもそも受け取れません。どのクエリが更新を出し、どのシンクがそれを受け取れるかを、計画から事前に読めてはじめて、設計の段階で事故を防げます。このラボの採点ツールは、クラスターに問い合わせません。皆さんが保存したsql-clientの出力とシンクのファイルを読み、チェンジログを元のCSVから1行ずつ再現して照合します(入力が追加のみで並列度が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に保存してください。
  2. /root/flink/dynamic/stream.sqlに、同じクエリをストリーミングモードで動かすSQLを作成し、出力を、/root/flink/dynamic/stream.outに保存してください。
  3. /root/flink/dynamic/nested.sqlに、クリック数(clicks)別のユーザー数(users)を出す、集計の上の集計をストリーミングで動かすSQLを作成し、出力を、/root/flink/dynamic/nested.outに保存してください。
  4. /root/flink/dynamic/explain.sqlに、EXPLAIN CHANGELOG_MODEの文を2つ(ユーザーごとのCOUNT集計、dwell > 60の行のuser_id, urlだけを選ぶフィルター)書き、出力を、/root/flink/dynamic/explain.outに保存してください。
  5. /root/flink/dynamic/fs-sink.sqlに、filesystemコネクターのテーブルuser_clicks_fs(user_id STRING, clicks BIGINT)を作成し、ユーザーごとのCOUNTをINSERT INTOするように書いたあと、拒否エラーが入った出力を、/root/flink/dynamic/fs-sink.outに保存してください。
  6. /root/flink/dynamic/to-changelog.sqlで、ユーザーごとのCOUNTのビューをTO_CHANGELOGで変換して、filesystemシンク(パスは/root/flink/dynamic/changelog、列はchange, user_id, clicks)に書き、出力を、/root/flink/dynamic/to-changelog.outに保存してください。
  7. /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に保存してください。
  8. /root/flink/dynamic/report.jsonに、users・retract_messages・upsert_messages・nested_deletesを書いてください。

参考

バッチで一度に数える

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 |」のような形を数えられます。