Machine learning practice based on Spark and TensorFlow
EMR E-Learning プラットフォーム
EMR E-Learning プラットフォームは、ビッグデータと AI テクノロジーを基盤としています。アルゴリズムを使用して既存データから機械学習モデルを構築し、トレーニングと予測を行います。現在、機械学習は顔認識、Natural Language Processing (NLP)、レコメンデーションシステム、コンピュータビジョンなど、多くの分野で広く利用されています。近年、ビッグデータとコンピューティング能力の向上により、AI テクノロジーは急速に発展しています。
機械学習における三つの重要な要素は、アルゴリズム、データ、コンピューティング能力です。EMR はビッグデータプラットフォームとして、従来のデータウェアハウスデータや画像データなど、多種多様なデータを保持しています。また、EMR は強力なスケジューリング機能を備えており、GPU および CPU リソースを適切に一時停止・スケジュールできます。機械学習アルゴリズムと組み合わせることで、優れた AI プラットフォームとなります。
典型的な AI 開発プロセスを次の図に示します。まず、データ収集として、携帯電話、ルーター、またはログデータをビッグデータフレームワークのデータレイクに取り込みます。次にデータ処理です。収集されたデータは、従来のビッグデータ ETL または特徴量エンジニアリングを通じて処理される必要があります。続いてモデルトレーニングです。特徴量エンジニアリングまたは ETL で処理されたデータがトレーニングに使用されます。最後に、トレーニング済みモデルの評価とデプロイメントを行います。モデル予測の結果はビッグデータプラットフォームに入力されて処理・分析され、このプロセス全体が繰り返されます。
次の図は AI 開発のプロセスを示しています。左側はシングルマシンまたはクラスターで、主に AI トレーニングと評価、およびデータストレージに使用されます。右側はビッグデータストレージで、主に特徴量エンジニアリングなどのビッグデータ処理に使用されます。同時に、左側でトレーニングされた機械学習モデルを予測に使用できます。
現在の AI 開発の課題は、主に以下の二点です。
・二つのクラスターの運用保守が複雑:図から分かるように、AI 開発に関わる二つのクラスターは独立しており、それぞれ別々にメンテナンスする必要があります。運用保守コストは複雑で、エラーが発生しやすくなります。
・トレーニング効率の低下:左右のクラスター間で大量のデータ転送とモデル転送が必要であり、エンドツーエンドのトレーニング遅延が大きくなります。
EMR は統合ビッグデータプラットフォームとして、多くの機能を備えています。最下層のインフラストラクチャ層は、GPU および CPU マシンをサポートしています。データストレージ層には HDFS と Alibaba Cloud OSS が含まれます。データアクセス層には Kafka と Flume が含まれます。リソーススケジューリング層の計算エンジンには YARN、Kubernetes、ZooKeeper が含まれます。計算エンジンの中核となる E-Learning プラットフォームは、現在広く利用されているオープンソースシステムの Spark を基盤としています。ここでの Spark は Jindo Spark を使用しており、EMR チームが Spark を修正・最適化した AI シナリオ向けのバージョンです。さらに、Spark 上の PAI TensorFlow も搭載されています。最後に、計算分析層ではデータ分析、特徴量エンジニアリング、AI トレーニング、ノートブック機能がユーザーに提供されます。
EMR プラットフォームには以下の特徴があります。
・統一されたリソース管理とスケジューリング:CPU、メモリ、GPU の詳細なリソーススケジューリングと割り当てをサポートし、YARN および Kubernetes のリソーススケジューリングフレームワークに対応しています。
・複数フレームワークのサポート:TensorFlow、MXNet、Caffe などに対応しています。
・Spark の汎用データ処理フレームワーク:データソース API を提供してさまざまなデータソースの読み取りを容易にし、特徴量エンジニアリングでは MLlib パイプラインが広く使用されています。
・Spark + ディープラーニングフレームワーク:Spark とディープラーニングフレームワークの統合サポートにより、Spark と TensorFlow 間の効率的なデータ転送、分散ディープラーニングトレーニングをサポートする Spark リソーススケジューリングモデルを実現しています。
・リソース監視とアラート:EMR APM システムは、複数のアラート方式で完全なアプリケーションおよびクラスター監視を提供します。
・使いやすさ:Jupyter Notebook および Python のマルチ環境デプロイメントをサポートし、エンドツーエンドの機械学習トレーニングプロセスなどを実現します。
EMR E-Learning には PAI TensorFlow が統合されており、ディープラーニングと大規模スパースデータシナリオの最適化をサポートします。
TensorFlow on Spark
市場調査の結果、ほとんどの顧客はディープラーニング前のデータ ETL と特徴量エンジニアリング段階でオープンソースの計算フレームワーク Spark を使用し、後段階で TensorFlow を広く使用していることが分かりました。そのため、TensorFlow と Spark を有機的に結合することを目標としています。TensorFlow on Spark には、主に以下の図にある 6 つの具体的な設計目標が含まれます。
TensorFlow on Spark は、ボトムレベルで PySpark アプリケーションフレームワークをカプセル化したものです。フレームワーク内で実装される主な機能は、まずユーザーの特徴量エンジニアリングタスクをスケジューリングし、次にディープラーニングの TensorFlow タスクをスケジューリングすることです。さらに、特徴量エンジニアリングで処理されたデータを下層の PAI TensorFlow Runtime に効率的かつ迅速に転送し、ディープラーニングと機械学習のトレーニングを実行する必要があります。Spark は現在ヘテロジニアススケジューリングをサポートしていないため、顧客が分散 TensorFlow を実行する場合、パラメータサーバタスクとワーカータスクの二つのタスクを同時に実行する必要があり、ユーザーが必要とするリソースに応じて異なる Spark エグゼキュータが生成されます。パラメータサーバタスクとワーカータスクは ZooKeeper を介してサービスを登録します。フレームワーク起動後、ユーザーが記述した特徴量エンジニアリングタスクがエグゼキュータにスケジュールされて実行されます。実行完了後、フレームワークはデータを下層の PAI TensorFlow Runtime に転送してトレーニングを行います。トレーニング後、データはデータレイクに保存され、後続のモデルリリースに備えます。
機械学習とディープラーニングにおいて、データのやり取りは効率化の重要なポイントです。そのため、TensorFlow on Spark はデータやり取り部分に一連の最適化を施しました。具体的には、Apache Arrow を使用した高速データ転送により、トレーニングデータを直接 TensorFlow Runtime API に供給し、トレーニングプロセス全体を高速化します。
TensorFlow on Spark のフォールトトレランスメカニズムを次の図に示します。ボトムレベルでは TensorFlow のチェックポイントメカニズムに依存し、ユーザーは定期的にトレーニングモデルをデータレイクに保存する必要があります。TensorFlow が再起動されると、最新のチェックポイントを読み込んでトレーニングを再開します。フォールトトレランスメカニズムはモードによって異なる処理方法をとります。分散タスクの場合、パラメータサーバとワーカータスクが起動され、デーモンプロセス内に両タスクが存在して対応するタスクの実行を監視します。MPI タスクの場合、Spark バリア実行メカニズムを使用してフォールトトレランスを実現します。タスクが失敗すると、失敗をマークしてすべてのタスクを再起動し、すべての環境変数を再設定します。TF タスクは最新のチェックポイントの読み取りを担当します。
TensorFlow on Spark の機能と使いやすさは、主に以下の点に表れています。
・複数のデプロイメント環境:Conda の指定、パッケージ Python ランタイム仮想環境のサポート、Docker の指定をサポート
・TensorFlow アーキテクチャのサポート:分散 TensorFlow のネイティブ PS アーキテクチャと分散 Horovod MPI アーキテクチャをサポート
・TensorFlow API のサポート:分散 TensorFlow Estimator 高レベル API と分散 TensorFlow Session 低レベル API をサポート
・各種フレームワークの迅速なサポート:顧客のニーズに応じて、MXNet などの新しい AI フレームワークを追加できます
EMR の顧客の多くはインターネット企業です。広告やプッシュ通知のビジネスシナリオは一般的です。次の図は典型的な広告プッシュのビジネスシナリオです。プロセス全体として、EMR の顧客はログデータを Kafka を通じてデータレイクにリアルタイムでプッシュします。TensorFlow on Spark はプロセスの前半を担当し、この段階でリアルタイムデータとオフラインデータを SparkSQL、MLlib などの Spark ツールを通じて ETL および特徴量エンジニアリングできます。データのトレーニング後、TensorFlow フレームワークを通じて PAI TensorFlow Runtime に効率的にデータを供給して大規模なトレーニングと最適化を実行し、モデルをデータレイクに保存します。
API レベルでは、TensorFlow on Spark は基底クラスを提供しており、ユーザーが実装する必要のある三つのメソッド(pre_Train、shutdown、train)が含まれています。pre_Train は、ユーザーが行うデータの読み込み、ETL、特徴量エンジニアリングなどのタスクで、Spark の DataFrame オブジェクトを返します。shutdown メソッドは、ユーザーの永続接続リソースの解放を実現します。train メソッドは、ユーザーが TensorFlow で実装するモデルの選択、オプティマイザ、最適化演算子などのコードです。最後に、pl_submit コマンドを通じて TensorFlow on Spark タスクを投入します。
次の図はレコメンデーションシステムの FM(Factorization Machine)の例です。これは比較的一般的なレコメンデーションアルゴリズムです。具体的なシナリオとして、ユーザーの過去の映画評価、映画ジャンル、公開日時に基づいて映画をスコアリングし、おすすめの映画をユーザーに推薦します。左側は特徴量エンジニアリングで、ユーザーは Spark のデータソース API を使用して映画情報と評価データを読み取れます。join や ETL 処理など、すべての Spark 操作はネイティブにサポートされています。右側は TensorFlow で、モデルとオプティマイザの選択に使用されます。現在、システム全体のコードは GitHub で公開されています。
最後に、EMR E-Learning プラットフォームは、ビッグデータ処理、ディープラーニング、機械学習、データレイク、GPU 機能を密接に組み合わせ、ワンストップのビッグデータ・機械学習プラットフォームを提供します。TensorFlow on Spark は効率的なデータやり取りプロセスと完全な機械学習トレーニングプロセスを提供し、Spark と TensorFlow を組み合わせ、PAI TensorFlow を使用したトレーニングの高速化を支援します。現在、E-Learning プラットフォームはパブリッククラウドでさまざまな顧客にサービスを提供しています。成功事例として、CPU クラスター規模は 1,000 台超、GPU クラスター規模は 100 台超に達しています。
EMR E-Learning プラットフォームは、ビッグデータと AI テクノロジーを基盤としています。アルゴリズムを使用して既存データから機械学習モデルを構築し、トレーニングと予測を行います。現在、機械学習は顔認識、Natural Language Processing (NLP)、レコメンデーションシステム、コンピュータビジョンなど、多くの分野で広く利用されています。近年、ビッグデータとコンピューティング能力の向上により、AI テクノロジーは急速に発展しています。
機械学習における三つの重要な要素は、アルゴリズム、データ、コンピューティング能力です。EMR はビッグデータプラットフォームとして、従来のデータウェアハウスデータや画像データなど、多種多様なデータを保持しています。また、EMR は強力なスケジューリング機能を備えており、GPU および CPU リソースを適切に一時停止・スケジュールできます。機械学習アルゴリズムと組み合わせることで、優れた AI プラットフォームとなります。
典型的な AI 開発プロセスを次の図に示します。まず、データ収集として、携帯電話、ルーター、またはログデータをビッグデータフレームワークのデータレイクに取り込みます。次にデータ処理です。収集されたデータは、従来のビッグデータ ETL または特徴量エンジニアリングを通じて処理される必要があります。続いてモデルトレーニングです。特徴量エンジニアリングまたは ETL で処理されたデータがトレーニングに使用されます。最後に、トレーニング済みモデルの評価とデプロイメントを行います。モデル予測の結果はビッグデータプラットフォームに入力されて処理・分析され、このプロセス全体が繰り返されます。
次の図は AI 開発のプロセスを示しています。左側はシングルマシンまたはクラスターで、主に AI トレーニングと評価、およびデータストレージに使用されます。右側はビッグデータストレージで、主に特徴量エンジニアリングなどのビッグデータ処理に使用されます。同時に、左側でトレーニングされた機械学習モデルを予測に使用できます。
現在の AI 開発の課題は、主に以下の二点です。
・二つのクラスターの運用保守が複雑:図から分かるように、AI 開発に関わる二つのクラスターは独立しており、それぞれ別々にメンテナンスする必要があります。運用保守コストは複雑で、エラーが発生しやすくなります。
・トレーニング効率の低下:左右のクラスター間で大量のデータ転送とモデル転送が必要であり、エンドツーエンドのトレーニング遅延が大きくなります。
EMR は統合ビッグデータプラットフォームとして、多くの機能を備えています。最下層のインフラストラクチャ層は、GPU および CPU マシンをサポートしています。データストレージ層には HDFS と Alibaba Cloud OSS が含まれます。データアクセス層には Kafka と Flume が含まれます。リソーススケジューリング層の計算エンジンには YARN、Kubernetes、ZooKeeper が含まれます。計算エンジンの中核となる E-Learning プラットフォームは、現在広く利用されているオープンソースシステムの Spark を基盤としています。ここでの Spark は Jindo Spark を使用しており、EMR チームが Spark を修正・最適化した AI シナリオ向けのバージョンです。さらに、Spark 上の PAI TensorFlow も搭載されています。最後に、計算分析層ではデータ分析、特徴量エンジニアリング、AI トレーニング、ノートブック機能がユーザーに提供されます。
EMR プラットフォームには以下の特徴があります。
・統一されたリソース管理とスケジューリング:CPU、メモリ、GPU の詳細なリソーススケジューリングと割り当てをサポートし、YARN および Kubernetes のリソーススケジューリングフレームワークに対応しています。
・複数フレームワークのサポート:TensorFlow、MXNet、Caffe などに対応しています。
・Spark の汎用データ処理フレームワーク:データソース API を提供してさまざまなデータソースの読み取りを容易にし、特徴量エンジニアリングでは MLlib パイプラインが広く使用されています。
・Spark + ディープラーニングフレームワーク:Spark とディープラーニングフレームワークの統合サポートにより、Spark と TensorFlow 間の効率的なデータ転送、分散ディープラーニングトレーニングをサポートする Spark リソーススケジューリングモデルを実現しています。
・リソース監視とアラート:EMR APM システムは、複数のアラート方式で完全なアプリケーションおよびクラスター監視を提供します。
・使いやすさ:Jupyter Notebook および Python のマルチ環境デプロイメントをサポートし、エンドツーエンドの機械学習トレーニングプロセスなどを実現します。
EMR E-Learning には PAI TensorFlow が統合されており、ディープラーニングと大規模スパースデータシナリオの最適化をサポートします。
TensorFlow on Spark
市場調査の結果、ほとんどの顧客はディープラーニング前のデータ ETL と特徴量エンジニアリング段階でオープンソースの計算フレームワーク Spark を使用し、後段階で TensorFlow を広く使用していることが分かりました。そのため、TensorFlow と Spark を有機的に結合することを目標としています。TensorFlow on Spark には、主に以下の図にある 6 つの具体的な設計目標が含まれます。
TensorFlow on Spark は、ボトムレベルで PySpark アプリケーションフレームワークをカプセル化したものです。フレームワーク内で実装される主な機能は、まずユーザーの特徴量エンジニアリングタスクをスケジューリングし、次にディープラーニングの TensorFlow タスクをスケジューリングすることです。さらに、特徴量エンジニアリングで処理されたデータを下層の PAI TensorFlow Runtime に効率的かつ迅速に転送し、ディープラーニングと機械学習のトレーニングを実行する必要があります。Spark は現在ヘテロジニアススケジューリングをサポートしていないため、顧客が分散 TensorFlow を実行する場合、パラメータサーバタスクとワーカータスクの二つのタスクを同時に実行する必要があり、ユーザーが必要とするリソースに応じて異なる Spark エグゼキュータが生成されます。パラメータサーバタスクとワーカータスクは ZooKeeper を介してサービスを登録します。フレームワーク起動後、ユーザーが記述した特徴量エンジニアリングタスクがエグゼキュータにスケジュールされて実行されます。実行完了後、フレームワークはデータを下層の PAI TensorFlow Runtime に転送してトレーニングを行います。トレーニング後、データはデータレイクに保存され、後続のモデルリリースに備えます。
機械学習とディープラーニングにおいて、データのやり取りは効率化の重要なポイントです。そのため、TensorFlow on Spark はデータやり取り部分に一連の最適化を施しました。具体的には、Apache Arrow を使用した高速データ転送により、トレーニングデータを直接 TensorFlow Runtime API に供給し、トレーニングプロセス全体を高速化します。
TensorFlow on Spark のフォールトトレランスメカニズムを次の図に示します。ボトムレベルでは TensorFlow のチェックポイントメカニズムに依存し、ユーザーは定期的にトレーニングモデルをデータレイクに保存する必要があります。TensorFlow が再起動されると、最新のチェックポイントを読み込んでトレーニングを再開します。フォールトトレランスメカニズムはモードによって異なる処理方法をとります。分散タスクの場合、パラメータサーバとワーカータスクが起動され、デーモンプロセス内に両タスクが存在して対応するタスクの実行を監視します。MPI タスクの場合、Spark バリア実行メカニズムを使用してフォールトトレランスを実現します。タスクが失敗すると、失敗をマークしてすべてのタスクを再起動し、すべての環境変数を再設定します。TF タスクは最新のチェックポイントの読み取りを担当します。
TensorFlow on Spark の機能と使いやすさは、主に以下の点に表れています。
・複数のデプロイメント環境:Conda の指定、パッケージ Python ランタイム仮想環境のサポート、Docker の指定をサポート
・TensorFlow アーキテクチャのサポート:分散 TensorFlow のネイティブ PS アーキテクチャと分散 Horovod MPI アーキテクチャをサポート
・TensorFlow API のサポート:分散 TensorFlow Estimator 高レベル API と分散 TensorFlow Session 低レベル API をサポート
・各種フレームワークの迅速なサポート:顧客のニーズに応じて、MXNet などの新しい AI フレームワークを追加できます
EMR の顧客の多くはインターネット企業です。広告やプッシュ通知のビジネスシナリオは一般的です。次の図は典型的な広告プッシュのビジネスシナリオです。プロセス全体として、EMR の顧客はログデータを Kafka を通じてデータレイクにリアルタイムでプッシュします。TensorFlow on Spark はプロセスの前半を担当し、この段階でリアルタイムデータとオフラインデータを SparkSQL、MLlib などの Spark ツールを通じて ETL および特徴量エンジニアリングできます。データのトレーニング後、TensorFlow フレームワークを通じて PAI TensorFlow Runtime に効率的にデータを供給して大規模なトレーニングと最適化を実行し、モデルをデータレイクに保存します。
API レベルでは、TensorFlow on Spark は基底クラスを提供しており、ユーザーが実装する必要のある三つのメソッド(pre_Train、shutdown、train)が含まれています。pre_Train は、ユーザーが行うデータの読み込み、ETL、特徴量エンジニアリングなどのタスクで、Spark の DataFrame オブジェクトを返します。shutdown メソッドは、ユーザーの永続接続リソースの解放を実現します。train メソッドは、ユーザーが TensorFlow で実装するモデルの選択、オプティマイザ、最適化演算子などのコードです。最後に、pl_submit コマンドを通じて TensorFlow on Spark タスクを投入します。
次の図はレコメンデーションシステムの FM(Factorization Machine)の例です。これは比較的一般的なレコメンデーションアルゴリズムです。具体的なシナリオとして、ユーザーの過去の映画評価、映画ジャンル、公開日時に基づいて映画をスコアリングし、おすすめの映画をユーザーに推薦します。左側は特徴量エンジニアリングで、ユーザーは Spark のデータソース API を使用して映画情報と評価データを読み取れます。join や ETL 処理など、すべての Spark 操作はネイティブにサポートされています。右側は TensorFlow で、モデルとオプティマイザの選択に使用されます。現在、システム全体のコードは GitHub で公開されています。
最後に、EMR E-Learning プラットフォームは、ビッグデータ処理、ディープラーニング、機械学習、データレイク、GPU 機能を密接に組み合わせ、ワンストップのビッグデータ・機械学習プラットフォームを提供します。TensorFlow on Spark は効率的なデータやり取りプロセスと完全な機械学習トレーニングプロセスを提供し、Spark と TensorFlow を組み合わせ、PAI TensorFlow を使用したトレーニングの高速化を支援します。現在、E-Learning プラットフォームはパブリッククラウドでさまざまな顧客にサービスを提供しています。成功事例として、CPU クラスター規模は 1,000 台超、GPU クラスター規模は 100 台超に達しています。
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
