Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
DataFrame は結果ではなくレシピなので毎回作り直される
一言でいうと
DataFrameは計算された結果ではなく、元データから結果までの道筋であるリネージ(lineage)なので、アクションを呼ぶたびに元データから計算し直します。キャッシュは、最初の計算結果を保持しておくもので、チェックポイントは、結果をファイルに書いてリネージを切るものです。
なぜ同じ計算を2回するのか
CSVを読んで整形し、結合したDataFrame 1つでcount()を呼び、続けて同じものをファイルに書くとします。コードだけを見ると、整形は1回です。実際には2回起きます。イベントログを見ると、2つのジョブがどちらもCSVスキャンから始まっています。
理由は、トランスフォーメーションが遅延評価だからです。filterやjoinはプランに1行を加えるだけで、何も計算しません。アクションが呼ばれると、Sparkはそのプランを元データから最後まで実行し、結果をドライバーやファイルに渡したあと、中間結果は捨てます。次のアクションは、再び元データから始まります。これは欠陥ではなく設計です。中間結果をすべて保持するとメモリがすぐに埋まり、リネージさえあれば、失った断片はいつでも作り直せます。
そのため、何度も使う中間結果だけを選んで保持するのは、人間の役目です。それがキャッシュです。
どう動くのか
cache()も遅延評価です。呼んだ瞬間には、「この結果を最初に計算するときに保持せよ」という印が付くだけです。最初のアクションが動くとき、各パーティションが計算されながら保存され、そのあとのアクションは、プランで元データのスキャンの代わりにInMemoryTableScanから始まります。RDDプログラミングガイドの説明のとおり、最初に計算されるときに保存され、そのあとのアクションがそれを再利用します。
from pyspark import StorageLevel
clean = raw.filter("qty > 0").join(products, "product_id")
clean.cache() # 표시만 한다 — 아직 아무것도 저장되지 않았다
clean.count() # 첫 행동: 계산하면서 저장
clean.write.parquet("/tmp/out") # 계획이 InMemoryTableScan 에서 시작한다
clean.unpersist() # 다음 행동은 다시 원본부터
slim = raw.select("order_id", "qty").persist(StorageLevel.DISK_ONLY)
ストレージレベルが、どこに置くかを決めます。混同しやすいのがデフォルトです。RDDのcache()は、メモリにだけ置くMEMORY_ONLYです。一方、DataFrameのcache()はMEMORY_AND_DISK_DESERで、メモリに収まらないパーティションはディスクに置きます。ドキュメントは、このデフォルトが3.0でScalaに合わせて変わったと書いています。別のレベルが欲しければpersist()に渡しますが、すでにストレージレベルが決まっているDataFrameには、新しいレベルを指定できません。先にunpersist()する必要があります。
DataFrameのキャッシュは、行のままではなくカラム形式でメモリに置かれます。パフォーマンスチューニングのドキュメントは、キャッシュされたテーブルから必要なカラムだけを読み、圧縮を自動で調整すると説明しています。
どのレベルを選ぶかについてのRDDガイドの助言は短いです。メモリに無理なく収まるならデフォルトのままにし、計算が高価か、絞り込む量が多い場合でなければ、ディスクに流さないでください。再計算のほうが、ディスクから読むのと同じくらい速いことがあるからです。元データがローカルのCSV 1つで、トランスフォーメーションが軽ければ、DISK_ONLYのキャッシュはほとんど利点がありません。
キャッシュが置かれる場所(統合メモリ)
キャッシュは、空いている場所に置かれるわけではありません。チューニングガイドによると、シャッフル・結合・ソート・集計が使う実行メモリと、キャッシュが使うストレージメモリは、1つの領域を分け合います。その領域のサイズがspark.memory.fraction(デフォルト0.6)で、JVMヒープから300MiBを引いた残りに対する割合です。その中でspark.memory.storageFraction(デフォルト0.5)の分は、実行メモリがキャッシュを追い出せない取り分です。
このラボのドライバーメモリ1gで概算すると、(1024 - 300) × 0.6、つまり400MiB強を、実行とキャッシュが分け合います。大きなテーブルをキャッシュすると、その分だけ結合とソートが使える場所が減り、キャッシュがあふれると、RDDガイドのとおり、最も長く使われていないパーティションから(LRU)追い出されます。追い出されたパーティションは、次に必要になったときにリネージをたどって再計算されます。そのため、キャッシュは間違った答えを出しません。遅くなるだけです。
unpersistとチェックポイント
使い終わったキャッシュは、unpersist()で解放します。RDDガイドは、この呼び出しがデフォルトでは待たないと書いています。リソースが実際に解放されるまで待つには、blocking引数を渡します。解放したあとのアクションは、プランで再び元データのスキャンから始まります。もう1つ、DataFrameのキャッシュは、クラスターのすべてのセッションが共有します。同じプランを別のセッションがキャッシュしていれば、自分のクエリもそれを使います。
キャッシュはリネージをそのまま残します。失ったパーティションを作り直せるのも、そのおかげです。ところが、反復アルゴリズムのように同じDataFrameにトランスフォーメーションを何十回も重ねると、リネージそのものが問題になります。チェックポイントのドキュメントは、このような場合、プランが指数関数的に膨らむことがあり、チェックポイントが論理プランを切り落とすと説明しています。結果をSparkContext.setCheckpointDir()かspark.checkpoint.dirで決めたディレクトリにファイルとして書き、そのあとのプランは元データではなく、そのファイルから始まります。プランには、元データのスキャンの代わりに、すでにあるRDDを読むノードが現れます。デフォルトは即時(eager)実行と呼ばれ、呼んだ瞬間にジョブが動きます。
localCheckpoint()は、ファイルの代わりにエグゼキューターのキャッシュ機構に保存します。速いですが、ドキュメント自身が言うように信頼できません。エグゼキューターを失うと、リネージも切れているので、作り直す方法がありません。
SQLにも、同じ道具があります。CACHE TABLEは、DataFrameのcache()と違って、デフォルトが即時で、LAZYを付けると最初に使うときにキャッシュします。ストレージレベルを別に指定しなければ、MEMORY_AND_DISKです。
現場での姿
1つ目は、1回しか使わないものをキャッシュすることです。アクションが1つだけのDataFrameをキャッシュすると、保存するコストだけが増えます。キャッシュは、2回以上使う分岐点にだけ置きます。
2つ目は、キャッシュを保持したまま忘れることです。長いノートブックやサービスの中でキャッシュが溜まると、実行メモリが減って、シャッフルがディスクにあふれ始めます。使い終わったら解放します。
3つ目は、キャッシュしたつもりなのにプランにないことです。cache()のあとに新しいトランスフォーメーションを付けた別のDataFrameでアクションを呼ぶと、そのDataFrameのプランの中でキャッシュされた部分だけが再利用されます。キャッシュが本当に使われたかは、アクションのあとのexplain()でInMemoryTableScanを探して確認します。Web UIのStorageタブも、実体化されたあとにしかキャッシュを表示しません。
実務で本当に大切なこと
- DataFrameはレシピです。アクションごとに元データから再計算します。
- cache()は遅延評価です。最初のアクションが動くときに保存されます。
- DataFrameのデフォルトのストレージレベルはMEMORY_AND_DISK_DESER、RDDはMEMORY_ONLYです。
- キャッシュは実行メモリと同じ領域を分け合います。使い終わったらunpersistします。
- リネージが長すぎるときは、チェックポイントで切ります。キャッシュはリネージを切りません。
次のラボですること
キャッシュなしで2つのアクションを呼び、2つのジョブがどちらもCSVを読み直すことを、イベントログとプランで確認します。cache()のあとは、2つ目のアクションのプランがInMemoryTableScanから始まることとストレージレベルを見て、persistでディスクだけに置くレベルを指定してみます。unpersistの前後でキャッシュの有無が変わることを確認し、トランスフォーメーションを10回重ねたDataFrameをチェックポイントで切って、同じ答えが出るかを見たあと、SQLのCACHE TABLEで同じことをしてみます。最後に、ブロック更新の記録をオンにしたアプリのイベントログで、キャッシュブロックのサイズを足して、キャッシュがメモリをどれだけ占めたかを数値で測ります。