Detailed explanation of Flink Connector
一、Connector の概要:Connector — Flink と外部システムをつなぐブリッジ
1. Connector コネクタ
Flink のデータの重要な入力元と出力先
コネクタは Flink と外部システムをつなぐブリッジです。たとえば、Kafka からデータを読み取り、Flink でデータを処理した後、Hive や Elasticsearch などの外部システムに書き戻す必要があります。
処理フローにおけるイベントコントロール:イベント処理、ウォーターマーク、チェックポイント整合、レコード追跡
負荷分散:異なる同時実行負荷に応じてデータパーティションを適切に割り当てる
データの解析とシリアル化:外部システムではデータがバイナリ形式で保存されていたり、データベース内のさまざまな列の形式で保存されていたりする場合があります。Flink に読み込んだ後は、後続のデータ処理を実行するために解析が必要です。同様に、外部システムに書き戻す際もシリアル化操作を実行し、外部システムの対応するストレージ形式に変換して保存します。
上図は非常に典型的な例を示しています。
まず、ソースを通じて Kafka から一部のレコードを読み取ります。次に、これらのレコードを Flink 内のオペレーターに送信して対応する処理を実行し、シンクを通じて Elasticsearch に書き出します。このように、ソースとシンクは Flink ジョブの両端でインターフェイスとして機能します。
二、Source API — Flink データの入力
1. Source インターフェイスの進化
Flink 1.10 より前は、Source は左側の 2 つのインターフェイスでした。SourceFunction API(ストリームデータ処理用)と InputFormat API(バッチデータ処理用)です。Flink 1.10 以降、コミュニティは新しい Source API を導入し、Source 全体をリファクタリングしました。では、なぜコミュニティはこのような変更を行ったのでしょうか。
バッチとストリームの実装の非整合性:生態系が拡大するにつれ、旧 API ではいくつかの問題が顕在化しました。最も直感的な問題は、バッチ処理とストリーム処理の実装が一貫していないことです。
インターフェイスはシンプルだが実装が複雑:以前の API はインターフェイス自体は比較的シンプルでしたが、実際には開発者がこのインターフェイスを実装する際、すべてのロジックと操作が非常に複雑で、十分に使いやすいものではありませんでした。
そこで、これらの問題に基づき、FLIP-27 で新しい Source API の設計が提案されました。これには 2 つの特徴があります。
バッチストリーム統一:ストリームデータ処理とバッチデータ処理で 2 組のコードを維持する必要がなく、1 組のコードで対応できます。
シンプルな実装:Source API は多くの概念的抽象化を定義しています。これらの抽象化は一見複雑に見えますが、実際には開発者の開発作業を簡素化します。
2. コア抽象化
1) レコードスプリット (Split)
番号付きのレコード集合
Kafka を例に取ると、Kafka のスプリットはパーティション全体として定義することも、パーティションの一部として定義することもできます。たとえば、オフセット 100 からデータの消費を開始する場合、100〜200 の範囲を 1 つのスプリット、201〜300 を別のスプリットとして定義できます。レコードの集合であり、一意の番号を付ければ、レコードスプリットとして定義できます。
進行状況の追跡が可能
このスプリット内の現在の処理位置を記録する必要があります。チェックポイントを記録する際、現在どこまで処理したかを知っておく必要があり、障害が発生した際にその位置から直接復旧できます。
スプリットのすべての情報を記録する
Kafka を例に取ると、パーティションの開始点や終了点などの情報がレコードスプリット全体に含まれている必要があります。チェックポイントもレコードスプリット単位で実行されるため、レコードスプリット内の情報に自己矛盾があってはなりません。
2) スプリット列挙子 (Split Enumerator)
レコードスプリットの検出:外部システム内のスプリットの存在を検出する
レコードスプリットの割り当て:列挙子はコーディネーターの役割を担います。ソースリーダーにタスクを割り当てる必要があります。
ソースリーダーの調整:たとえば、一部のリーダーの進行状況が速すぎる場合、ウォーターマークのおおむねの整合性を保つために速度を落とすよう調整します。
3) ソースリーダー (Source Reader)
レコードスプリットからのデータ読み取り:列挙子から割り当てられたレコードスプリットに従ってデータを読み取る
イベント時間のウォーターマーク処理:外部システムから読み取ったデータからイベント時間を抽出し、対応するウォーターマーク送信操作を実行する必要がある
データ解析:外部システムから読み取ったデータを逆シリアル化し、下流オペレーターに送信する
3. Enumerator-Reader アーキテクチャ
スプリット列挙子は Job Master 上で動作し、ソースリーダーは Task Executor 上で動作します。したがって、列挙子はリーダー役、コーディネーター役であり、リーダーは実行役です。
それぞれのチェックポイントストレージも分離されていますが、両者間には通信があります。たとえば、列挙子はリーダーにタスクを割り当て、処理すべきスプリットがもうないことを通知する必要があります。異なる動作環境のため、両者間にはネットワーク通信が必然的に発生します。そこで、以下のような通信スタックの定義があります。
この通信スタック上では、開発者が独自に実装できるよう、いくつかのイベントが定義されています。
まず、最上位層は Source Event で、開発者がカスタマイズ操作を定義するために用意されています。たとえば、ある Source の設計で、特定の条件で読み取りを一時停止する場合、SplitEnumerator からこの Source Event を通じて Source Reader に通知を送信できます。
次に、下層は Operator Coordinator(オペレーターコーディネーター)と呼ばれます。Operator Event(オペレーターイベント)を通じて、実際にタスクを実行するオペレーターと通信します。スプリットの追加や新スプリットがないことの通知など、事前に定義されたオペレーターイベントがあります。これらのすべての Source に共通するイベントは、Operator Event レベルで抽象化されています。
Address Lookup は、メッセージをどのオペレーターに送信すべきかを特定するために使用されます。Flink のジョブ全体が実行されると有向非循環グラフが生成され、異なるオペレーターが異なる Task Manager 上で動作する可能性があるため、対応するタスクとオペレーターを見つけるのがこのレイヤーの役割です。
ネットワーク通信が存在するため、Job Master と Task Executor の間には RPC Gateway があります。すべてのイベントは最終的に RPC Gateway と RPC 呼び出しを通じてネットワーク上で送信されます。
4. Source リーダーの設計
Source Reader の実装手順を簡素化し、開発者の作業を軽減するため、コミュニティは SourceReaderBase を提供しています。ユーザーは開発時に SourceReaderBase クラスを直接継承でき、開発作業を大幅に簡素化できます。次に SourceReaderBase を分析します。この図には多くのコンポーネントが含まれているように見えますが、実際には 2 つの部分に分割して理解できます。
中間の elementQueue キューを境界として、左側の青でマークされた部分は外部システムとのやり取りを担当するコンポーネント、右側のオレンジでマークされた部分は Flink のエンジン側とのやり取りを担当するコンポーネントです。
まず、左側は 1 つ以上のスプリットリーダーで構成されます。各リーダーは Fetcher によって駆動され、複数の Fetcher は Fetcher Manager によって管理されます。ここにも多くの実装方法があります。たとえば、1 つのスレッドと 1 つの SplitReader のみを起動し、この 1 つのリーダーで複数のパーティションを消費できます。また、要件に応じてマルチスレッドを起動し、各スレッドで 1 つの Fetcher と 1 つのリーダーを実行し、各リーダーが 1 つのパーティションを担当して並列でデータを消費することもできます。これらはすべて、ユーザーの実装と選択に委ねられています。
パフォーマンスの観点から、各 SplitReader は外部システムからバッチデータを取得し、elementQueue に格納します。図に示すように、青いボックス内は毎回取得するバッチデータで、オレンジのボックスはこのバッチデータ内の各データピースです。
次に、elementQueue の右側は RecordEmitter と SourceOutput で構成されます。RecordEmitter は各レコードを下流の SourceOutput に送信してレコードを出力します。RecordEmitter は毎回、中間の elementQueue からバッチデータを取得し、1 つずつ下流に送信します。RecordEmitter はメインスレッドによって駆動されるため、現在のメインスレッドの設計はロックフリーのメールボックスモデルを使用しています。このモデルは実行すべき作業をメール単位に分割し、ワーカースレッドが毎回メールボックスからメールを取り出して処理するため、ここの実装はノンブロッキングである必要があることに注意してください。
RecordEmitter は下流にデータを送信するたびに、後続の処理データがあるかどうかを下流に報告します。同時に、SplitStates に現在のスプリット処理の進行状況を記録し、現在の状態と処理位置を記録します。
SplitEnumerator が外部システムで新しいスプリットを検出すると、RPC 経由で addSplits メソッドを呼び出して新しいスプリットをリーダーに追加する必要があります。SplitFetchermanager 側では、事前に選択したスレッドモデルに従って新しいスプリットが割り当てられます(単一スレッドのみの場合、そのスレッドに新しいタスクを割り当て、リーダーに新しいスプリットを読み取らせます。マルチスレッド実装の場合、新しいスレッドとリーダーを作成してスプリットを別々に処理します)。同様に、SplitStates に現在の処理進行状況を記録する必要があります。
5. チェックポイントの作成
次に、新しい Source API でチェックポイントがどのように処理されるかを見ていきます。
まず、左側のコーディネーター、スプリット列挙子です。図に示すように、現在まだ割り当てられていないスプリット(Split_5)があります。中間の矢印は転送中のスプリットです。破線はこのチェックポイントの境界です。2 番目のスプリットはチェックポイントより前に、4 番目のスプリットはチェックポイントより後ろにあり、下部のリーダーは SplitEnumerator に新しいスプリットを要求していることが分かります。次にリーダーを見ると、3 つのリーダーにはそれぞれスプリットが割り当てられ、処理済みで、既に位置情報(Position)を持っています。では、チェックポイント時に列挙子とリーダーが保存すべきものを見ていきましょう。
列挙子:未割り当てのレコードスプリット(Split_5)、割り当て済みだが未チェックポイントのレコードスプリット(Split_4)
リーダー:割り当て済みレコードスプリット(Split_0、1、3)、スプリット割り当て状態(Split_2)
6. Source を実装する 3 つの簡単なステップ
1) Split/SplitState
Split:外部システムのフラグメンテーション
SplitSerializer:Split をシリアル化/逆シリアル化して SourceReader に渡す
SplitState:Split 状態、チェックポイントとリカバリに使用
2) SplitEnumerator
Split の検出とサブスクライブ
EnumState:Enumerator の状態、チェックポイントとリカバリに使用
EnumStateSerializer:EnumState をシリアル化/逆シリアル化
3) SourceReader
SplitReader:外部システムとのデータやり取りのインターフェイス
FetcherManager:スレッドモデルの選択(現在利用可能な実装あり)
RecordEmitter:メッセージタイプの変換とイベント時間の処理
よく考えると、これらのほとんどは実際には外部システムとのやり取りを担当しており、Flink エンジン自体とのやり取りはごくわずかであることが分かります。ユーザーはチェックポイントロックやマルチスレッドの問題などを心配する必要がなくなり、外部システムとの開発とやり取りにより多くの開発エネルギーを集中できます。したがって、新しい Source API はこれらの抽象化を通じて、開発者の開発作業を大幅に簡素化します。
三、Sink API — Flink データのエクスポート
Flink についてある程度理解していれば、exactly-once セマンティクスを実現でき、データは重複も損失もないことが分かります。この「正確に 1 回」を実現するために、Flink は多くの重要な作業を行っており、その 1 つがシンク側での 2 フェーズコミットの実装です。
1. プリコミットフェーズ
プリコミットフェーズでは、分散システムには一般的に「コーディネーター 1 + エグゼキュータ N」のモードがあるため、まずコーディネーターがコミット要求を送信する必要があります。つまり、すべてのエグゼキュータにコミット要求メッセージを送信し、2 フェーズコミット全体を開始します。
エグゼキュータがコミット要求を受信すると、コミットの準備作業を行います。すべての準備作業が完了したら、エグゼキュータはコーディネーターに「次のコミットの準備が完了した」と応答します。コーディネーターがすべてのエグゼキュータから「続行可能」の応答を受け取ると、プリコミットフェーズが終了し、コミット実行フェーズに入ります。
2. コミット実行フェーズ
コーディネーターはコミット決定メッセージをエグゼキュータに送信し、エグゼキュータは準備したコミット関連処理を実際に実行します。完了後、結果をコーディネーターに応答し、コミットが正常に実行されたことを報告します。
第 2 フェーズであるコミット実行フェーズに入ると、すべてのエグゼキュータはこの決定を無条件で実行しなければなりません。つまり、あるコーディネーターがこの段階で問題が発生しても、復旧後にこの決定を必ず実行する必要があります。つまり、一度コミットが決定されたら、エグゼキュータは必ずコミットアクションを実行しなければなりません。
プリコミットフェーズで、エグゼキュータがコミット直前に障害を経験し、正しいコミットアクションを実行できない場合があります。その場合、ネットワーク切断やタイムアウトなどにより、コーディネーターにエラーを応答することがあります。一定期間経過後にコーディネーターが第 3 のエグゼキュータからの応答を受信しなかった場合、コーディネーターは第 2 段階のロールバックアクションをトリガーします。つまり、すべてのエグゼキュータに「このコミット試行は失敗したため、全員が前の状態にロールバックする必要がある」と通知します。そして、エグゼキュータはロールバックアクションを実行して以前の操作を元に戻します。
3. Flink における 2 フェーズコミットの実践
1) プリコミットフェーズ
ファイルシステムのシンクを例に取ります。
ファイルシステムのシンクは、チェックポイント境界を受信した後にプリコミットアクションを実行します(現在のデータをディスク上の一時ファイルに書き込みます)。プリコミットフェーズが完了すると、すべてのオペレーターはコーディネーターに「コミット準備完了」と応答します。
2) コミット実行フェーズ
第 2 フェーズ、コミット実行フェーズが開始されます。JobManager はすべてのオペレーターにコミット指示を送信し、シンクはこれを受信すると最終的なコミットアクションを実行します。
ファイルシステムを例に取ると、前述の通り、プリコミットフェーズでデータは一時ファイルに書き込まれます。実際のコミット実行時に、一時ファイルは事前に定義された正規の名前にリネームされます。このリネーム操作がコミットの実行に相当します。
ここで注意すべきは、一時ファイルの設定は無駄ではなく、後続のロールバックなどの状況に備える重要な役割を果たしている点です。2 フェーズコミットメカニズムを巧みに活用することで、exactly-once セマンティクスを保証しています。
4. Sink モデル
1) Writer:上流からの連続データを中間状態に書き込む、またはプリコミットフェーズでの処理を担当します。
2) Committable:前述の「中間状態」は、コミット操作を実行可能なコンポーネントです。
3) Committer:Committable を実際にコミットします
4) Global Committer:グローバルコミッター。このコンポーネントはオプションで、外部システムに依存します。例:Iceberg。
四、将来の開発
新しい Source の改善
Source と Sink はリリースされたばかりで、まだいくつかの問題が残っています。開発者から新しい要件や機能改善の要望があり、現在は比較的安定した状態ですが、継続的な改善が必要です。
既存のコネクタを新しい API に移行
ストリームバッチ統合コネクタの継続的な進展に伴い、すべてのコネクタが新しい API に移行されます。
コネクタテストフレームワーク
コネクタテストフレームワークは、すべてのコネクタに対して一貫した統一テスト基準を提供しようとします。テスト開発者は個別にテストケースを書いたり、さまざまなテスト環境やテストシナリオを考慮したりする必要がなくなります。開発者がブロックを組み立てるように異なるシナリオやユースケースでコードを素早くテストできるようにし、ロジック自体の開発により多くの開発エネルギーを集中できるようにして、開発者のテスト負荷を大幅に削減します。これは Source API と Sink API、および後続フレームワーク開発の一貫した目標でもあります。コネクタ開発をより簡単にして参入障壁を下げ、より多くの開発者に Flink エコシステムへの貢献を促すことです。
1. Connector コネクタ
Flink のデータの重要な入力元と出力先
コネクタは Flink と外部システムをつなぐブリッジです。たとえば、Kafka からデータを読み取り、Flink でデータを処理した後、Hive や Elasticsearch などの外部システムに書き戻す必要があります。
処理フローにおけるイベントコントロール:イベント処理、ウォーターマーク、チェックポイント整合、レコード追跡
負荷分散:異なる同時実行負荷に応じてデータパーティションを適切に割り当てる
データの解析とシリアル化:外部システムではデータがバイナリ形式で保存されていたり、データベース内のさまざまな列の形式で保存されていたりする場合があります。Flink に読み込んだ後は、後続のデータ処理を実行するために解析が必要です。同様に、外部システムに書き戻す際もシリアル化操作を実行し、外部システムの対応するストレージ形式に変換して保存します。
上図は非常に典型的な例を示しています。
まず、ソースを通じて Kafka から一部のレコードを読み取ります。次に、これらのレコードを Flink 内のオペレーターに送信して対応する処理を実行し、シンクを通じて Elasticsearch に書き出します。このように、ソースとシンクは Flink ジョブの両端でインターフェイスとして機能します。
二、Source API — Flink データの入力
1. Source インターフェイスの進化
Flink 1.10 より前は、Source は左側の 2 つのインターフェイスでした。SourceFunction API(ストリームデータ処理用)と InputFormat API(バッチデータ処理用)です。Flink 1.10 以降、コミュニティは新しい Source API を導入し、Source 全体をリファクタリングしました。では、なぜコミュニティはこのような変更を行ったのでしょうか。
バッチとストリームの実装の非整合性:生態系が拡大するにつれ、旧 API ではいくつかの問題が顕在化しました。最も直感的な問題は、バッチ処理とストリーム処理の実装が一貫していないことです。
インターフェイスはシンプルだが実装が複雑:以前の API はインターフェイス自体は比較的シンプルでしたが、実際には開発者がこのインターフェイスを実装する際、すべてのロジックと操作が非常に複雑で、十分に使いやすいものではありませんでした。
そこで、これらの問題に基づき、FLIP-27 で新しい Source API の設計が提案されました。これには 2 つの特徴があります。
バッチストリーム統一:ストリームデータ処理とバッチデータ処理で 2 組のコードを維持する必要がなく、1 組のコードで対応できます。
シンプルな実装:Source API は多くの概念的抽象化を定義しています。これらの抽象化は一見複雑に見えますが、実際には開発者の開発作業を簡素化します。
2. コア抽象化
1) レコードスプリット (Split)
番号付きのレコード集合
Kafka を例に取ると、Kafka のスプリットはパーティション全体として定義することも、パーティションの一部として定義することもできます。たとえば、オフセット 100 からデータの消費を開始する場合、100〜200 の範囲を 1 つのスプリット、201〜300 を別のスプリットとして定義できます。レコードの集合であり、一意の番号を付ければ、レコードスプリットとして定義できます。
進行状況の追跡が可能
このスプリット内の現在の処理位置を記録する必要があります。チェックポイントを記録する際、現在どこまで処理したかを知っておく必要があり、障害が発生した際にその位置から直接復旧できます。
スプリットのすべての情報を記録する
Kafka を例に取ると、パーティションの開始点や終了点などの情報がレコードスプリット全体に含まれている必要があります。チェックポイントもレコードスプリット単位で実行されるため、レコードスプリット内の情報に自己矛盾があってはなりません。
2) スプリット列挙子 (Split Enumerator)
レコードスプリットの検出:外部システム内のスプリットの存在を検出する
レコードスプリットの割り当て:列挙子はコーディネーターの役割を担います。ソースリーダーにタスクを割り当てる必要があります。
ソースリーダーの調整:たとえば、一部のリーダーの進行状況が速すぎる場合、ウォーターマークのおおむねの整合性を保つために速度を落とすよう調整します。
3) ソースリーダー (Source Reader)
レコードスプリットからのデータ読み取り:列挙子から割り当てられたレコードスプリットに従ってデータを読み取る
イベント時間のウォーターマーク処理:外部システムから読み取ったデータからイベント時間を抽出し、対応するウォーターマーク送信操作を実行する必要がある
データ解析:外部システムから読み取ったデータを逆シリアル化し、下流オペレーターに送信する
3. Enumerator-Reader アーキテクチャ
スプリット列挙子は Job Master 上で動作し、ソースリーダーは Task Executor 上で動作します。したがって、列挙子はリーダー役、コーディネーター役であり、リーダーは実行役です。
それぞれのチェックポイントストレージも分離されていますが、両者間には通信があります。たとえば、列挙子はリーダーにタスクを割り当て、処理すべきスプリットがもうないことを通知する必要があります。異なる動作環境のため、両者間にはネットワーク通信が必然的に発生します。そこで、以下のような通信スタックの定義があります。
この通信スタック上では、開発者が独自に実装できるよう、いくつかのイベントが定義されています。
まず、最上位層は Source Event で、開発者がカスタマイズ操作を定義するために用意されています。たとえば、ある Source の設計で、特定の条件で読み取りを一時停止する場合、SplitEnumerator からこの Source Event を通じて Source Reader に通知を送信できます。
次に、下層は Operator Coordinator(オペレーターコーディネーター)と呼ばれます。Operator Event(オペレーターイベント)を通じて、実際にタスクを実行するオペレーターと通信します。スプリットの追加や新スプリットがないことの通知など、事前に定義されたオペレーターイベントがあります。これらのすべての Source に共通するイベントは、Operator Event レベルで抽象化されています。
Address Lookup は、メッセージをどのオペレーターに送信すべきかを特定するために使用されます。Flink のジョブ全体が実行されると有向非循環グラフが生成され、異なるオペレーターが異なる Task Manager 上で動作する可能性があるため、対応するタスクとオペレーターを見つけるのがこのレイヤーの役割です。
ネットワーク通信が存在するため、Job Master と Task Executor の間には RPC Gateway があります。すべてのイベントは最終的に RPC Gateway と RPC 呼び出しを通じてネットワーク上で送信されます。
4. Source リーダーの設計
Source Reader の実装手順を簡素化し、開発者の作業を軽減するため、コミュニティは SourceReaderBase を提供しています。ユーザーは開発時に SourceReaderBase クラスを直接継承でき、開発作業を大幅に簡素化できます。次に SourceReaderBase を分析します。この図には多くのコンポーネントが含まれているように見えますが、実際には 2 つの部分に分割して理解できます。
中間の elementQueue キューを境界として、左側の青でマークされた部分は外部システムとのやり取りを担当するコンポーネント、右側のオレンジでマークされた部分は Flink のエンジン側とのやり取りを担当するコンポーネントです。
まず、左側は 1 つ以上のスプリットリーダーで構成されます。各リーダーは Fetcher によって駆動され、複数の Fetcher は Fetcher Manager によって管理されます。ここにも多くの実装方法があります。たとえば、1 つのスレッドと 1 つの SplitReader のみを起動し、この 1 つのリーダーで複数のパーティションを消費できます。また、要件に応じてマルチスレッドを起動し、各スレッドで 1 つの Fetcher と 1 つのリーダーを実行し、各リーダーが 1 つのパーティションを担当して並列でデータを消費することもできます。これらはすべて、ユーザーの実装と選択に委ねられています。
パフォーマンスの観点から、各 SplitReader は外部システムからバッチデータを取得し、elementQueue に格納します。図に示すように、青いボックス内は毎回取得するバッチデータで、オレンジのボックスはこのバッチデータ内の各データピースです。
次に、elementQueue の右側は RecordEmitter と SourceOutput で構成されます。RecordEmitter は各レコードを下流の SourceOutput に送信してレコードを出力します。RecordEmitter は毎回、中間の elementQueue からバッチデータを取得し、1 つずつ下流に送信します。RecordEmitter はメインスレッドによって駆動されるため、現在のメインスレッドの設計はロックフリーのメールボックスモデルを使用しています。このモデルは実行すべき作業をメール単位に分割し、ワーカースレッドが毎回メールボックスからメールを取り出して処理するため、ここの実装はノンブロッキングである必要があることに注意してください。
RecordEmitter は下流にデータを送信するたびに、後続の処理データがあるかどうかを下流に報告します。同時に、SplitStates に現在のスプリット処理の進行状況を記録し、現在の状態と処理位置を記録します。
SplitEnumerator が外部システムで新しいスプリットを検出すると、RPC 経由で addSplits メソッドを呼び出して新しいスプリットをリーダーに追加する必要があります。SplitFetchermanager 側では、事前に選択したスレッドモデルに従って新しいスプリットが割り当てられます(単一スレッドのみの場合、そのスレッドに新しいタスクを割り当て、リーダーに新しいスプリットを読み取らせます。マルチスレッド実装の場合、新しいスレッドとリーダーを作成してスプリットを別々に処理します)。同様に、SplitStates に現在の処理進行状況を記録する必要があります。
5. チェックポイントの作成
次に、新しい Source API でチェックポイントがどのように処理されるかを見ていきます。
まず、左側のコーディネーター、スプリット列挙子です。図に示すように、現在まだ割り当てられていないスプリット(Split_5)があります。中間の矢印は転送中のスプリットです。破線はこのチェックポイントの境界です。2 番目のスプリットはチェックポイントより前に、4 番目のスプリットはチェックポイントより後ろにあり、下部のリーダーは SplitEnumerator に新しいスプリットを要求していることが分かります。次にリーダーを見ると、3 つのリーダーにはそれぞれスプリットが割り当てられ、処理済みで、既に位置情報(Position)を持っています。では、チェックポイント時に列挙子とリーダーが保存すべきものを見ていきましょう。
列挙子:未割り当てのレコードスプリット(Split_5)、割り当て済みだが未チェックポイントのレコードスプリット(Split_4)
リーダー:割り当て済みレコードスプリット(Split_0、1、3)、スプリット割り当て状態(Split_2)
6. Source を実装する 3 つの簡単なステップ
1) Split/SplitState
Split:外部システムのフラグメンテーション
SplitSerializer:Split をシリアル化/逆シリアル化して SourceReader に渡す
SplitState:Split 状態、チェックポイントとリカバリに使用
2) SplitEnumerator
Split の検出とサブスクライブ
EnumState:Enumerator の状態、チェックポイントとリカバリに使用
EnumStateSerializer:EnumState をシリアル化/逆シリアル化
3) SourceReader
SplitReader:外部システムとのデータやり取りのインターフェイス
FetcherManager:スレッドモデルの選択(現在利用可能な実装あり)
RecordEmitter:メッセージタイプの変換とイベント時間の処理
よく考えると、これらのほとんどは実際には外部システムとのやり取りを担当しており、Flink エンジン自体とのやり取りはごくわずかであることが分かります。ユーザーはチェックポイントロックやマルチスレッドの問題などを心配する必要がなくなり、外部システムとの開発とやり取りにより多くの開発エネルギーを集中できます。したがって、新しい Source API はこれらの抽象化を通じて、開発者の開発作業を大幅に簡素化します。
三、Sink API — Flink データのエクスポート
Flink についてある程度理解していれば、exactly-once セマンティクスを実現でき、データは重複も損失もないことが分かります。この「正確に 1 回」を実現するために、Flink は多くの重要な作業を行っており、その 1 つがシンク側での 2 フェーズコミットの実装です。
1. プリコミットフェーズ
プリコミットフェーズでは、分散システムには一般的に「コーディネーター 1 + エグゼキュータ N」のモードがあるため、まずコーディネーターがコミット要求を送信する必要があります。つまり、すべてのエグゼキュータにコミット要求メッセージを送信し、2 フェーズコミット全体を開始します。
エグゼキュータがコミット要求を受信すると、コミットの準備作業を行います。すべての準備作業が完了したら、エグゼキュータはコーディネーターに「次のコミットの準備が完了した」と応答します。コーディネーターがすべてのエグゼキュータから「続行可能」の応答を受け取ると、プリコミットフェーズが終了し、コミット実行フェーズに入ります。
2. コミット実行フェーズ
コーディネーターはコミット決定メッセージをエグゼキュータに送信し、エグゼキュータは準備したコミット関連処理を実際に実行します。完了後、結果をコーディネーターに応答し、コミットが正常に実行されたことを報告します。
第 2 フェーズであるコミット実行フェーズに入ると、すべてのエグゼキュータはこの決定を無条件で実行しなければなりません。つまり、あるコーディネーターがこの段階で問題が発生しても、復旧後にこの決定を必ず実行する必要があります。つまり、一度コミットが決定されたら、エグゼキュータは必ずコミットアクションを実行しなければなりません。
プリコミットフェーズで、エグゼキュータがコミット直前に障害を経験し、正しいコミットアクションを実行できない場合があります。その場合、ネットワーク切断やタイムアウトなどにより、コーディネーターにエラーを応答することがあります。一定期間経過後にコーディネーターが第 3 のエグゼキュータからの応答を受信しなかった場合、コーディネーターは第 2 段階のロールバックアクションをトリガーします。つまり、すべてのエグゼキュータに「このコミット試行は失敗したため、全員が前の状態にロールバックする必要がある」と通知します。そして、エグゼキュータはロールバックアクションを実行して以前の操作を元に戻します。
3. Flink における 2 フェーズコミットの実践
1) プリコミットフェーズ
ファイルシステムのシンクを例に取ります。
ファイルシステムのシンクは、チェックポイント境界を受信した後にプリコミットアクションを実行します(現在のデータをディスク上の一時ファイルに書き込みます)。プリコミットフェーズが完了すると、すべてのオペレーターはコーディネーターに「コミット準備完了」と応答します。
2) コミット実行フェーズ
第 2 フェーズ、コミット実行フェーズが開始されます。JobManager はすべてのオペレーターにコミット指示を送信し、シンクはこれを受信すると最終的なコミットアクションを実行します。
ファイルシステムを例に取ると、前述の通り、プリコミットフェーズでデータは一時ファイルに書き込まれます。実際のコミット実行時に、一時ファイルは事前に定義された正規の名前にリネームされます。このリネーム操作がコミットの実行に相当します。
ここで注意すべきは、一時ファイルの設定は無駄ではなく、後続のロールバックなどの状況に備える重要な役割を果たしている点です。2 フェーズコミットメカニズムを巧みに活用することで、exactly-once セマンティクスを保証しています。
4. Sink モデル
1) Writer:上流からの連続データを中間状態に書き込む、またはプリコミットフェーズでの処理を担当します。
2) Committable:前述の「中間状態」は、コミット操作を実行可能なコンポーネントです。
3) Committer:Committable を実際にコミットします
4) Global Committer:グローバルコミッター。このコンポーネントはオプションで、外部システムに依存します。例:Iceberg。
四、将来の開発
新しい Source の改善
Source と Sink はリリースされたばかりで、まだいくつかの問題が残っています。開発者から新しい要件や機能改善の要望があり、現在は比較的安定した状態ですが、継続的な改善が必要です。
既存のコネクタを新しい API に移行
ストリームバッチ統合コネクタの継続的な進展に伴い、すべてのコネクタが新しい API に移行されます。
コネクタテストフレームワーク
コネクタテストフレームワークは、すべてのコネクタに対して一貫した統一テスト基準を提供しようとします。テスト開発者は個別にテストケースを書いたり、さまざまなテスト環境やテストシナリオを考慮したりする必要がなくなります。開発者がブロックを組み立てるように異なるシナリオやユースケースでコードを素早くテストできるようにし、ロジック自体の開発により多くの開発エネルギーを集中できるようにして、開発者のテスト負荷を大幅に削減します。これは Source API と Sink API、および後続フレームワーク開発の一貫した目標でもあります。コネクタ開発をより簡単にして参入障壁を下げ、より多くの開発者に Flink エコシステムへの貢献を促すことです。
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
