Fault-tolerance in Flink
1. ステートフルなフローコンピューティング
ストリームコンピューティング
ストリームコンピューティングとは、メッセージを継続的に送信できるデータソースがあり、同時にコードを実行する常駐プログラムが存在することを意味します。データソースからメッセージを受信すると、それを処理して結果を下流に出力します。
分散ストリームコンピューティング
分散ストリームコンピューティングとは、入力ストリームを何らかの方法で分割し、複数の分散インスタンスを使用してストリームを処理することを指します。
ストリームコンピューティングにおける状態
コンピューティングはステートフルとステートレスに分けることができます。ステートレスコンピューティングは単一のイベントを処理するだけで済みますが、ステートフルコンピューティングは複数のイベントを記録して処理する必要があります。
簡単な例を挙げます。たとえば、イベントがイベント ID とイベント値の 2 つの部分で構成されているとします。処理ロジックが、イベントを取得するたびにそのイベント値を解析して出力するだけであれば、これはステートレスコンピューティングです。一方、イベントを取得するたびに、その値を解析した後に前のイベント値と比較し、前のイベント値より大きい場合にのみ出力する必要がある場合、これはステートフルコンピューティングです。
ストリームコンピューティングには多くの状態があります。たとえば、重複排除シナリオではすべての主キーが記録されます。また、ウィンドウ計算では、ウィンドウに入りまだトリガーされていないデータもフローコンピューティングの状態です。機械学習 / ディープラーニングシナリオでは、学習済みモデルとパラメーターデータがストリームコンピューティングの状態となります。
2. グローバル整合性スナップショット
グローバル整合性スナップショットは、分散システムのバックアップと障害復旧に使用できるメカニズムです。
グローバルスナップショット
グローバルスナップショットとは
グローバルスナップショットは、まず複数のサーバーに分散した複数のプロセスを持つ分散アプリケーションです。次に、アプリケーション内部に独自の処理ロジックと状態を持っています。3 つ目に、アプリケーション間で相互に通信できます。4 つ目に、この種の分散アプリケーションには内部状態があり、ハードウェア間で通信可能な場合、ある時点のグローバル状態をグローバルスナップショットと呼びます。
グローバルスナップショットが必要な理由
1 つ目に、チェックポイントとして使用し、グローバル状態を定期的にバックアップでき、アプリケーションに障害が発生した際に復元に使用できます。
2 つ目に、デッドロック検出を行います。スナップショットを取得した後、現在のプログラムを継続して実行し、スナップショットを分析してアプリケーションがデッドロック状態にあるかどうかを確認できます。デッドロック状態であれば、それに応じて対処できます。
グローバルスナップショットの例
下図は分散システムにおけるグローバルスナップショットの例です。
P1 と P2 は 2 つのプロセスであり、その間でメッセージを送信するためのパイプラインがあり、それぞれ C12 と C21 と呼ばれます。P1 プロセスにとって、C12 はメッセージを送信するチャンネルで、出力チャンネルと呼ばれます。C21 はメッセージを受信するチャンネルで、入力チャンネルと呼ばれます。
各プロセスはパイプラインの他にローカル状態を持っています。たとえば、P1 と P2 の各プロセスのメモリには 3 つの変数 XYZ と対応する値があります。P1 と P2 のプロセスのローカル状態と、その間でメッセージを送信するパイプライン状態が初期グローバル状態を形成し、これをグローバルスナップショットと呼ぶこともできます。
P1 が P2 にメッセージを送信し、x の状態値を 4 から 7 に変更するよう要求したとします。ただし、このメッセージはパイプライン内にあり、まだ P2 に到達していません。この状態もグローバルスナップショットです。
次に、P2 が P1 からのメッセージを受信しましたが、まだ処理していない状態です。この状態もグローバルスナップショットです。
最後に、メッセージを受信した P2 が X のローカル値を 4 から 7 に変更します。これもグローバルスナップショットです。
つまり、イベントが発生するとグローバル状態が変化します。イベントには、プロセスによるメッセージの送信、プロセスによるメッセージの受信、およびプロセスによる自身の状態の変更が含まれます。
2. グローバル整合性スナップショット
2 つのイベント a と b があり、絶対時間で a が b より前に発生し、b がスナップショットに含まれている場合、a もそのスナップショットに含まれます。この条件を満たすグローバルスナップショットをグローバル整合性スナップショットと呼びます。
2.1 グローバル整合性スナップショットの実装方法
時計の同期ではグローバル整合性スナップショットを実現できません。グローバル同期は可能ですが、その欠点も非常に明白です。すべてのアプリケーションを停止させ、グローバルパフォーマンスに影響を与えます。
3. 非同期グローバル整合性スナップショットアルゴリズム -- Chandy-Lamport
非同期グローバル整合性スナップショットアルゴリズム Chandy-Lamport は、アプリケーションの実行に影響を与えずにグローバル整合性スナップショットを実現できます。
Chandy-Lamport のシステム必要条件は以下のとおりです。
1 つ目に、アプリケーションの動作に影響を与えないこと。つまり、メッセージの送受信に影響を与えず、アプリケーションを停止する必要がありません。
2 つ目に、各プロセスがローカル状態を記録できること。
3 つ目に、記録された状態を分散方式で収集できること。
4 つ目に、任意のプロセスがスナップショットを開始できること。
同時に、Chandy-Lamport アルゴリズムの実行には前提条件があります。メッセージが順序付けられており、重複せず、メッセージの信頼性が保証されていることです。
3.1 Chandy-Lamport アルゴリズムの流れ
Chandy-Lamport のアルゴリズムフローは、主に開始スナップショット、分散実行スナップショット、終了スナップショットの 3 つの部分に分かれます。
スナップショットの開始
任意のプロセスがスナップショットを開始できます。下図に示すように、P1 がスナップショットを開始するとき、最初のステップはローカル状態を記録すること、つまりローカルのスナップショットを取得することです。その後、時間をおかずにすべての出力チャンネルにマーカーメッセージを即座に送信します。マーカーメッセージは特別なメッセージであり、アプリケーション間で渡されるメッセージとは異なります。
マーカーメッセージの送信後、P1 はすべての入力チャンネルメッセージの記録を開始します。これが図に示されている C21 パイプラインのメッセージです。
分散実行スナップショット
下図に示すように、Pi が Cki からマーカーメッセージを受信したと仮定します。これは Pk が Pi に送信したマーカーメッセージです。2 つの場合に分けられます。
ケース 1:これが Pi が他のパイプラインから受信した最初のマーカーメッセージです。まずローカル状態を記録し、C12 パイプラインを空として記録します。つまり、これ以降 P1 からメッセージが送信されても、このスナップショットには含まれません。同時に、すべての出力チャンネルにマーカーメッセージを即座に送信します。最後に、Cki を除くすべての入力チャンネルからのメッセージの記録を開始します。
上記で Cki のメッセージはリアルタイムスナップショットに含まれないと述べましたが、リアルタイムメッセージは引き続き発生します。そのため、2 つ目のケースとして、Pi が既にマーカーメッセージを受信していた場合、Cki メッセージの記録を停止し、それまでに記録されたすべての Cki メッセージをこのスナップショットにおける Cki の最終状態として保存します。
スナップショットの終了
スナップショットの終了条件は 2 つあります。
1 つ目に、すべてのプロセスがマーカーメッセージを受信し、ローカルスナップショットに記録したこと。
2 つ目に、すべてのプロセスが自身の n-1 個の入力チャンネルからマーカーメッセージを受信し、パイプライン状態を記録したこと。
スナップショットが終了すると、スナップショットコレクター (Central Server) が各部分のスナップショットを収集し、グローバル整合性スナップショットを形成します。
例の紹介
下図の例では、A のように他のプロセスとやり取りしない内部状態が発生しています。内部状態は P1 が自身に送信したメッセージであり、A は C11=[A->] と考えることができます。
Chandy-Lamport グローバル整合性スナップショットアルゴリズムはどのように実装されるのでしょうか。
p1 からスナップショットを開始すると仮定します。スナップショットを開始する際、まずローカル状態のスナップショットを取得します。これを S1 と呼びます。その後、すべての出力チャンネル、つまり P2 と P3 にマーカーメッセージを即座に送信し、すべての入力チャンネルのメッセージを記録します。つまり、P2 と P3 および自身からのメッセージです。
凡例に示すように、縦軸は絶対時間です。絶対時間に基づいて、P3 と P2 がマーカーメッセージを受信する際に時間差が生じるのはなぜでしょうか。実際の物理環境での分散プロセスの場合、異なるノード間のネットワーク状況が異なるため、メッセージの配信時間に差が生じるからです。
P3 が最初にマーカーメッセージを受信し、これは P3 が受信する最初のマーカーメッセージです。メッセージを受信した後、まずローカル状態のスナップショットを取得し、C13 パイプラインをクローズとしてマークし、同時にすべての出力チャンネルにマーカーメッセージの送信を開始し、最後に C13 を除くすべての入力チャンネルからのメッセージの記録を開始します。
P3 が送信したマーカーメッセージを受信するのは P1 ですが、これは P1 が受信する最初のマーカーではありません。C31 チャンネルからのパイプラインを即座にクローズし、現在の記録メッセージをこのチャンネルのスナップショットとします。これ以降 P3 から受信するメッセージは、このスナップショット状態では更新されません。
次に P2 は P3 からのメッセージを受信しますが、これは P2 が受信する最初のマーカーメッセージです。メッセージを受信した後、まずローカル状態のスナップショットを取得し、C32 パイプラインをクローズとしてマークし、同時にすべての出力チャンネルにマーカーメッセージの送信を開始し、最後に C32 を除くすべての入力チャンネルからのメッセージの記録を開始します。
P2 が P1 から受信したメッセージを見てみましょう。これは P2 が受信する最初のマーカーメッセージではないため、すべての入力チャンネルをクローズし、チャンネルの状態を記録します。
次に P1 が P2 からメッセージを受信しますが、これも P1 が受信する最初のメッセージではありません。すべての入力チャンネルをクローズし、記録されたメッセージを状態として使用します。その中には 2 つの状態があります。1 つは C11 で、自身から送信されたメッセージです。もう 1 つは C21 で、P2 の H から P1 に送信された D です。
最後の時点で、P3 は P2 からメッセージを受信しますが、これは P3 が受信する最初のメッセージではなく、操作は上記と同じです。この間に P3 はローカルでイベント J を持ち、J を自身の状態として記録します。
すべてのプロセスがローカル状態を記録し、各プロセスのすべての入力パイプがクローズされると、グローバル整合性スナップショットが完了します。つまり、過去の特定時点のグローバル状態の記録が完了します。
3.3 Chandy-Lamport と Flink の関係
Flink は分散システムであるため、Flink はグローバル整合性スナップショットを使用してチェックポイントを形成し、障害復旧をサポートします。Flink の非同期グローバル整合性スナップショットアルゴリズムと Chandy-Lamport アルゴリズムの主な違いは以下のとおりです。
1 つ目に、Chandy-Lamport は強連結グラフをサポートしますが、Flink は弱連結グラフをサポートします。
2 つ目に、Flink はカスタマイズされた (Tailored) Chandy-Lamport 非同期スナップショットアルゴリズムを使用します。
3 つ目に、Flink の非同期スナップショットアルゴリズムは DAG シナリオでチャンネル状態を保存する必要がないため、スナップショットのストレージ容量を大幅に削減できます。
3. Flink のフォールトトレランスメカニズム
フォールトトレランスとは、エラー発生前の状態に復元することです。ストリームコンピューティングのフォールトトレランス整合性保証には、Exactly once、At least once、At most once の 3 種類があります。
Exactly once は、各イベントが状態に一度だけ影響を与えることを意味します。ここでの「1 回」は厳密なエンドツーエンドの 1 回ではなく、Flink 内部で一度だけ処理されることを意味し、source と sink の処理は含みません。
At least once は、各イベントが状態に少なくとも一度影響を与えることを意味します。つまり、重複処理の可能性があります。
At most once は、各イベントが状態に最大で一度だけ影響を与えることを意味します。つまり、エラー発生時に状態が失われる可能性があります。
エンドツーエンド Exactly once
Exactly once は、ジョブの結果が常に正しいことを意味しますが、複数回生成される可能性があります。そのため、リプレイ可能なソースが必要です。
エンドツーエンド Exactly once は、ジョブの結果が正しく、一度だけ生成されることを意味します。リプレイ可能なソースだけでなく、トランザクション対応シンクと冪等な出力結果も必要です。
Flink の状態フォールトトレランス
多くのシナリオで Exactly once のセマンティクス、つまり一度だけ処理することが求められます。セマンティクスをどのように保証するのでしょうか。
単純なシナリオのための Exactly once フォールトトレランスアプローチ
単純なシナリオのための方法は下図に示す通りです。方法は、ローカル状態を記録し、ソースのオフセット、つまりイベントログの位置を記録することです。
分散シナリオにおける状態フォールトトレランス
分散シナリオでは、ローカル状態を持つ複数のオペレーターに対して、動作を中断せずにグローバル整合性スナップショットを生成する必要があります。Flink 分散シナリオのジョブトポロジーはかなり特殊です。有向非巡回で弱連結なグラフです。カスタマイズされた Chandy-Lamport を使用でき、すべての入力オフセットと各オペレーターの状態のみを記録し、巻き戻し可能なソース (バックトラッキング可能なソース、つまりオフセットを通じて以前の時点を読み取れるソース) に依存するため、チャンネル状態を保存する必要がなく、集約ロジックがある場合に大量のストレージ容量を節約できます。
最後に、復元です。復元とは、データソースの位置をリセットし、チェックポイントから各オペレーターの状態を復元することです。
3. Flink の分散スナップショット方式
まず、ソースデータストリームにチェックポイントバリアを挿入します。これは前述の Chandy-Lamport アルゴリズムにおけるマーカーメッセージです。異なるチェックポイントバリアはストリームを自然に複数のセグメントに分割し、各セグメントにはチェックポイントデータが含まれます。
Flink にはグローバルコーディネーターがあります。任意のプロセスがスナップショットを開始できる Chandy-Lamport とは異なり、この集中型コーディネーターが各ソースにチェックポイントバリアを注入し、スナップショットを開始します。各ノードがバリアを受信した後、Flink ではチャンネル状態を保存しないため、ローカル状態のみを保存すればよいです。
チェックポイント完了後、各オペレーターの各同時実行数はコーディネーターに確認メッセージを送信します。すべてのタスクの確認メッセージがチェックポイントコーディネーターに受信されると、スナップショットが終了します。
4. プロセスデモ
下図に示すように、チェックポイント N がソースに注入されたと仮定すると、ソースはまず処理中のパーティションのオフセットを記録します。
時間の経過とともに、チェックポイントバリアが 2 つの下流の同時実行数に送信されます。バリアが 2 つの同時実行数に到達すると、それらはそれぞれローカル状態をチェックポイントに記録します。
最後に、バリアが最終サブタスクに到達し、スナップショットが完了します。
これは比較的シンプルなシーンのデモです。各オペレーターは単一ストリーム入力のみを持っています。下図のより複雑なシーンを見てみましょう。オペレーターが複数の入力ストリームを持つ場合です。
オペレーターが複数の入力を持つ場合、バリアのアライメントが必要です。バリアのアライメントはどのように行うのでしょうか。下図に示すように、左側の初期状態で、一方のバリアが到達したとき、もう一方のバリアはまだパイプラインにあり到達していません。既に到達した方のストリームを直接ブロックし、もう一方のストリームのデータ処理を待ちます。もう一方のストリームが到達すると、前のストリームのブロックが解除され、バリアがオペレーターに送信されます。
このプロセスで、ストリームの 1 つをブロックする効果は、バックプレッシャーの発生を引き起こすことです。バリアアライメントによりバックプレッシャーが発生し、オペレーターのデータ処理が一時停止します。
アライメントプロセス中に既にバリアを受信したデータパイプラインがブロックされず、データが流入し続ける場合、次のチェックポイントに属するデータが現在のチェックポイントに含まれます。障害が発生してソースが巻き戻されると、一部のデータが失われ、重複処理が発生します。これが At least once です。At least once を受け入れ可能な場合、バリアアライメントを回避できる他の副作用を選択できます。さらに、非同期スナップショットを使用してタスクの一時停止を最小限に抑え、複数のチェックポイントを同時にサポートすることも可能です。
5. スナップショットのトリガー
ローカルスナップショットをシステムに同期的にアップロードするには、状態のコピーオンライトメカニズムが必要です。
メタデータ情報のスナップショットを取得した後にデータ処理を再開する場合、アップロードプロセス中に復元されたアプリケーションロジックがアップロード中のデータを変更しないことをどのように保証するのでしょうか。実際、異なる状態ストレージバックエンドの処理は異なります。Heap バックエンドはデータのコピーオンライトをトリガーし、RocksDB バックエンドの場合、LSM の特性によりスナップショット取得後のデータが変更されないことが保証されます。
4. Flink の状態管理
1. Flink の状態管理
まず、状態を定義する必要があります。以下の例では、まず ValueState を定義します。
状態を定義する際、以下の情報を指定する必要があります。
*状態識別 ID
*状態のデータ型
*ローカル状態バックエンドでの状態の登録
*ローカル状態バックエンドでの状態の読み取りと書き込み
2. Flink の状態バックエンド
状態バックエンドとも呼ばれ、Flink の状態バックエンドには 2 種類あります。
1 つ目のタイプ、JVM Heap は、データが Java オブジェクトの形式で存在し、読み取りと書き込みもオブジェクトの形式で行われるため、速度が非常に速いです。ただし、2 つの欠点もあります。1 つ目の欠点は、オブジェクトストレージに必要なスペースがディスク上のシリアル化されて圧縮されたデータのサイズの何倍にもなり、大量のメモリスペースを消費することです。2 つ目の欠点は、読み取りと書き込み時にシリアル化は不要ですが、スナップショット形成時にシリアル化が必要で、非同期スナップショット処理が遅くなることです。
2 つ目のタイプ、RocksDB は、読み取りと書き込み時にシリアル化が必要なため、読み取りと書き込みの速度が比較的遅いです。ただし、利点が 1 つあります。LSM ベースのデータ構造はスナップショット形成後に SST ファイルを形成します。非同期チェックポイント処理はファイルコピーの処理であり、CPU 消費が比較的低くなります。
ストリームコンピューティング
ストリームコンピューティングとは、メッセージを継続的に送信できるデータソースがあり、同時にコードを実行する常駐プログラムが存在することを意味します。データソースからメッセージを受信すると、それを処理して結果を下流に出力します。
分散ストリームコンピューティング
分散ストリームコンピューティングとは、入力ストリームを何らかの方法で分割し、複数の分散インスタンスを使用してストリームを処理することを指します。
ストリームコンピューティングにおける状態
コンピューティングはステートフルとステートレスに分けることができます。ステートレスコンピューティングは単一のイベントを処理するだけで済みますが、ステートフルコンピューティングは複数のイベントを記録して処理する必要があります。
簡単な例を挙げます。たとえば、イベントがイベント ID とイベント値の 2 つの部分で構成されているとします。処理ロジックが、イベントを取得するたびにそのイベント値を解析して出力するだけであれば、これはステートレスコンピューティングです。一方、イベントを取得するたびに、その値を解析した後に前のイベント値と比較し、前のイベント値より大きい場合にのみ出力する必要がある場合、これはステートフルコンピューティングです。
ストリームコンピューティングには多くの状態があります。たとえば、重複排除シナリオではすべての主キーが記録されます。また、ウィンドウ計算では、ウィンドウに入りまだトリガーされていないデータもフローコンピューティングの状態です。機械学習 / ディープラーニングシナリオでは、学習済みモデルとパラメーターデータがストリームコンピューティングの状態となります。
2. グローバル整合性スナップショット
グローバル整合性スナップショットは、分散システムのバックアップと障害復旧に使用できるメカニズムです。
グローバルスナップショット
グローバルスナップショットとは
グローバルスナップショットは、まず複数のサーバーに分散した複数のプロセスを持つ分散アプリケーションです。次に、アプリケーション内部に独自の処理ロジックと状態を持っています。3 つ目に、アプリケーション間で相互に通信できます。4 つ目に、この種の分散アプリケーションには内部状態があり、ハードウェア間で通信可能な場合、ある時点のグローバル状態をグローバルスナップショットと呼びます。
グローバルスナップショットが必要な理由
1 つ目に、チェックポイントとして使用し、グローバル状態を定期的にバックアップでき、アプリケーションに障害が発生した際に復元に使用できます。
2 つ目に、デッドロック検出を行います。スナップショットを取得した後、現在のプログラムを継続して実行し、スナップショットを分析してアプリケーションがデッドロック状態にあるかどうかを確認できます。デッドロック状態であれば、それに応じて対処できます。
グローバルスナップショットの例
下図は分散システムにおけるグローバルスナップショットの例です。
P1 と P2 は 2 つのプロセスであり、その間でメッセージを送信するためのパイプラインがあり、それぞれ C12 と C21 と呼ばれます。P1 プロセスにとって、C12 はメッセージを送信するチャンネルで、出力チャンネルと呼ばれます。C21 はメッセージを受信するチャンネルで、入力チャンネルと呼ばれます。
各プロセスはパイプラインの他にローカル状態を持っています。たとえば、P1 と P2 の各プロセスのメモリには 3 つの変数 XYZ と対応する値があります。P1 と P2 のプロセスのローカル状態と、その間でメッセージを送信するパイプライン状態が初期グローバル状態を形成し、これをグローバルスナップショットと呼ぶこともできます。
P1 が P2 にメッセージを送信し、x の状態値を 4 から 7 に変更するよう要求したとします。ただし、このメッセージはパイプライン内にあり、まだ P2 に到達していません。この状態もグローバルスナップショットです。
次に、P2 が P1 からのメッセージを受信しましたが、まだ処理していない状態です。この状態もグローバルスナップショットです。
最後に、メッセージを受信した P2 が X のローカル値を 4 から 7 に変更します。これもグローバルスナップショットです。
つまり、イベントが発生するとグローバル状態が変化します。イベントには、プロセスによるメッセージの送信、プロセスによるメッセージの受信、およびプロセスによる自身の状態の変更が含まれます。
2. グローバル整合性スナップショット
2 つのイベント a と b があり、絶対時間で a が b より前に発生し、b がスナップショットに含まれている場合、a もそのスナップショットに含まれます。この条件を満たすグローバルスナップショットをグローバル整合性スナップショットと呼びます。
2.1 グローバル整合性スナップショットの実装方法
時計の同期ではグローバル整合性スナップショットを実現できません。グローバル同期は可能ですが、その欠点も非常に明白です。すべてのアプリケーションを停止させ、グローバルパフォーマンスに影響を与えます。
3. 非同期グローバル整合性スナップショットアルゴリズム -- Chandy-Lamport
非同期グローバル整合性スナップショットアルゴリズム Chandy-Lamport は、アプリケーションの実行に影響を与えずにグローバル整合性スナップショットを実現できます。
Chandy-Lamport のシステム必要条件は以下のとおりです。
1 つ目に、アプリケーションの動作に影響を与えないこと。つまり、メッセージの送受信に影響を与えず、アプリケーションを停止する必要がありません。
2 つ目に、各プロセスがローカル状態を記録できること。
3 つ目に、記録された状態を分散方式で収集できること。
4 つ目に、任意のプロセスがスナップショットを開始できること。
同時に、Chandy-Lamport アルゴリズムの実行には前提条件があります。メッセージが順序付けられており、重複せず、メッセージの信頼性が保証されていることです。
3.1 Chandy-Lamport アルゴリズムの流れ
Chandy-Lamport のアルゴリズムフローは、主に開始スナップショット、分散実行スナップショット、終了スナップショットの 3 つの部分に分かれます。
スナップショットの開始
任意のプロセスがスナップショットを開始できます。下図に示すように、P1 がスナップショットを開始するとき、最初のステップはローカル状態を記録すること、つまりローカルのスナップショットを取得することです。その後、時間をおかずにすべての出力チャンネルにマーカーメッセージを即座に送信します。マーカーメッセージは特別なメッセージであり、アプリケーション間で渡されるメッセージとは異なります。
マーカーメッセージの送信後、P1 はすべての入力チャンネルメッセージの記録を開始します。これが図に示されている C21 パイプラインのメッセージです。
分散実行スナップショット
下図に示すように、Pi が Cki からマーカーメッセージを受信したと仮定します。これは Pk が Pi に送信したマーカーメッセージです。2 つの場合に分けられます。
ケース 1:これが Pi が他のパイプラインから受信した最初のマーカーメッセージです。まずローカル状態を記録し、C12 パイプラインを空として記録します。つまり、これ以降 P1 からメッセージが送信されても、このスナップショットには含まれません。同時に、すべての出力チャンネルにマーカーメッセージを即座に送信します。最後に、Cki を除くすべての入力チャンネルからのメッセージの記録を開始します。
上記で Cki のメッセージはリアルタイムスナップショットに含まれないと述べましたが、リアルタイムメッセージは引き続き発生します。そのため、2 つ目のケースとして、Pi が既にマーカーメッセージを受信していた場合、Cki メッセージの記録を停止し、それまでに記録されたすべての Cki メッセージをこのスナップショットにおける Cki の最終状態として保存します。
スナップショットの終了
スナップショットの終了条件は 2 つあります。
1 つ目に、すべてのプロセスがマーカーメッセージを受信し、ローカルスナップショットに記録したこと。
2 つ目に、すべてのプロセスが自身の n-1 個の入力チャンネルからマーカーメッセージを受信し、パイプライン状態を記録したこと。
スナップショットが終了すると、スナップショットコレクター (Central Server) が各部分のスナップショットを収集し、グローバル整合性スナップショットを形成します。
例の紹介
下図の例では、A のように他のプロセスとやり取りしない内部状態が発生しています。内部状態は P1 が自身に送信したメッセージであり、A は C11=[A->] と考えることができます。
Chandy-Lamport グローバル整合性スナップショットアルゴリズムはどのように実装されるのでしょうか。
p1 からスナップショットを開始すると仮定します。スナップショットを開始する際、まずローカル状態のスナップショットを取得します。これを S1 と呼びます。その後、すべての出力チャンネル、つまり P2 と P3 にマーカーメッセージを即座に送信し、すべての入力チャンネルのメッセージを記録します。つまり、P2 と P3 および自身からのメッセージです。
凡例に示すように、縦軸は絶対時間です。絶対時間に基づいて、P3 と P2 がマーカーメッセージを受信する際に時間差が生じるのはなぜでしょうか。実際の物理環境での分散プロセスの場合、異なるノード間のネットワーク状況が異なるため、メッセージの配信時間に差が生じるからです。
P3 が最初にマーカーメッセージを受信し、これは P3 が受信する最初のマーカーメッセージです。メッセージを受信した後、まずローカル状態のスナップショットを取得し、C13 パイプラインをクローズとしてマークし、同時にすべての出力チャンネルにマーカーメッセージの送信を開始し、最後に C13 を除くすべての入力チャンネルからのメッセージの記録を開始します。
P3 が送信したマーカーメッセージを受信するのは P1 ですが、これは P1 が受信する最初のマーカーではありません。C31 チャンネルからのパイプラインを即座にクローズし、現在の記録メッセージをこのチャンネルのスナップショットとします。これ以降 P3 から受信するメッセージは、このスナップショット状態では更新されません。
次に P2 は P3 からのメッセージを受信しますが、これは P2 が受信する最初のマーカーメッセージです。メッセージを受信した後、まずローカル状態のスナップショットを取得し、C32 パイプラインをクローズとしてマークし、同時にすべての出力チャンネルにマーカーメッセージの送信を開始し、最後に C32 を除くすべての入力チャンネルからのメッセージの記録を開始します。
P2 が P1 から受信したメッセージを見てみましょう。これは P2 が受信する最初のマーカーメッセージではないため、すべての入力チャンネルをクローズし、チャンネルの状態を記録します。
次に P1 が P2 からメッセージを受信しますが、これも P1 が受信する最初のメッセージではありません。すべての入力チャンネルをクローズし、記録されたメッセージを状態として使用します。その中には 2 つの状態があります。1 つは C11 で、自身から送信されたメッセージです。もう 1 つは C21 で、P2 の H から P1 に送信された D です。
最後の時点で、P3 は P2 からメッセージを受信しますが、これは P3 が受信する最初のメッセージではなく、操作は上記と同じです。この間に P3 はローカルでイベント J を持ち、J を自身の状態として記録します。
すべてのプロセスがローカル状態を記録し、各プロセスのすべての入力パイプがクローズされると、グローバル整合性スナップショットが完了します。つまり、過去の特定時点のグローバル状態の記録が完了します。
3.3 Chandy-Lamport と Flink の関係
Flink は分散システムであるため、Flink はグローバル整合性スナップショットを使用してチェックポイントを形成し、障害復旧をサポートします。Flink の非同期グローバル整合性スナップショットアルゴリズムと Chandy-Lamport アルゴリズムの主な違いは以下のとおりです。
1 つ目に、Chandy-Lamport は強連結グラフをサポートしますが、Flink は弱連結グラフをサポートします。
2 つ目に、Flink はカスタマイズされた (Tailored) Chandy-Lamport 非同期スナップショットアルゴリズムを使用します。
3 つ目に、Flink の非同期スナップショットアルゴリズムは DAG シナリオでチャンネル状態を保存する必要がないため、スナップショットのストレージ容量を大幅に削減できます。
3. Flink のフォールトトレランスメカニズム
フォールトトレランスとは、エラー発生前の状態に復元することです。ストリームコンピューティングのフォールトトレランス整合性保証には、Exactly once、At least once、At most once の 3 種類があります。
Exactly once は、各イベントが状態に一度だけ影響を与えることを意味します。ここでの「1 回」は厳密なエンドツーエンドの 1 回ではなく、Flink 内部で一度だけ処理されることを意味し、source と sink の処理は含みません。
At least once は、各イベントが状態に少なくとも一度影響を与えることを意味します。つまり、重複処理の可能性があります。
At most once は、各イベントが状態に最大で一度だけ影響を与えることを意味します。つまり、エラー発生時に状態が失われる可能性があります。
エンドツーエンド Exactly once
Exactly once は、ジョブの結果が常に正しいことを意味しますが、複数回生成される可能性があります。そのため、リプレイ可能なソースが必要です。
エンドツーエンド Exactly once は、ジョブの結果が正しく、一度だけ生成されることを意味します。リプレイ可能なソースだけでなく、トランザクション対応シンクと冪等な出力結果も必要です。
Flink の状態フォールトトレランス
多くのシナリオで Exactly once のセマンティクス、つまり一度だけ処理することが求められます。セマンティクスをどのように保証するのでしょうか。
単純なシナリオのための Exactly once フォールトトレランスアプローチ
単純なシナリオのための方法は下図に示す通りです。方法は、ローカル状態を記録し、ソースのオフセット、つまりイベントログの位置を記録することです。
分散シナリオにおける状態フォールトトレランス
分散シナリオでは、ローカル状態を持つ複数のオペレーターに対して、動作を中断せずにグローバル整合性スナップショットを生成する必要があります。Flink 分散シナリオのジョブトポロジーはかなり特殊です。有向非巡回で弱連結なグラフです。カスタマイズされた Chandy-Lamport を使用でき、すべての入力オフセットと各オペレーターの状態のみを記録し、巻き戻し可能なソース (バックトラッキング可能なソース、つまりオフセットを通じて以前の時点を読み取れるソース) に依存するため、チャンネル状態を保存する必要がなく、集約ロジックがある場合に大量のストレージ容量を節約できます。
最後に、復元です。復元とは、データソースの位置をリセットし、チェックポイントから各オペレーターの状態を復元することです。
3. Flink の分散スナップショット方式
まず、ソースデータストリームにチェックポイントバリアを挿入します。これは前述の Chandy-Lamport アルゴリズムにおけるマーカーメッセージです。異なるチェックポイントバリアはストリームを自然に複数のセグメントに分割し、各セグメントにはチェックポイントデータが含まれます。
Flink にはグローバルコーディネーターがあります。任意のプロセスがスナップショットを開始できる Chandy-Lamport とは異なり、この集中型コーディネーターが各ソースにチェックポイントバリアを注入し、スナップショットを開始します。各ノードがバリアを受信した後、Flink ではチャンネル状態を保存しないため、ローカル状態のみを保存すればよいです。
チェックポイント完了後、各オペレーターの各同時実行数はコーディネーターに確認メッセージを送信します。すべてのタスクの確認メッセージがチェックポイントコーディネーターに受信されると、スナップショットが終了します。
4. プロセスデモ
下図に示すように、チェックポイント N がソースに注入されたと仮定すると、ソースはまず処理中のパーティションのオフセットを記録します。
時間の経過とともに、チェックポイントバリアが 2 つの下流の同時実行数に送信されます。バリアが 2 つの同時実行数に到達すると、それらはそれぞれローカル状態をチェックポイントに記録します。
最後に、バリアが最終サブタスクに到達し、スナップショットが完了します。
これは比較的シンプルなシーンのデモです。各オペレーターは単一ストリーム入力のみを持っています。下図のより複雑なシーンを見てみましょう。オペレーターが複数の入力ストリームを持つ場合です。
オペレーターが複数の入力を持つ場合、バリアのアライメントが必要です。バリアのアライメントはどのように行うのでしょうか。下図に示すように、左側の初期状態で、一方のバリアが到達したとき、もう一方のバリアはまだパイプラインにあり到達していません。既に到達した方のストリームを直接ブロックし、もう一方のストリームのデータ処理を待ちます。もう一方のストリームが到達すると、前のストリームのブロックが解除され、バリアがオペレーターに送信されます。
このプロセスで、ストリームの 1 つをブロックする効果は、バックプレッシャーの発生を引き起こすことです。バリアアライメントによりバックプレッシャーが発生し、オペレーターのデータ処理が一時停止します。
アライメントプロセス中に既にバリアを受信したデータパイプラインがブロックされず、データが流入し続ける場合、次のチェックポイントに属するデータが現在のチェックポイントに含まれます。障害が発生してソースが巻き戻されると、一部のデータが失われ、重複処理が発生します。これが At least once です。At least once を受け入れ可能な場合、バリアアライメントを回避できる他の副作用を選択できます。さらに、非同期スナップショットを使用してタスクの一時停止を最小限に抑え、複数のチェックポイントを同時にサポートすることも可能です。
5. スナップショットのトリガー
ローカルスナップショットをシステムに同期的にアップロードするには、状態のコピーオンライトメカニズムが必要です。
メタデータ情報のスナップショットを取得した後にデータ処理を再開する場合、アップロードプロセス中に復元されたアプリケーションロジックがアップロード中のデータを変更しないことをどのように保証するのでしょうか。実際、異なる状態ストレージバックエンドの処理は異なります。Heap バックエンドはデータのコピーオンライトをトリガーし、RocksDB バックエンドの場合、LSM の特性によりスナップショット取得後のデータが変更されないことが保証されます。
4. Flink の状態管理
1. Flink の状態管理
まず、状態を定義する必要があります。以下の例では、まず ValueState を定義します。
状態を定義する際、以下の情報を指定する必要があります。
*状態識別 ID
*状態のデータ型
*ローカル状態バックエンドでの状態の登録
*ローカル状態バックエンドでの状態の読み取りと書き込み
2. Flink の状態バックエンド
状態バックエンドとも呼ばれ、Flink の状態バックエンドには 2 種類あります。
1 つ目のタイプ、JVM Heap は、データが Java オブジェクトの形式で存在し、読み取りと書き込みもオブジェクトの形式で行われるため、速度が非常に速いです。ただし、2 つの欠点もあります。1 つ目の欠点は、オブジェクトストレージに必要なスペースがディスク上のシリアル化されて圧縮されたデータのサイズの何倍にもなり、大量のメモリスペースを消費することです。2 つ目の欠点は、読み取りと書き込み時にシリアル化は不要ですが、スナップショット形成時にシリアル化が必要で、非同期スナップショット処理が遅くなることです。
2 つ目のタイプ、RocksDB は、読み取りと書き込み時にシリアル化が必要なため、読み取りと書き込みの速度が比較的遅いです。ただし、利点が 1 つあります。LSM ベースのデータ構造はスナップショット形成後に SST ファイルを形成します。非同期チェックポイント処理はファイルコピーの処理であり、CPU 消費が比較的低くなります。
Related Articles
-
A detailed explanation of Hadoop core architecture HDFS
Knowledge Base Team
-
What Does IOT Mean
Knowledge Base Team
-
6 Optional Technologies for Data Storage
Knowledge Base Team
-
What Is Blockchain Technology
Knowledge Base Team
Explore More Special Offers
-
Short Message Service(SMS) & Mail Service
50,000 email package starts as low as USD 1.99, 120 short messages start at only USD 1.00
