Introduction to Flink Sort-Shuffle Implementation
1. データシャッフルの概要
データシャッフルは、バッチデータ処理ジョブの重要なフェーズです。
このフェーズでは、上流の処理ノードの出力データが外部ストレージに永続化され、下流の計算ノードがそのデータを読み取って処理します。
これらの永続データは、計算ノード間のデータ交換の一形態であるだけでなく、エラー復旧においても重要な役割を果たします。
現在、既存の大規模分散コンピューティングシステムでは、ハッシュベース方式とソートベース方式の 2 つのバッチデータシャッフルモデルが採用されています。
ハッシュベース方式の基本的な考え方は、下流の異なる同時実行コンシューマータスクに送信されるデータをそれぞれ別々のファイルに書き込むことです。
これにより、ファイル自体が異なるデータパーティションを区別する自然な境界となります。
ソートベース方式の基本的な考え方は、すべてのパーティションのデータをまずまとめて書き込み、ソートを使用して異なるデータパーティションの境界を区別することです。
Flink ではバージョン 1.12 でソートベースのバッチシャッフル実装を導入し、その後も継続的にパフォーマンスと安定性の最適化を進めました。
バージョン 1.13 時点で、sort-shuffle は本番環境で利用可能な状態になっています。
2. Sort-Shuffle を導入する意義
Flink に sort-shuffle 実装を導入した重要な理由の 1 つは、Flink の従来のハッシュベース実装が大規模バッチジョブに対応できなかったためです。
この事実は、他の既存の大規模分散コンピューティングシステムによっても裏付けられています。
安定性の面では、同時実行度の高いバッチジョブにおいて、ハッシュベース実装は大量のファイルを生成し、これらのファイルを同時に読み書きするため、多くのリソースを消費し、ファイルシステムに大きな負荷をかけます。
ファイルシステムは大量のファイルメタデータを管理する必要があり、ファイルハンドルや inode の枯渇といった不安定化リスクが生じます。
パフォーマンスの面では、同時実行度の高いバッチジョブにおいて、大量のファイルを同時に読み書きすることは大量のランダム IO を意味し、各 IO で実際に読み書きされるデータ量は非常に小さくなる可能性があります。
これは IO パフォーマンスにとって極めて重要な課題です。
ディスク上では、データシャッフルがバッチジョブのパフォーマンスボトルネックになりやすくなります。
ソートベースのバッチデータシャッフル実装を導入することで、同時に読み書きするファイル数を大幅に削減でき、データの順次読み書きに有利に働き、Flink の大規模バッチ処理ジョブの安定性とパフォーマンスが向上します。
さらに、新しい sort-shuffle 実装はメモリーバッファーの消費量も削減できます。
ハッシュベース実装では、各データパーティションに読み書きバッファーが必要で、メモリーバッファーの消費量は同時実行数に正比例します。
ソートベース実装では、メモリーバッファー消費量とジョブの同時実行数のデカップリングを実現できます(ただし、より大きなメモリーはより高いパフォーマンスをもたらす場合があります)。
さらに重要なのは、読み書きに新たなストレージ構造と IO 最適化を実装したことで、Flink のバッチデータシャッフルが他の大規模分散データ処理システムよりも優れた優位性を持つようになった点です。
以降のセクションでは、Flink の sort-shuffle 実装とその成果について詳しく紹介します。
3. Flink Sort-Shuffle の実装
他の分散システムのバッチデータ sort-shuffle 実装と同様に、Flink のシャッフルプロセス全体も複数の重要なフェーズに分かれています。
これには、メモリーバッファーへのデータ書き込み、メモリーバッファーのソート、ソート済みデータのファイルへの書き出し、シャッフルデータの読み取り、および下流への送信が含まれます。
ただし、他のシステムと比較して、Flink の実装にはいくつかの根本的な違いがあります。
これには、マルチセグメントデータストレージフォーマット、データマージプロセスの省略、およびデータ読み取り IO スケジューリングが含まれ、これらにより Flink の実装はより優れたパフォーマンスを実現しています。
1. 設計目標
Flink の sort-shuffle 実装全体を通じて、以下のポイントを主な設計目標として考慮しました。
1.1 ファイル数の削減
前述の通り、ハッシュベース実装では大量のファイルが生成され、ファイル数を減らすことは安定性とパフォーマンスの向上に寄与します。
Sort-Spill-Merge 方式は、この目標を達成するために分散コンピューティングシステムで広く採用されています。
まず、データをメモリーバッファーに書き込み、メモリーバッファーがいっぱいになったらデータをソートし、ソート済みデータをファイルに書き出します。
このとき、総ファイル数は「総データ量 ÷ メモリーバッファーサイズ」となり、ファイル数が削減されます。
すべてのデータを書き出し終えたら、生成されたファイルを 1 つのファイルにマージし、ファイル数のさらなる削減と各データパーティションのサイズ拡大を実現します(順次読み取りに有利です)。
他のシステムの実装と比較して、Flink の実装には重要な違いがあります。
Flink は複数のファイルを書き出してからマージするのではなく、常に同じファイルにデータを追加していきます。
この利点は、ファイルが常に 1 つだけであり、ファイル数の最小化が実現されることです。
1.2 開くファイル数の削減
同時に開くファイル数が多すぎると、より多くのリソースを消費し、ファイルハンドルの不足問題を招きやすくなり、安定性が低下します。
したがって、開くファイル数を減らすことはシステムの安定性向上に寄与します。
データ書き込みについては、前述の通り、各同時実行タスクは常に同じファイルに追加するため、開くファイルは常に 1 つだけです。
データ読み取りについては、各ファイルは多数の下流同時実行タスクに読み取られますが、Flink はファイルを 1 回だけ開き、ファイルハンドルをこれらの同時実行読み取りタスク間で共有することで、各ファイルを 1 回だけ開くという目標を達成しています。
1.3 順次読み書きの最大化
ファイルの順次読み書きは、ファイルの IO パフォーマンスにとって極めて重要です。
シャッフルファイル数を削減することで、ファイルのランダム IO をある程度軽減しました。
さらに、Flink のバッチデータ sort-shuffle は、ファイルの順次読み書きを最大化するための IO 最適化も実装しています。
データ書き込みフェーズでは、書き出すデータバッファーをより大きなバッチに集約し、writev システムコールで書き出すことで、より優れた順次書き込みを実現しています。
データ読み取りフェーズでは、読み取り IO スケジューリングを導入し、データ読み取りリクエストが常にファイルオフセット順序で処理されるようにすることで、ファイルの順次読み取りを最大化しています。
実験により、これらの最適化がバッチデータシャッフルのパフォーマンスを大幅に向上させることが示されています。
1.4 読み書き IO 増幅の削減
従来の Sort-Spill-Merge 方式では、生成された複数のファイルをより大きなファイルにマージして読み取りデータブロックのサイズを大きくします。
この実装方式には利点がある一方で、いくつかの欠点もあります。
最も大きな問題は読み書き IO の増幅です。
計算ノード間のデータシャッフルにおいて、エラーが発生しなければデータは 1 回だけ書き込まれ、1 回だけ読み取られれば十分です。
しかし、データマージにより同じデータが複数回読み書きされ、総 IO 量が増加し、ストレージスペースも大量に消費されます。
Flink の実装は、同じファイルへの継続的なデータ追加と独自のストレージ構造により、ファイルマージプロセスを回避しています。
単一データブロックのサイズはマージ後のサイズより小さいものの、ファイルマージのオーバーヘッドを回避し、Flink 独自の IO スケジューリングと組み合わせることで、最終的に Sort-Spill-Merge 方式よりも高いパフォーマンスを実現できます。
1.5 メモリーバッファー消費量の削減
他の分散コンピューティングシステムの sort-shuffle 実装と同様に、Flink は固定サイズのメモリーバッファーをデータのキャッシングとソートに使用します。
このメモリーバッファーのサイズは同時実行数に依存しないため、上流のシャッフルデータ書き込みに必要なメモリーバッファーのサイズが同時実行数から切り離されます。
さらに、別のメモリー管理最適化 FLINK-16428 と組み合わせることで、下流のシャッフルデータ読み取りにおいても同時実行数に依存しないメモリーバッファー消費を実現し、大規模バッチジョブのメモリーバッファー消費量を削減できます(注:FLINK-16428 はバッチジョブとストリーミングジョブの両方に適用されます)。
2. 実装の詳細
2.1 メモリーデータのソート
シャッフルデータの sort-spill フェーズでは、各データはまずシリアライズされてソートバッファーに書き込まれます。
バッファーがいっぱいになったら、バッファー内のすべてのバイナリデータがデータパーティションの順序でソートされます。
その後、ソート済みデータがデータパーティション順にファイルに書き出されます。
現時点ではデータ自体はソートされていませんが、ソートバッファーのインターフェイスは十分に一般化されており、将来的により複雑なソート要件にも対応できます。
ソートバッファーのインターフェイスは以下のように定義されています。
ソートアルゴリズムには、複雑性の低いバケットソートを採用しました。
具体的には、各シリアライズデータの先頭に 16 バイトのメタデータを挿入します。
これには、4 バイトの長さ、4 バイトのデータ型、および同じデータパーティション内の次のデータへの 8 バイトのポインターが含まれます。
構造を下図に示します。
バッファーからデータを読み取る際は、各データパーティションのチェーンインデックス構造に従って、そのデータパーティションに属するすべてのデータを読み取るだけでよく、これらのデータは書き込み時の順序を保持します。
これにより、データパーティション順序ですべてのデータを読み取ることで、データパーティションに基づくソートの目的を達成できます。
2.2 ファイルストレージ構造
前述の通り、各並列タスクが生成するシャッフルデータは 1 つの物理ファイルに書き込まれます。
各物理ファイルには複数のデータ領域が含まれており、各データ領域はデータバッファーの 1 回の sort-spill によって生成されます。
各データブロックでは、異なるデータパーティション(下流の計算ノードの異なる並列タスクによって消費される)に属するすべてのデータが、データパーティションの番号順にソートされて集約されています。
下図はシャッフルデータファイルの詳細構造を示しています。
ここで (R1, R2, R3) は 3 つの異なるデータブロックで、3 回の sort-spill 書き込みに対応します。
各データブロックには 3 つの異なるデータパーティションがあり、(C1, C2, C3) の 3 つの異なる並列コンシューマータスクによって読み取られます。
つまり、データ B1.1、B2.1、B3.1 は C1 で処理され、データ B1.2、B2.2、B3.2 は C2 で処理され、データ B1.3、B2.3、B3.3 は C3 で処理されます。
他の分散処理システムの実装と同様に、Flink では各データファイルに対応するインデックスファイルが存在します。
インデックスファイルは、読み取り時に各コンシューマーが自身に属するデータ(データパーティション)を検索するために使用されます。
インデックスファイルはデータファイルと同じデータ領域を持ち、各データ領域にはデータパーティション数と同じ数のインデックスエントリが含まれます。
各インデックスエントリは 2 つの部分からなり、データファイル内のオフセットとデータ長に対応します。
最適化として、Flink は各インデックスファイルに対して最大 4 MB のインデックスデータをキャッシュします。
データファイルとインデックスファイルの対応関係は以下の通りです。
2.3 読み取り IO スケジューリング
ファイル IO パフォーマンスをさらに向上させるため、上記のストレージ構造に基づき、Flink は IO スケジューリングメカニズムを導入しています。
これはディスクスケジューリングのエレベーターアルゴリズムに類似しています。
Flink の IO スケジューリングは、IO リクエストのファイルオフセット順序に従って常にスケジューリングされます。
具体的には、あるデータファイルに n 個のデータ領域があり、各データ領域に m 個のデータパーティションがあり、m 個の下流計算タスクがそのデータファイルを読み取る場合、以下の擬似コードは Flink の IO スケジューリングアルゴリズムのワークフローを示しています。
2.4 データブロードキャスト最適化
データブロードキャストとは、同じデータを下流の計算ノードのすべての並列タスクに送信することを指します。
一般的なアプリケーションシナリオはブロードキャスト結合です。
Flink の sort-shuffle 実装ではこのプロセスを最適化し、ブロードキャストデータをメモリーソートバッファーとシャッフルファイル内に 1 コピーのみ保存することで、データブロードキャストのパフォーマンスを大幅に向上できます。
具体的には、ブロードキャストデータをソートバッファーに書き込む際、そのデータは 1 回だけシリアライズされてコピーされ、データがシャッフルファイルに書き出される際も 1 コピーのみが書き込まれます。
インデックスファイルでは、異なるデータパーティションのデータインデックスエントリは、すべてデータファイル内の同じデータを指します。
下図はデータブロードキャスト最適化のすべての詳細を示しています。
2.5 データ圧縮
データ圧縮はシンプルかつ効果的な最適化手法です。
テスト結果によると、データ圧縮により TPC-DS の全体パフォーマンスが 30% 以上向上します。
Flink のハッシュベースバッチシャッフル実装と同様に、データ圧縮はネットワークバッファー単位で実行され、データ圧縮はデータパーティションをまたがりません。
つまり、異なる下流並列タスクに送信されるデータはそれぞれ個別に圧縮されます。
圧縮はデータのソート完了後かつ書き出し前に実行され、下流のコンシューマータスクはデータ受信後に伸張処理を行います。
下図はデータ圧縮の全プロセスを示しています。
4. テスト結果
1. 安定性
新しい sort-shuffle 実装により、Flink のバッチジョブ実行時の安定性が大幅に向上しました。
潜在的なファイルハンドルや inode の枯渇による不安定化問題を解決したほか、Flink の従来のハッシュシャッフルの既知の問題も解決しています。
たとえば、FLINK-21201(大量のファイル作成によりメインスレッドがブロックされる問題)や FLINK-19925(ネットワーク Netty スレッド内で実行される IO 操作によりネットワークの安定性が影響を受ける問題)などがあります。
2. パフォーマンス
1,000 同時実行規模の環境で TPC-DS 10 TB データ規模のテストを実行しました。
その結果、Flink の従来のバッチデータシャッフル実装と比較して、新しいデータシャッフル実装は 2〜6 倍のパフォーマンス向上を実現できることが示されました。
計算時間を除外し、データシャッフル時間のみを計測すると、最大 10 倍のパフォーマンス向上が見込めます。
以下の表はパフォーマンス向上の詳細データを示しています。
テストクラスターでは、各機械式ハードディスクのデータ読み書き帯域幅は 160 MB/s に達します。
注:テスト環境の構成は以下の通りです。大量のメモリーがあるため、シャッフルデータ量が少ない一部のジョブでは、実際のデータシャッフルがメモリーへの読み書きのみで行われます。そのため、上記の表にはデータ量が大きくパフォーマンス向上が顕著な一部のクエリのみを掲載しています。
5. チューニングパラメーター
Flink では、sort-shuffle はデフォルトで有効になっていません。有効にするには、パラメーター taskmanager.network.sort-shuffle.min-parallelism を設定する必要があります。このパラメーターの意味は、データパーティション数(1 つの計算タスクが下流の計算ノードにデータを送信する同時実行数)がこの値を下回る場合は hash-shuffle 実装が使用され、この値を超える場合は sort-shuffle が有効になることです。実際のアプリケーションでは、機械式ハードディスク上では 1 に設定して、常に sort-shuffle を使用できます。
Flink はデータ圧縮をデフォルトで有効にしていません。バッチジョブでは、データ圧縮率が低い場合を除き、ほとんどのシナリオで有効にすることが推奨されます。有効化パラメーターは taskmanager.network.blocking-shuffle.compression.enabled です。
シャッフルデータの書き込みとデータ読み取りには、いずれもメモリーバッファーが必要です。データ書き込みバッファーのサイズは taskmanager.network.sort-shuffle.min-buffers で制御され、データ読み取りバッファーは taskmanager.memory.framework.off-heap.batch-shuffle.size で制御されます。データ書き込みバッファーはネットワークメモリーから分割されます。データ書き込みバッファーを増やす場合は、ネットワークメモリーの総サイズも増やす必要がある場合があります。これにより、ネットワークメモリー不足エラーを回避できます。データ読み取りバッファーはフレームワークのオフヒープメモリーから分割されます。データ読み取りバッファーを増やす場合は、フレームワークのオフヒープメモリーも増やす必要がある場合があります。これにより、直接メモリーの OOM エラーを回避できます。一般的に、より大きなメモリーバッファーはより良いパフォーマンスをもたらします。大規模バッチジョブでは、数百 MB 程度のデータ書き込みバッファーと読み取りバッファーがあれば十分です。
データシャッフルは、バッチデータ処理ジョブの重要なフェーズです。
このフェーズでは、上流の処理ノードの出力データが外部ストレージに永続化され、下流の計算ノードがそのデータを読み取って処理します。
これらの永続データは、計算ノード間のデータ交換の一形態であるだけでなく、エラー復旧においても重要な役割を果たします。
現在、既存の大規模分散コンピューティングシステムでは、ハッシュベース方式とソートベース方式の 2 つのバッチデータシャッフルモデルが採用されています。
ハッシュベース方式の基本的な考え方は、下流の異なる同時実行コンシューマータスクに送信されるデータをそれぞれ別々のファイルに書き込むことです。
これにより、ファイル自体が異なるデータパーティションを区別する自然な境界となります。
ソートベース方式の基本的な考え方は、すべてのパーティションのデータをまずまとめて書き込み、ソートを使用して異なるデータパーティションの境界を区別することです。
Flink ではバージョン 1.12 でソートベースのバッチシャッフル実装を導入し、その後も継続的にパフォーマンスと安定性の最適化を進めました。
バージョン 1.13 時点で、sort-shuffle は本番環境で利用可能な状態になっています。
2. Sort-Shuffle を導入する意義
Flink に sort-shuffle 実装を導入した重要な理由の 1 つは、Flink の従来のハッシュベース実装が大規模バッチジョブに対応できなかったためです。
この事実は、他の既存の大規模分散コンピューティングシステムによっても裏付けられています。
安定性の面では、同時実行度の高いバッチジョブにおいて、ハッシュベース実装は大量のファイルを生成し、これらのファイルを同時に読み書きするため、多くのリソースを消費し、ファイルシステムに大きな負荷をかけます。
ファイルシステムは大量のファイルメタデータを管理する必要があり、ファイルハンドルや inode の枯渇といった不安定化リスクが生じます。
パフォーマンスの面では、同時実行度の高いバッチジョブにおいて、大量のファイルを同時に読み書きすることは大量のランダム IO を意味し、各 IO で実際に読み書きされるデータ量は非常に小さくなる可能性があります。
これは IO パフォーマンスにとって極めて重要な課題です。
ディスク上では、データシャッフルがバッチジョブのパフォーマンスボトルネックになりやすくなります。
ソートベースのバッチデータシャッフル実装を導入することで、同時に読み書きするファイル数を大幅に削減でき、データの順次読み書きに有利に働き、Flink の大規模バッチ処理ジョブの安定性とパフォーマンスが向上します。
さらに、新しい sort-shuffle 実装はメモリーバッファーの消費量も削減できます。
ハッシュベース実装では、各データパーティションに読み書きバッファーが必要で、メモリーバッファーの消費量は同時実行数に正比例します。
ソートベース実装では、メモリーバッファー消費量とジョブの同時実行数のデカップリングを実現できます(ただし、より大きなメモリーはより高いパフォーマンスをもたらす場合があります)。
さらに重要なのは、読み書きに新たなストレージ構造と IO 最適化を実装したことで、Flink のバッチデータシャッフルが他の大規模分散データ処理システムよりも優れた優位性を持つようになった点です。
以降のセクションでは、Flink の sort-shuffle 実装とその成果について詳しく紹介します。
3. Flink Sort-Shuffle の実装
他の分散システムのバッチデータ sort-shuffle 実装と同様に、Flink のシャッフルプロセス全体も複数の重要なフェーズに分かれています。
これには、メモリーバッファーへのデータ書き込み、メモリーバッファーのソート、ソート済みデータのファイルへの書き出し、シャッフルデータの読み取り、および下流への送信が含まれます。
ただし、他のシステムと比較して、Flink の実装にはいくつかの根本的な違いがあります。
これには、マルチセグメントデータストレージフォーマット、データマージプロセスの省略、およびデータ読み取り IO スケジューリングが含まれ、これらにより Flink の実装はより優れたパフォーマンスを実現しています。
1. 設計目標
Flink の sort-shuffle 実装全体を通じて、以下のポイントを主な設計目標として考慮しました。
1.1 ファイル数の削減
前述の通り、ハッシュベース実装では大量のファイルが生成され、ファイル数を減らすことは安定性とパフォーマンスの向上に寄与します。
Sort-Spill-Merge 方式は、この目標を達成するために分散コンピューティングシステムで広く採用されています。
まず、データをメモリーバッファーに書き込み、メモリーバッファーがいっぱいになったらデータをソートし、ソート済みデータをファイルに書き出します。
このとき、総ファイル数は「総データ量 ÷ メモリーバッファーサイズ」となり、ファイル数が削減されます。
すべてのデータを書き出し終えたら、生成されたファイルを 1 つのファイルにマージし、ファイル数のさらなる削減と各データパーティションのサイズ拡大を実現します(順次読み取りに有利です)。
他のシステムの実装と比較して、Flink の実装には重要な違いがあります。
Flink は複数のファイルを書き出してからマージするのではなく、常に同じファイルにデータを追加していきます。
この利点は、ファイルが常に 1 つだけであり、ファイル数の最小化が実現されることです。
1.2 開くファイル数の削減
同時に開くファイル数が多すぎると、より多くのリソースを消費し、ファイルハンドルの不足問題を招きやすくなり、安定性が低下します。
したがって、開くファイル数を減らすことはシステムの安定性向上に寄与します。
データ書き込みについては、前述の通り、各同時実行タスクは常に同じファイルに追加するため、開くファイルは常に 1 つだけです。
データ読み取りについては、各ファイルは多数の下流同時実行タスクに読み取られますが、Flink はファイルを 1 回だけ開き、ファイルハンドルをこれらの同時実行読み取りタスク間で共有することで、各ファイルを 1 回だけ開くという目標を達成しています。
1.3 順次読み書きの最大化
ファイルの順次読み書きは、ファイルの IO パフォーマンスにとって極めて重要です。
シャッフルファイル数を削減することで、ファイルのランダム IO をある程度軽減しました。
さらに、Flink のバッチデータ sort-shuffle は、ファイルの順次読み書きを最大化するための IO 最適化も実装しています。
データ書き込みフェーズでは、書き出すデータバッファーをより大きなバッチに集約し、writev システムコールで書き出すことで、より優れた順次書き込みを実現しています。
データ読み取りフェーズでは、読み取り IO スケジューリングを導入し、データ読み取りリクエストが常にファイルオフセット順序で処理されるようにすることで、ファイルの順次読み取りを最大化しています。
実験により、これらの最適化がバッチデータシャッフルのパフォーマンスを大幅に向上させることが示されています。
1.4 読み書き IO 増幅の削減
従来の Sort-Spill-Merge 方式では、生成された複数のファイルをより大きなファイルにマージして読み取りデータブロックのサイズを大きくします。
この実装方式には利点がある一方で、いくつかの欠点もあります。
最も大きな問題は読み書き IO の増幅です。
計算ノード間のデータシャッフルにおいて、エラーが発生しなければデータは 1 回だけ書き込まれ、1 回だけ読み取られれば十分です。
しかし、データマージにより同じデータが複数回読み書きされ、総 IO 量が増加し、ストレージスペースも大量に消費されます。
Flink の実装は、同じファイルへの継続的なデータ追加と独自のストレージ構造により、ファイルマージプロセスを回避しています。
単一データブロックのサイズはマージ後のサイズより小さいものの、ファイルマージのオーバーヘッドを回避し、Flink 独自の IO スケジューリングと組み合わせることで、最終的に Sort-Spill-Merge 方式よりも高いパフォーマンスを実現できます。
1.5 メモリーバッファー消費量の削減
他の分散コンピューティングシステムの sort-shuffle 実装と同様に、Flink は固定サイズのメモリーバッファーをデータのキャッシングとソートに使用します。
このメモリーバッファーのサイズは同時実行数に依存しないため、上流のシャッフルデータ書き込みに必要なメモリーバッファーのサイズが同時実行数から切り離されます。
さらに、別のメモリー管理最適化 FLINK-16428 と組み合わせることで、下流のシャッフルデータ読み取りにおいても同時実行数に依存しないメモリーバッファー消費を実現し、大規模バッチジョブのメモリーバッファー消費量を削減できます(注:FLINK-16428 はバッチジョブとストリーミングジョブの両方に適用されます)。
2. 実装の詳細
2.1 メモリーデータのソート
シャッフルデータの sort-spill フェーズでは、各データはまずシリアライズされてソートバッファーに書き込まれます。
バッファーがいっぱいになったら、バッファー内のすべてのバイナリデータがデータパーティションの順序でソートされます。
その後、ソート済みデータがデータパーティション順にファイルに書き出されます。
現時点ではデータ自体はソートされていませんが、ソートバッファーのインターフェイスは十分に一般化されており、将来的により複雑なソート要件にも対応できます。
ソートバッファーのインターフェイスは以下のように定義されています。
ソートアルゴリズムには、複雑性の低いバケットソートを採用しました。
具体的には、各シリアライズデータの先頭に 16 バイトのメタデータを挿入します。
これには、4 バイトの長さ、4 バイトのデータ型、および同じデータパーティション内の次のデータへの 8 バイトのポインターが含まれます。
構造を下図に示します。
バッファーからデータを読み取る際は、各データパーティションのチェーンインデックス構造に従って、そのデータパーティションに属するすべてのデータを読み取るだけでよく、これらのデータは書き込み時の順序を保持します。
これにより、データパーティション順序ですべてのデータを読み取ることで、データパーティションに基づくソートの目的を達成できます。
2.2 ファイルストレージ構造
前述の通り、各並列タスクが生成するシャッフルデータは 1 つの物理ファイルに書き込まれます。
各物理ファイルには複数のデータ領域が含まれており、各データ領域はデータバッファーの 1 回の sort-spill によって生成されます。
各データブロックでは、異なるデータパーティション(下流の計算ノードの異なる並列タスクによって消費される)に属するすべてのデータが、データパーティションの番号順にソートされて集約されています。
下図はシャッフルデータファイルの詳細構造を示しています。
ここで (R1, R2, R3) は 3 つの異なるデータブロックで、3 回の sort-spill 書き込みに対応します。
各データブロックには 3 つの異なるデータパーティションがあり、(C1, C2, C3) の 3 つの異なる並列コンシューマータスクによって読み取られます。
つまり、データ B1.1、B2.1、B3.1 は C1 で処理され、データ B1.2、B2.2、B3.2 は C2 で処理され、データ B1.3、B2.3、B3.3 は C3 で処理されます。
他の分散処理システムの実装と同様に、Flink では各データファイルに対応するインデックスファイルが存在します。
インデックスファイルは、読み取り時に各コンシューマーが自身に属するデータ(データパーティション)を検索するために使用されます。
インデックスファイルはデータファイルと同じデータ領域を持ち、各データ領域にはデータパーティション数と同じ数のインデックスエントリが含まれます。
各インデックスエントリは 2 つの部分からなり、データファイル内のオフセットとデータ長に対応します。
最適化として、Flink は各インデックスファイルに対して最大 4 MB のインデックスデータをキャッシュします。
データファイルとインデックスファイルの対応関係は以下の通りです。
2.3 読み取り IO スケジューリング
ファイル IO パフォーマンスをさらに向上させるため、上記のストレージ構造に基づき、Flink は IO スケジューリングメカニズムを導入しています。
これはディスクスケジューリングのエレベーターアルゴリズムに類似しています。
Flink の IO スケジューリングは、IO リクエストのファイルオフセット順序に従って常にスケジューリングされます。
具体的には、あるデータファイルに n 個のデータ領域があり、各データ領域に m 個のデータパーティションがあり、m 個の下流計算タスクがそのデータファイルを読み取る場合、以下の擬似コードは Flink の IO スケジューリングアルゴリズムのワークフローを示しています。
2.4 データブロードキャスト最適化
データブロードキャストとは、同じデータを下流の計算ノードのすべての並列タスクに送信することを指します。
一般的なアプリケーションシナリオはブロードキャスト結合です。
Flink の sort-shuffle 実装ではこのプロセスを最適化し、ブロードキャストデータをメモリーソートバッファーとシャッフルファイル内に 1 コピーのみ保存することで、データブロードキャストのパフォーマンスを大幅に向上できます。
具体的には、ブロードキャストデータをソートバッファーに書き込む際、そのデータは 1 回だけシリアライズされてコピーされ、データがシャッフルファイルに書き出される際も 1 コピーのみが書き込まれます。
インデックスファイルでは、異なるデータパーティションのデータインデックスエントリは、すべてデータファイル内の同じデータを指します。
下図はデータブロードキャスト最適化のすべての詳細を示しています。
2.5 データ圧縮
データ圧縮はシンプルかつ効果的な最適化手法です。
テスト結果によると、データ圧縮により TPC-DS の全体パフォーマンスが 30% 以上向上します。
Flink のハッシュベースバッチシャッフル実装と同様に、データ圧縮はネットワークバッファー単位で実行され、データ圧縮はデータパーティションをまたがりません。
つまり、異なる下流並列タスクに送信されるデータはそれぞれ個別に圧縮されます。
圧縮はデータのソート完了後かつ書き出し前に実行され、下流のコンシューマータスクはデータ受信後に伸張処理を行います。
下図はデータ圧縮の全プロセスを示しています。
4. テスト結果
1. 安定性
新しい sort-shuffle 実装により、Flink のバッチジョブ実行時の安定性が大幅に向上しました。
潜在的なファイルハンドルや inode の枯渇による不安定化問題を解決したほか、Flink の従来のハッシュシャッフルの既知の問題も解決しています。
たとえば、FLINK-21201(大量のファイル作成によりメインスレッドがブロックされる問題)や FLINK-19925(ネットワーク Netty スレッド内で実行される IO 操作によりネットワークの安定性が影響を受ける問題)などがあります。
2. パフォーマンス
1,000 同時実行規模の環境で TPC-DS 10 TB データ規模のテストを実行しました。
その結果、Flink の従来のバッチデータシャッフル実装と比較して、新しいデータシャッフル実装は 2〜6 倍のパフォーマンス向上を実現できることが示されました。
計算時間を除外し、データシャッフル時間のみを計測すると、最大 10 倍のパフォーマンス向上が見込めます。
以下の表はパフォーマンス向上の詳細データを示しています。
テストクラスターでは、各機械式ハードディスクのデータ読み書き帯域幅は 160 MB/s に達します。
注:テスト環境の構成は以下の通りです。大量のメモリーがあるため、シャッフルデータ量が少ない一部のジョブでは、実際のデータシャッフルがメモリーへの読み書きのみで行われます。そのため、上記の表にはデータ量が大きくパフォーマンス向上が顕著な一部のクエリのみを掲載しています。
5. チューニングパラメーター
Flink では、sort-shuffle はデフォルトで有効になっていません。有効にするには、パラメーター taskmanager.network.sort-shuffle.min-parallelism を設定する必要があります。このパラメーターの意味は、データパーティション数(1 つの計算タスクが下流の計算ノードにデータを送信する同時実行数)がこの値を下回る場合は hash-shuffle 実装が使用され、この値を超える場合は sort-shuffle が有効になることです。実際のアプリケーションでは、機械式ハードディスク上では 1 に設定して、常に sort-shuffle を使用できます。
Flink はデータ圧縮をデフォルトで有効にしていません。バッチジョブでは、データ圧縮率が低い場合を除き、ほとんどのシナリオで有効にすることが推奨されます。有効化パラメーターは taskmanager.network.blocking-shuffle.compression.enabled です。
シャッフルデータの書き込みとデータ読み取りには、いずれもメモリーバッファーが必要です。データ書き込みバッファーのサイズは taskmanager.network.sort-shuffle.min-buffers で制御され、データ読み取りバッファーは taskmanager.memory.framework.off-heap.batch-shuffle.size で制御されます。データ書き込みバッファーはネットワークメモリーから分割されます。データ書き込みバッファーを増やす場合は、ネットワークメモリーの総サイズも増やす必要がある場合があります。これにより、ネットワークメモリー不足エラーを回避できます。データ読み取りバッファーはフレームワークのオフヒープメモリーから分割されます。データ読み取りバッファーを増やす場合は、フレームワークのオフヒープメモリーも増やす必要がある場合があります。これにより、直接メモリーの OOM エラーを回避できます。一般的に、より大きなメモリーバッファーはより良いパフォーマンスをもたらします。大規模バッチジョブでは、数百 MB 程度のデータ書き込みバッファーと読み取りバッファーがあれば十分です。
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
