参考
メモ
前提
Download によると、Sparkは2.4系まで対応しているようだ。 Issue-3070 によると、Spark3対応も進んでいるようだ。
Spark起動
Run で記載の通り、以下のように起動した。
1 | $ export SPARK_HOME=/usr/local/spark/default |
1 | scala> import com.intel.analytics.bigdl.utils.Engine |
とりあえず初期化までは動いた。
サンプルを動かす
Examples に記載のサンプルを動かす。 LeNet Train にあるLeNetの例が良さそう。
MNISTの画像をダウンロードして、展開した。
1 | $ mkdir ~/tmp |
展開したディレクトリを指定しながら、LeNetの学習を実行。
1 | $ spark-submit \ |
実行中の様子は以下の通り。
今回はローカルモードで実行したが、入力された学習データと同様のサイズのキャッシュがメモリ上に展開されていることがわかる。
サンプルの中身
2020/11/15時点のmasterブランチを確認する。
Trainクラス
上記サンプルで実行されているTrainを見る。 モデルを定義している箇所は以下の通り。
com/intel/analytics/bigdl/models/lenet/Train.scala:48
1 | val model = if (param.modelSnapshot.isDefined) { |
Modelインスタンスは以下の通り、Optimizerオブジェクトに渡され、 データセットに合わせたOptimizerが返される。 (例:データセットが分散データセットかどうか、など)
com/intel/analytics/bigdl/models/lenet/Train.scala:83
1 | val optimizer = Optimizer( |
参考までに、Optimizerの種類と判定の処理は以下の通り。
com/intel/analytics/bigdl/optim/Optimizer.scala:688
1 | dataset match { |
返されたOptimizerの、Optimizer#optimizeメソッドを利用し学習が実行される。
com/intel/analytics/bigdl/models/lenet/Train.scala:98
1 | optimizer |
Optimizerの種類
上記の内容を見るに、Optimizerにはいくつか種類がありそうだ。
- DistriOptimizer
- DistriOptimizerV2
- LocalOptimizer
DistriOptimizer
ひとまずメモリ使用に関連する箇所ということで、入力データの準備の処理を確認する。
com.intel.analytics.bigdl.optim.DistriOptimizer#optimize
メソッドには 以下のような箇所がある。
com/intel/analytics/bigdl/optim/DistriOptimizer.scala:870
1 | prepareInput() |
これは
com.intel.analytics.bigdl.optim.DistriOptimizer#prepareInput
メソッドであり、 内部的に
com.intel.analytics.bigdl.optim.AbstractOptimizer#prepareInput
メソッドを呼び出し、
入力データをSparkのキャッシュに載せるように処理する。
com/intel/analytics/bigdl/optim/DistriOptimizer.scala:808
1 | if (!dataset.toDistributed().isCached) { |
キャシュに載せると箇所は以下の通り。
com/intel/analytics/bigdl/optim/AbstractOptimizer.scala:279
1 | private[bigdl] def prepareInput[T: ClassTag](dataset: DataSet[MiniBatch[T]], |
上記の DistributedDataSet の chache
メソッドは以下の通り。
com/intel/analytics/bigdl/dataset/DataSet.scala:216
1 | def cache(): Unit = { |
originRDD の戻り値に対して、count
を読んでいる。 ここで count を呼ぶのは、入力データである
originRDD
の戻り値に入っているRDDをメモリ上にマテリアライズするためである。
count
を呼ぶだけでマテリアライズできるのは、予め入力データを定義したときに
Spark RDDの cache
を利用してキャッシュ化することを指定されているからである。
今回の例では、Optimizerオブジェクトのapplyを利用する際に渡されるデータセット
trainSet を 定義する際に予め cache
が呼ばれる。
com/intel/analytics/bigdl/models/lenet/Train.scala:79
1 | val trainSet = DataSet.array(load(trainData, trainLabel), sc) -> |
trainSet を定義する際、
com.intel.analytics.bigdl.dataset.DataSet$#array(T[], org.apache.spark.SparkContext)
メソッドが 呼ばれるのだが、その中で以下のように Sparkの
RDD#cache が呼ばれていることがわかる。
com/intel/analytics/bigdl/dataset/DataSet.scala:343
1 | def array[T: ClassTag](localData: Array[T], sc: SparkContext): DistributedDataSet[T] = { |
以下、一例。
具体的には、Optimizerのインスタンスを生成するための
apply メソッドはいくつかあるが、
以下のように引数にデータセットを指定する箇所がある。(再掲)
com/intel/analytics/bigdl/optim/Optimizer.scala:619
1 | _dataset = (DataSet.rdd(sampleRDD) -> |
ここで用いられている
com.intel.analytics.bigdl.dataset.DataSet#rdd
メソッドは以下の通り。
com/intel/analytics/bigdl/dataset/DataSet.scala:363
1 | def rdd[T: ClassTag](data: RDD[T], partitionNum: Int = Engine.nodeNumber() |
com.intel.analytics.bigdl.dataset.CachedDistriDataSet#CachedDistriDataSet
のコンストラクタ引数に、 org.apache.spark.rdd.RDD#cache
を用いてキャッシュ化することを指定したRDDを渡していることがわかる。