TT Lab
はじめる
学ぶ 学習パス コース

Apache Spark — 遅いジョブの答えは実行計画とイベントログにある

DataFrame は結果ではなくレシピなので毎回作り直される

TT Labで続きを見る

一言でいうと

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タブも、実体化されたあとにしかキャッシュを表示しません。

実務で本当に大切なこと

次のラボですること

キャッシュなしで2つのアクションを呼び、2つのジョブがどちらもCSVを読み直すことを、イベントログとプランで確認します。cache()のあとは、2つ目のアクションのプランがInMemoryTableScanから始まることとストレージレベルを見て、persistでディスクだけに置くレベルを指定してみます。unpersistの前後でキャッシュの有無が変わることを確認し、トランスフォーメーションを10回重ねたDataFrameをチェックポイントで切って、同じ答えが出るかを見たあと、SQLのCACHE TABLEで同じことをしてみます。最後に、ブロック更新の記録をオンにしたアプリのイベントログで、キャッシュブロックのサイズを足して、キャッシュがメモリをどれだけ占めたかを数値で測ります。