Analysis of Spark Codegen

背景

SparkSQL の優れたパフォーマンスを支える技術的な柱は、オプティマイザーとランタイムの 2 つです。オプティマイザーは最適な実行計画の発見に注力し、ランタイムは定められた実行計画をいかに高速に実行するかに注力します。ランタイムのさまざまな最適化は、以下の 2 つのレベルに集約できます。

1. グローバル最適化。グローバルなリソース使用率の向上、データスキューの解消、IO の削減という観点から最適化を行います。アダプティブ実行、シャッフル削除などが含まれます。

2. ローカル最適化。特定のタスクの実行効率を最適化します。主に式レベルと WholeStage レベルのコード生成技術に依存しています。

本記事では、Spark のコード生成の技術原理について解説します。

ケーススタディ
このセクションでは、2 つの具体的なケースを通じてコード生成のアプローチを紹介します。

式レベル
次の式の計算を考えます:x + (1 + 2)。Scala コードで以下のように表現されます。

Add(Attribute(x), Add(Literal(1), Literal(2)))
構文ツリーは以下のようになります。

expression

この構文ツリーを再帰的に評価する一般的なコードは以下の通りです。

tree. transformUp {
case Attribute(idx) => Literal(row. getValue(idx))
case Add(Literal(c1), Literal(c2)) => Literal(c1+c2)
case Literal(c) => Literal(c)
}

上記のコードを実行するには、型マッチング、仮想関数呼び出し、オブジェクト生成など、多くの追加ロジックが必要です。これらのオーバーヘッドは、式評価自体のコストをはるかに超えます。

これらのオーバーヘッドを排除するため、Spark Codegen は評価式の Java コードを直接生成し、その場でコンパイルします。具体的には、以下の 3 つのステップに分けられます。

1. コード生成。構文ツリーに基づいて Java コードを生成し、ラッパークラスにカプセル化します。

... // class wrapper
row.getValue(idx) + (1 + 2)
... // class wrapper

2. 実行時コンパイル。Janino フレームワークを使用して、生成されたコードをクラスファイルにコンパイルします。

3. ロードと実行。最後にロードして実行します。

最適化の前後で、パフォーマンスは桁違いに向上します。

expression

WholeStage レベル
次の SQL ステートメントを考えます。

select count(*) from store_sales
where ss_item_sk=1000;
生成される物理プランは以下の通りです。

expression

このプランを実装する一般的な方法は、ボルケーノモデルを使用することです。各オペレーターは Iterator インターフェイスを継承し、next() メソッドはまず上流の実行を駆動して入力を取得し、その後自身のロジックを実行します。コード例は以下の通りです。

class Agg extends Iterator[Row] {
def doAgg() {
while (child. hasNext()) {
val row = child. next();
// do aggregation
...
}
}
def next(): Row {
if (。doneAgg) {
doAgg();
}
return aggIter. next();
}
}

class Filter extends Iterator[Row] {
def next(): Row {
var current = child. next()
while (current 。= null && 。predicate(current)) {
current = child. next()
}
return current;
}
}

上記のコードから分かるように、ボルケーノモデルでは大量の型変換と仮想関数呼び出しが発生します。仮想関数呼び出しは CPU の分岐予測を失敗させ、深刻なパフォーマンス低下を引き起こします。

これらのオーバーヘッドを排除するため、Spark WholeStageCodegen は物理プランに対して型が確定した Java コードを生成し、Expression と同様にリアルタイムでコンパイルおよびロードを行います。この例で生成される Java コードの例は以下の通りです(実際のコードではありません。実際のコードスニペットは後述を参照してください)。

var count = 0
for (ss_item_sk in store_sales) {
if (ss_item_sk == 1000) {
count += 1
}
}

最適化前後のパフォーマンス向上データは以下の通りです。

expression

Spark Codegen フレームワーク
Spark Codegen フレームワークには、3 つのコアコンポーネントがあります。
1. コアインターフェイス/クラス
2. CodegenContext
3. Produce-Consume パターン

以下で詳しく説明します。

インターフェイス/クラス

4 つのコアインターフェイスを紹介します。

1. CodegenSupport (インターフェイス)
このインターフェイスを実装するオペレーターは、自身のロジックを Java コードに組み込めます。重要なメソッドは以下の通りです。

produce() // このノードが生成する Row の Java コードを出力
consume() // このノードが上流ノードから入力される Row を消費する Java コードを出力
実装クラスには、ProjectExec、FilterExec、HashAggregateExec、SortMergeJoinExec などがあります。

2. WholeStageCodegenExec (クラス)
CodegenSupport の実装クラスの一つです。Stage 内で CodegenSupport インターフェイスを実装するすべての隣接オペレーターを統合し、出力コードは融合されたすべてのオペレーターの実行ロジックをラッパークラスにカプセル化します。このクラスは Janino の実行時コンパイルの入力として使用されます。

3. InputAdapter (クラス)
CodegenSupport の実装クラスの一つで、接着剤的な役割を持つクラスです。WholeStageCodegenExec ノードと、CodegenSupport を実装しない上流ノードを接続するために使用されます。

4. BufferedRowIterator (インターフェイス)
WholeStageCodegenExec が生成する Java コードの親クラスです。重要なメソッドは以下の通りです。

public InternalRow next() // 次の Row を返す
public void append(InternalRow row) // Row を追加する
CodegenContext

生成コードを管理するコアクラスです。主に以下の機能を担当します。

1. 名前管理。同じスコープ内で変数名が競合しないようにします。
2. 変数管理。クラス変数を管理し、変数の型(独立した変数として宣言するか、型付き配列に圧縮するか)、変数の初期化ロジックなどを維持します。
3. メソッド管理。クラスメソッドを管理します。
4. 内部クラス管理。内部クラスを管理します。
5. 同一式管理。同じ部分式を管理し、二重計算を回避します。
6. サイズ管理。メソッドやクラスのサイズが大きくなりすぎないようにし、クラス変数が増えすぎないようにします。たとえば、式ブロックを複数の関数に分割したり、関数と変数の定義を複数の内部クラスに分割したりします。
7. 依存関係管理。このクラスが依存する外部オブジェクト(ブロードキャストオブジェクト、ユーティリティオブジェクト、計測オブジェクトなど)を管理します。
8. 汎用テンプレート管理。genComp、nullSafeExec などの共通コードテンプレートを提供します。

Produce-Consume パターン
隣接するオペレーターは、Produce-Consume パターンを通じてコードを生成します。
Produce は、全体の処理のフレームワークコードを生成します。たとえば、集約が生成するコードフレームワークは以下の通りです。

if (。initialized) {
# create a hash map, then build the aggregation hash map
# call child. produce()
initialized = true;
}
while (hashmap. hasNext()) {
row = hashmap. next();
# build the aggregation results
# create variables for results
# call consume(), which will call parent.doConsume()
if (shouldStop()) return;
}
Consume は、現在のノードが上流から入力される Row を処理するロジックを生成します。たとえば、Filter が生成するコードは以下の通りです。

# code to evaluate the predicate expression, result is isNull1 and value2
if (。isNull1 && value2) {
# call consume(), which will call parent.doConsume()
}

Related Articles

Explore More Special Offers

  1. 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

phone お問い合わせ
Hi, I'm Alibaba Cloud AI Assistant!
I can help with questions and solutions.