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

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

止めて復元しても連番は続くか

TT Labで続きを見る

目標

連番ソース(1、2、3 …)をファイルシンクにつなぎ、チェックポイント・セーブポイントで停止して復元したあと、確定したファイルのidが、欠落・重複なしに続くかを、ディスクで直接確認します。セーブポイントなしでもう一度動かすと何が重なるか、キャンセルしても残したチェックポイントが何を守ってくれるかも数えます。

なぜ重要なのか

再デプロイ・障害復旧のとき、ソースの位置、オペレーターの状態、シンクの未確定ファイルが、1つの瞬間として合っていないと、結果が抜けたり、2回出力されたりします。チェックポイントは、この3つをバリア1つで合わせ、ファイルシンクは、チェックポイントが完了してはじめてファイルを確定します。このラボの採点ツールは、クラスターに問い合わせません。停止した時点の行数は実行ごとに違うので、正解の数字がありません。代わりに、確定ファイルのidが1から欠落・重複なしに続くか、復元したジョブがすぐ次のidから始まるか、セーブポイント・チェックポイントのディレクトリに_metadataが残っているかを見ます。

ステップ

  1. flink-upでクラスターを起動してください。/root/flink/checkpoint/run1.sqlに、チェックポイント間隔1秒・チェックポイントのディレクトリfile:///root/flink/checkpoint/ckpt・セーブポイントのディレクトリfile:///root/flink/checkpoint/sp・ジョブ名flk-seq-1で、datagenの連番(id 1..100000、毎秒20行)をfile:///root/flink/checkpoint/out1にcsvで書くSQLを作成して提出し、出力を、/root/flink/checkpoint/run1.outに保存してください。
  2. ジョブが動いている間に、ls -A /root/flink/checkpoint/out1の結果を、/root/flink/checkpoint/files-running.txtに保存してください。確定ファイル(part-…)と書き込み中のファイル(.part-…inprogress…)が、一緒に見える必要があります。
  3. そのジョブの/jobs/<jid>/checkpointsを、/root/flink/checkpoint/checkpoints.jsonに、/jobs/<jid>/checkpoints/configを、/root/flink/checkpoint/checkpoint-config.jsonに保存してください。
  4. /root/flink/checkpoint/stop.sqlでSTOP JOB '<jid>' WITH SAVEPOINTを実行し、出力を、/root/flink/checkpoint/stop.outに保存してください。停止したあと、out1にドットファイルが残っていてはいけません。
  5. /root/flink/checkpoint/run2.sqlに、セーブポイントから復元するSQL(ジョブ名flk-seq-2、シンクfile:///root/flink/checkpoint/out2)を作成して提出し、出力を、/root/flink/checkpoint/run2.outに置き、数秒後にRESTでセーブポイントを取りながら停止して、完了した状態のレスポンスを、/root/flink/checkpoint/stop2.jsonに保存してください。
  6. /root/flink/checkpoint/run3.sqlに、復元なしで最初から動くSQL(ジョブ名flk-seq-3、シンクfile:///root/flink/checkpoint/out3、チェックポイントの保持RETAIN_ON_CANCELLATION)を作成して提出し、出力を、/root/flink/checkpoint/run3.outに残し、ファイルが確定したあと、ジョブをキャンセルしてください。キャンセルされたジョブの/jobs/<jid>を、/root/flink/checkpoint/job3.jsonに、/jobs/<jid>/checkpointsを、/root/flink/checkpoint/checkpoints3.jsonに保存してください。
  7. /root/flink/checkpoint/run4.sqlに、残ったチェックポイントから復元して、同じout3に続きを書くSQL(ジョブ名flk-seq-4)を作成して提出し、出力を、/root/flink/checkpoint/run4.outに残し、数秒後にRESTでセーブポイントを取りながら停止して、状態のレスポンスを、/root/flink/checkpoint/stop4.jsonに保存してください。
  8. /root/flink/checkpoint/report.jsonに、savepoint_path・last_id_before_stop・first_id_after_resume・fresh_run_overlap・retained_checkpoint・leftover_inprogress_filesを書いてください。

参考

チェックポイントを有効にした連番ジョブを提出する

flink-upのあと、/root/flink/checkpoint/run1.sqlに、チェックポイント間隔1 s、execution.checkpointing.dir = file:///root/flink/checkpoint/ckpt、execution.checkpointing.savepoint-dir = file:///root/flink/checkpoint/sp、ジョブ名flk-seq-1、datagenの連番ソース、file:///root/flink/checkpoint/out1に書くcsvシンクとINSERTを書いて、sql-client.sh -f run1.sql > run1.out 2>&1で、/root/flink/checkpoint/run1.outを作成してください。

SET文は、INSERTより前に置かないと、そのジョブに効きません。連番ソースはfields.id.kind = sequenceで、開始・終わりをfields.id.start・fields.id.endで指定します。毎秒の行数を小さく(20)すると、ファイルが少ししかできず、目で追いやすいです。出力にJob IDが見えれば、ジョブは裏で動き続けています。

動作中のシンクのディレクトリを見る

ジョブが動いている間に、ls -A /root/flink/checkpoint/out1の出力を、/root/flink/checkpoint/files-running.txtに保存してください。確定したpart-…と書き込み中の.part-…inprogress…が、一緒に見える必要があります。

チェックポイントが1回完了してはじめて、最初のpartファイルが確定します。1秒間隔なら、数秒待てばよいです。ドットで始まるファイルは、-Aなしでは見えず、このPodで測ると、チェックポイントの間にほんの一瞬だけ現れます。0.1秒間隔で何回も一覧を取っておき、両方が一緒に見えたときのものを保存してください。採点ツールは、一覧の確定ファイルが今もout1にあるかを照合します。

チェックポイントの記録をRESTで取得する

run1のジョブの/jobs/<jid>/checkpointsを、/root/flink/checkpoint/checkpoints.jsonに、/jobs/<jid>/checkpoints/configを、/root/flink/checkpoint/checkpoint-config.jsonに保存してください。完了したチェックポイントが1つ以上あり、最後のもののパスが、このジョブのckpt/<jid>/chk-<번호>(プレースホルダーは番号です)である必要があります。

ジョブidは、run1.outのJob IDの行にあります。checkpointsのレスポンスのlatest.completedに、最後に完了したチェックポイントの番号とexternal_pathが、configのレスポンスに、モード(正確に1回)・間隔(ミリ秒)・保持の有無(externalization)があります。

セーブポイントを取って停止する

/root/flink/checkpoint/stop.sqlに、セーブポイントのディレクトリのSETと、STOP JOB '<run1 의 jid>' WITH SAVEPOINT;(プレースホルダーはrun1のjidです)を書いて、出力を、/root/flink/checkpoint/stop.outに保存してください。停止したあと、out1にはドットファイルがなく、確定ファイルのidが1から欠落・重複なしに続いている必要があります。

STOP JOBはsql-clientで動く文で、セーブポイントのパスを、表の1つのセルとして返します。パスは、sp/savepoint-(ジョブidの先頭6桁)-…という形です。セーブポイントが完了する瞬間に、書き込み中だったファイルも確定するので、停止したあとにドットファイルはないはずです。残っていれば、キャンセルで止まったのです。

セーブポイントから復元して続きを書く

/root/flink/checkpoint/run2.sqlに、先頭にSET 'execution.state-recovery.path' = '<stop.out 의 세이브포인트>';(プレースホルダーはstop.outのセーブポイントです)を置き、ジョブ名をflk-seq-2、シンクをfile:///root/flink/checkpoint/out2に変えたSQLを書いて提出し、出力を、/root/flink/checkpoint/run2.outに残してください。数秒後にREST(POST /jobs/<jid>/stop)でセーブポイントを取りながら停止し、COMPLETEDになった状態のレスポンスを、/root/flink/checkpoint/stop2.jsonに保存してください。out2の最初のidは、out1の最後のすぐ次である必要があります。

セーブポイントには、ソースが次に出力する番号も状態として入っています。そのため、シンクのパスを変えても、番号は続きます。RESTで停止すると、すぐにrequest-idだけが返ってくるので、/jobs//savepoints/をCOMPLETEDになるまで取得し直して、保存します。

セーブポイントなしで新しく動かしてキャンセルする

/root/flink/checkpoint/run3.sqlに、復元パスなしで、ジョブ名flk-seq-3・シンクfile:///root/flink/checkpoint/out3・SET 'execution.checkpointing.externalized-checkpoint-retention' = 'RETAIN_ON_CANCELLATION';を入れたSQLを書いて提出し、出力を、/root/flink/checkpoint/run3.outに残してください。out3に確定ファイルができたあと、ジョブをキャンセルし、キャンセルされたジョブの/jobs/<jid>を、/root/flink/checkpoint/job3.jsonに、/jobs/<jid>/checkpointsを、/root/flink/checkpoint/checkpoints3.jsonに保存してください。

復元しなかったジョブは、ソースが1からもう一度数えます。out1と同じ番号が出ます。キャンセルはセーブポイントを取りません。保持の設定がないと、キャンセルとともにチェックポイントのディレクトリが削除されるので、採点ツールは、checkpoints3.jsonが指すchkディレクトリに_metadataが残っているかを見ます。

残したチェックポイントから続きを書く

/root/flink/checkpoint/run4.sqlに、checkpoints3.jsonの最後の完了したチェックポイントを復元パスにして、ジョブ名flk-seq-4・シンクは同じfile:///root/flink/checkpoint/out3のSQLを書いて提出し、出力を、/root/flink/checkpoint/run4.outに残してください。数秒後にRESTでセーブポイントを取りながら停止して、状態のレスポンスを、/root/flink/checkpoint/stop4.jsonに保存してください。out3の確定ファイルは、1から欠落・重複なしに続く必要があります。

チェックポイントのディレクトリ(chk-N)も、セーブポイントのように復元パスに使えます。新しいジョブは、チェックポイント時点の番号から書き直し、キャンセルのときに書きかけだったドットファイルは、結果ではありません。採点ツールは、ドットで始まらないファイルだけを集めて見ます。ファイル名のラベル(uuid)が違う2つの実行のファイルが、混ざっている必要があります。

報告書: ディスクで数え直す

/root/flink/checkpoint/report.jsonに、savepoint_path(stop.outのセーブポイント)、last_id_before_stop(out1の確定ファイルの最後のid)、first_id_after_resume(out2の最初のid)、fresh_run_overlap(run3、つまり1からもう一度動いた実行が、out3に確定したidのうち、out1・out2にもあるものの数)、retained_checkpoint(checkpoints3.jsonの最後の完了したパス)、leftover_inprogress_files(現在out3に残っているドットファイルの数)を書いてください。

すべての値は、ディスクで数え直せます。確定ファイルは名前がpart-で始まり、1回の実行のファイルは、同じラベル(uuid)を共有します。run3のファイルは、1が入ったファイルとラベルが同じです。重なりは、2つのidの一覧の積集合の大きさです(comm -12)。