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

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

チェックポイントとセーブポイント — 止めて再開しても重複しない理由

TT Labで続きを見る

一言でいうと

チェックポイントは、ソースがどこまで読んだか、オペレーターが何を保持しているか、シンクがまだ何を確定していないかを、1つの瞬間に合わせて撮った写真です。ファイルシンクは、その写真が完成したという通知を受けてはじめてファイルを確定するので、ジョブが死んで復活しても、確定した結果には欠落も重複もありません。セーブポイントは、同じ仕組みで人が撮っておく写真です。

なぜ必要なのか

ストリーミングジョブは、何週間も動きます。その間に、TaskManagerが死に、コードを直して再デプロイし、クラスターを移します。再開する瞬間に、3つのことを同時に決めなければなりません。ソースのどこから再び読むか、これまでに積み上げた合計・ウィンドウのような状態をどう復元するか、すでに外に出した結果をどうするかです。

3つを別々に決めると、必ずずれます。最初から読み直すと、すでに書いた結果が2回出力され、最後に読んだ位置から続きを読むと、状態が空なので合計が間違います。状態だけを別に保存しておいても、保存した瞬間とソースの位置が何件でもずれると、そのぶん2回数えたり、取りこぼしたりします。必要なのは、「3つが同じ瞬間だった」という保証です。Flinkはこれを、データの流れの中にマーカーを流す方法で解決します。

どう動くのか

上の行は、連番ソースが1から順に出力する行で、その間にチェックポイントバリアnとn+1が挟まって流れます。下の行はファイルシンクのディレクトリです。バリアnがシンクに届いてチェックポイントnが完了すると、そこまで書いていたドットファイル(.part-…inprogress)がpart-…に名前が変わって確定し、新しいドットファイルに次の行が書かれます。右は、STOP JOB WITH SAVEPOINTでセーブポイントを取りながら停止すると、書き込み中だったファイルまで確定し、そのセーブポイントから復元した新しいジョブは、ソースの次の連番から続きを書くことを、下の分岐は、セーブポイントなしで新しく始めると、1からもう一度書いて重なることを示しています

バリアは、次のとおりです。JobManagerのチェックポイントコーディネーターが、ソースにバリアnを挟み込みます(公式ドキュメントのStateful Stream Processing)。バリアはレコードと同じ道を流れ下り、オペレーターはバリアを受け取った瞬間に、自分の状態を取得して保存したあと、バリアを下流に渡します。ソースの「状態」は、読んだ位置です。連番ソースなら、次に出力する番号です。すべてのタスクが写真を撮り終えたと報告すると、チェックポイントnが完了し、コーディネーターが完了を改めて通知します。

ファイルシンクは、この通知を待ちます。シンクが書くファイルは、3つの段階を経ます。書き込み中は、名前がドットで始まる.part-…inprogress…で、バリアを受け取ると閉じられて確定を待ち、完了の通知が来ると、part-…に名前が変わって確定します。ドットファイルは、慣例として、読み取る側がスキップする隠しファイルです。そのため、下流が見る結果は、常に「どのチェックポイントまで」の結果であり、チェックポイントの間隔が、そのまま結果が見えるまでの遅延になります。このPodで1秒間隔で動かすと、1秒ごとにpartファイルが1つずつ増え、ドットファイルは、チェックポイントの間にほんの一瞬だけ見えて消えます。

保存場所は、次のとおりです。execution.checkpointing.dirを指定すると、ファイルシステムに保存されます。ドキュメントが記す構造は<dir>/<job-id>/chk-<n>/で、その中の_metadataが写真の目次です。デフォルトでは、チェックポイントを保持しません。新しいチェックポイントが完了すると古いものを削除し、ジョブをキャンセルするとすべて削除します。障害復旧用であり、人が使うために残すものではありません。残したいなら、execution.checkpointing.externalized-checkpoint-retentionをRETAIN_ON_CANCELLATIONにします。

セーブポイントは、同じ仕組みで撮りますが、所有者は人です。<savepoint-dir>/savepoint-<잡 id 앞 6자리>-<무작위>/(プレースホルダーは順に、ジョブidの先頭6桁とランダム文字列です)にできて、Flinkが勝手に削除することはありません。ドキュメントが強調する落とし穴が1つあります。Flink 1.15から、停止せずに取った途中のセーブポイントは、副作用をコミットしません。ファイルを確定してくれるのは、チェックポイントとSTOP ... WITH SAVEPOINTです。

SET 'execution.checkpointing.savepoint-dir' = 'file:///root/flink/checkpoint/sp';
STOP JOB '<jid>' WITH SAVEPOINT;                       -- 찍고 멈춘다. 쓰던 파일까지 확정
SET 'execution.state-recovery.path' = 'file:/.../savepoint-xxxxxx-yyyy';
INSERT INTO sink SELECT id FROM seq;                   -- 같은 질의를 되살린다 — 다음 순번부터

このコードブロックの2つの韓国語コメントは、順に、セーブポイントを取って停止し、書き込み中だったファイルまで確定するという意味と、同じクエリを復元して次の連番から続けるという意味です。

復元したジョブは、セーブポイントのソースの位置とシンクの状態を一緒に受け取ります。そのため、シンクのパスを変えても連番は続き、2つのディレクトリを合わせると、欠落も重複もありません。保持されたチェックポイントからも、同じように復元できます(ドキュメント: チェックポイントのメタデータファイルで、セーブポイントのように再開)。

現場での姿

最もよくある事故は、セーブポイントなしの再デプロイです。SQLを少し直してもう一度INSERTすると、新しいジョブは空の状態で、ソースの最初(または設定された開始位置)から読みます。連番ソースなら、1からもう一度書きます。結果のディレクトリが同じなら、すでにある行が、もう1回入ります。このラボで、その重なりを自分で数えます。

2つ目は、キャンセルで止まったジョブです。キャンセルはセーブポイントを取らないので、最後のチェックポイントのあとに書いていたドットファイルが、ディレクトリにそのまま残ります。保持したチェックポイントから復元すると、新しいジョブはチェックポイント時点から書き直し、結果(ドットで始まらないファイル)には、今も欠落・重複がありません。ところが、ドットファイルまで読むツールを下流につないでおくと、その瞬間に重複が生まれます。「正確に1回」は、確定したものだけを読むという約束とセットです。

3つ目は、復元の失敗です。セーブポイントは、オペレーターidごとに状態を保持します。ドキュメントは、自動で作られたidがプログラムの構造に敏感だと警告しています。SQLの形を大きく変えると、状態を当てはめるオペレーターが見つかりません。状態のあるジョブを直すときは、「この変更のあとでも、セーブポイントが入るか」を、先に試します。

再起動戦略は、最初のモジュールで見たとおり、チェックポイントを有効にするとデフォルトが再起動する側に変わります。このモジュールは、その再起動がどこへ戻るのかを扱います。

次のラボですること

1秒ごとにチェックポイントを取る連番ジョブをファイルシンクにつなぎ、動いている間に、確定ファイルとドットファイルが一緒にある一覧を残します。/checkpointsのレスポンスから最後のチェックポイントの位置を読み、STOP JOB ... WITH SAVEPOINTで停止したあと、セーブポイントから復元して、連番が続くかを確認します。続けて、セーブポイントなしで新しく動かして重なりを作り、キャンセルしても残したチェックポイントから続きを書いたあと、すべての数字をディスクから数え直して、報告書にまとめます。