EMR Spark-SQL performance optimization reveals the secret
はじめに
Alibaba Cloud E-MapReduce チームは、最近 TPCDS-Perf に最新の結果を提出しました。2 位 (2019 年に同じく EMR チームが提出した記録) と比較して、パフォーマンスとコストパフォーマンスの両面で 2 倍以上の改善を達成しています。詳細は TPCDS Perf をご覧ください。
tpcds
Alibaba Cloud E-MapReduce チームは、プロダクト、ユーザビリティ、セキュリティなどさまざまな側面に多くの研究開発リソースとエネルギーを投入し、EMR のような高い評価を受けるビッグデータプロダクトを創り出してきました。同時に、エンジンレベルでも長期間にわたり投資を続けており、オープンソースソフトウェアとの 100% の互換性を維持しつつ、チームの技術力でプロダクトに技術的な優位性を構築することを目的としてきました。これにより、お客様はオープンソースソフトウェアスタックをよりコスト効率よく利用でき、クラウド上のコストを極限まで削減できます。お客様がクラウド移行のプロセスにおいて、いかなる不安や懸念も抱くことなく進められるようにすることが私たちの目標です。
TPCDS Perf での Alibaba Cloud E-MapReduce チームの成果は、Spark エンジンにおけるチームの技術力と深い専門性を十分に証明するものです。今後、2020 年のランキング挑戦プロセスで行った最適化ポイントを紹介します。コミュニティの Spark エンジン開発者や Spark アプリケーション開発者の皆様、ぜひご注目ください。意見交換も歓迎します。そして何より、Alibaba Cloud E-MapReduce チームへの応募をお待ちしています。優秀な人材を切に求めています。
三度目のランキング更新に向けた目標
TPCDS Perf のページから分かるように、EMR チームは 10 TB スケールで 3 つの結果を提出しています。3 回目のランキング挑戦には、少しばかりのストーリーがあります。TPCDS Perf では、最終的にパフォーマンス指標とコストパフォーマンス指標という 2 つの指標に注目します。プロジェクト開始時に、私たちは困難な目標を掲げました。物理ハードウェアを変更しない条件下で、ソフトウェアの最適化により 2 倍以上の改善を実現し、パフォーマンス指標とコストパフォーマンス指標の両方を倍増させるというものです。
オープンソース Spark との比較データ
スコア提出後、オープンソースの Spark V2.4.3 を使用して TPCDS 99 クエリのテストを実施しました。以下がパフォーマンスデータの比較です。
ロードフェーズで約 3 倍のパフォーマンス向上
load
PT ステージで約 6 倍のパフォーマンス向上
pt
補足:コミュニティ版 Spark V2.4.3 では、Query 14 と Query 95 が OOM のため実行できず、集計から除外されています。
コミュニティ版 Spark で実行時間が 200 秒を超えるクエリの比較
_200
補足:これらのクエリの中で最も改善幅が小さい Query 78 でも 3 倍のパフォーマンス向上があり、Query 57 では 100 倍近いパフォーマンス向上を達成しています。
最適化ポイントの概要
オプティマイザー
InMemoryTable Cache を活用した CTE マテリアライゼーション
簡単に言えば、InMemoryTable Cache を最大限に活用して、不要な重複計算を削減する手法です。たとえば、Query 23A/B のスカラー計算は非常に重い操作で、繰り返し実行する必要があります。CTE 最適化パターンにより、繰り返し計算が必要な時間集約的な操作を特定し、InMemoryTable Cache を使用して E2E 時間を全体的に短縮します。
より効果的なフィルター関連の最適化
動的パーティションプルーニング:この機能はコミュニティの最新バージョン 3.0 でのみ利用可能です
小さなテーブルのブロードキャスト多重化:フィルター条件を持つ小さなテーブルが 2 つ以上のテーブルのデータを絞り込める場合、この小さなテーブルのフィルター効果を再利用できます。Query 64 が好例です
SMJ 前の BloomFilter:SMJ (SortMergeJoin) が実際に実行される前に BloomFilter を適用することで、Join プロセスのデータをさらに削減し、SpillDisk の問題を最大限に解消します
PK/FK 制約の最適化:主キーと外部キーの情報を通じて、オプティマイザーにさらなる最適化のヒントを提供します
RI-Join:主キーと外部キーに基づくファクトテーブルとディメンションテーブルの Join を削除します。ディメンションテーブルのカラムが射影されない場合、この Join 自体を実行する必要がありません
GroupBy Keys からの非主キー列の削除:GroupBy Keys に主キー列と非主キー列の両方が含まれる場合、主キー列が既にユニーク情報を保持しているため、非主キー列は GroupBy の結果に影響しません
Join 前の GroupBy Push Down
Fast Decimal
ランタイム時のテーブル分析と統計情報に基づき、オプティマイザーは特定の Decimal 演算を Long または Int 演算に最適化できます。TPCDS 99 クエリには多数の Decimal 演算が含まれているため、大きなパフォーマンス向上が期待できます。
ランタイム
今回の最適化では、もう一つ非常に興味深い最適化があります。それが Native Runtime の導入です。前述のオプティマイザー最適化が特定のケースでの切り札だとすれば、Native Runtime は広範な効果をもたらす切り札です。その後の統計によると、Native Runtime の導入により、SQL クエリの E2E 実行時間が全体的に 15〜20% 改善されることが確認されており、TPCDS Perf においても大きなパフォーマンス改善ポイントとなっています。
Native Runtime の概要
オープンソース版の WholeStageCodeGeneration フレームワークをベースに、生成される Java コードを Weld IR に置き換えて実行します。Weld の詳細は http://weld.stanford.edu/ をご覧ください。プロジェクト全体において、Weld IR への置き換えは実際にはごく一部の作業に過ぎません。Weld IR を実行可能にするために、以下の作業が必要でした。
式の Weld IR CodeGen (TPCDS 内で完全サポート)
演算子の Weld IR CodeGen (SortMergeJoin は C++ で実装。それ以外は Weld IR で置き換え可能)
統一メモリレイアウト (OffHeap UnsafeRow => C++ & Weld Runtime)
バッチベースの実行フレームワーク (Java のようにレコードごとに生成コード内を流れる方式では、NativeRuntime における JNI と WeldRuntime 間のオーバーヘッドが大きすぎるため)
その他の高性能ネイティブ演算子:SortMergeJoin、PartitionBy、CSV 解析。これらの演算子は現時点では Weld IR のインターフェイスで直接実装できないため、C++ を使用してネイティブ実行を実現しています。
Alibaba Cloud E-MapReduce チームは、最近 TPCDS-Perf に最新の結果を提出しました。2 位 (2019 年に同じく EMR チームが提出した記録) と比較して、パフォーマンスとコストパフォーマンスの両面で 2 倍以上の改善を達成しています。詳細は TPCDS Perf をご覧ください。
tpcds
Alibaba Cloud E-MapReduce チームは、プロダクト、ユーザビリティ、セキュリティなどさまざまな側面に多くの研究開発リソースとエネルギーを投入し、EMR のような高い評価を受けるビッグデータプロダクトを創り出してきました。同時に、エンジンレベルでも長期間にわたり投資を続けており、オープンソースソフトウェアとの 100% の互換性を維持しつつ、チームの技術力でプロダクトに技術的な優位性を構築することを目的としてきました。これにより、お客様はオープンソースソフトウェアスタックをよりコスト効率よく利用でき、クラウド上のコストを極限まで削減できます。お客様がクラウド移行のプロセスにおいて、いかなる不安や懸念も抱くことなく進められるようにすることが私たちの目標です。
TPCDS Perf での Alibaba Cloud E-MapReduce チームの成果は、Spark エンジンにおけるチームの技術力と深い専門性を十分に証明するものです。今後、2020 年のランキング挑戦プロセスで行った最適化ポイントを紹介します。コミュニティの Spark エンジン開発者や Spark アプリケーション開発者の皆様、ぜひご注目ください。意見交換も歓迎します。そして何より、Alibaba Cloud E-MapReduce チームへの応募をお待ちしています。優秀な人材を切に求めています。
三度目のランキング更新に向けた目標
TPCDS Perf のページから分かるように、EMR チームは 10 TB スケールで 3 つの結果を提出しています。3 回目のランキング挑戦には、少しばかりのストーリーがあります。TPCDS Perf では、最終的にパフォーマンス指標とコストパフォーマンス指標という 2 つの指標に注目します。プロジェクト開始時に、私たちは困難な目標を掲げました。物理ハードウェアを変更しない条件下で、ソフトウェアの最適化により 2 倍以上の改善を実現し、パフォーマンス指標とコストパフォーマンス指標の両方を倍増させるというものです。
オープンソース Spark との比較データ
スコア提出後、オープンソースの Spark V2.4.3 を使用して TPCDS 99 クエリのテストを実施しました。以下がパフォーマンスデータの比較です。
ロードフェーズで約 3 倍のパフォーマンス向上
load
PT ステージで約 6 倍のパフォーマンス向上
pt
補足:コミュニティ版 Spark V2.4.3 では、Query 14 と Query 95 が OOM のため実行できず、集計から除外されています。
コミュニティ版 Spark で実行時間が 200 秒を超えるクエリの比較
_200
補足:これらのクエリの中で最も改善幅が小さい Query 78 でも 3 倍のパフォーマンス向上があり、Query 57 では 100 倍近いパフォーマンス向上を達成しています。
最適化ポイントの概要
オプティマイザー
InMemoryTable Cache を活用した CTE マテリアライゼーション
簡単に言えば、InMemoryTable Cache を最大限に活用して、不要な重複計算を削減する手法です。たとえば、Query 23A/B のスカラー計算は非常に重い操作で、繰り返し実行する必要があります。CTE 最適化パターンにより、繰り返し計算が必要な時間集約的な操作を特定し、InMemoryTable Cache を使用して E2E 時間を全体的に短縮します。
より効果的なフィルター関連の最適化
動的パーティションプルーニング:この機能はコミュニティの最新バージョン 3.0 でのみ利用可能です
小さなテーブルのブロードキャスト多重化:フィルター条件を持つ小さなテーブルが 2 つ以上のテーブルのデータを絞り込める場合、この小さなテーブルのフィルター効果を再利用できます。Query 64 が好例です
SMJ 前の BloomFilter:SMJ (SortMergeJoin) が実際に実行される前に BloomFilter を適用することで、Join プロセスのデータをさらに削減し、SpillDisk の問題を最大限に解消します
PK/FK 制約の最適化:主キーと外部キーの情報を通じて、オプティマイザーにさらなる最適化のヒントを提供します
RI-Join:主キーと外部キーに基づくファクトテーブルとディメンションテーブルの Join を削除します。ディメンションテーブルのカラムが射影されない場合、この Join 自体を実行する必要がありません
GroupBy Keys からの非主キー列の削除:GroupBy Keys に主キー列と非主キー列の両方が含まれる場合、主キー列が既にユニーク情報を保持しているため、非主キー列は GroupBy の結果に影響しません
Join 前の GroupBy Push Down
Fast Decimal
ランタイム時のテーブル分析と統計情報に基づき、オプティマイザーは特定の Decimal 演算を Long または Int 演算に最適化できます。TPCDS 99 クエリには多数の Decimal 演算が含まれているため、大きなパフォーマンス向上が期待できます。
ランタイム
今回の最適化では、もう一つ非常に興味深い最適化があります。それが Native Runtime の導入です。前述のオプティマイザー最適化が特定のケースでの切り札だとすれば、Native Runtime は広範な効果をもたらす切り札です。その後の統計によると、Native Runtime の導入により、SQL クエリの E2E 実行時間が全体的に 15〜20% 改善されることが確認されており、TPCDS Perf においても大きなパフォーマンス改善ポイントとなっています。
Native Runtime の概要
オープンソース版の WholeStageCodeGeneration フレームワークをベースに、生成される Java コードを Weld IR に置き換えて実行します。Weld の詳細は http://weld.stanford.edu/ をご覧ください。プロジェクト全体において、Weld IR への置き換えは実際にはごく一部の作業に過ぎません。Weld IR を実行可能にするために、以下の作業が必要でした。
式の Weld IR CodeGen (TPCDS 内で完全サポート)
演算子の Weld IR CodeGen (SortMergeJoin は C++ で実装。それ以外は Weld IR で置き換え可能)
統一メモリレイアウト (OffHeap UnsafeRow => C++ & Weld Runtime)
バッチベースの実行フレームワーク (Java のようにレコードごとに生成コード内を流れる方式では、NativeRuntime における JNI と WeldRuntime 間のオーバーヘッドが大きすぎるため)
その他の高性能ネイティブ演算子:SortMergeJoin、PartitionBy、CSV 解析。これらの演算子は現時点では Weld IR のインターフェイスで直接実装できないため、C++ を使用してネイティブ実行を実現しています。
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
