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

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

同じ計算を二度しないために — キャッシュが実行計画に現れる様子

TT Labで続きを見る

目標

同じDataFrameでアクションを2回呼ぶとき、キャッシュがなければ元データから再計算することをプランで確認し、cache・persist・unpersist・チェックポイント・CACHE TABLEが、プランとイベントログにどう現れるかを見ます。キャッシュが実際にメモリをどれだけ使うかも、ブロック更新イベントで測ります。

なぜ重要なのか

DataFrameは結果ではなく、作り方(リネージ)です。アクションを呼ぶたびに、Sparkはその作り方を最初からやり直します。同じ中間結果を2回使うコードは、元データを2回読んで、結合を2回行います。エラーも警告もなしに。 cache()は、その中間結果をエグゼキューターに残しておくようにという印です。印にすぎないので、最初のアクションが埋め、そのあとのアクションは、プランで元データを読む代わりにInMemoryTableScanでキャッシュを読みます。ストレージレベル(persist)は、メモリ・ディスク・シリアライズの有無を選ぶもので、メモリが足りないとき何をあきらめるかを決めます。使い終わったキャッシュはunpersistで解放して、はじめてほかの作業のメモリになります。 リネージが長くなりすぎると(反復計算)、プラン自体が大きくなってドライバーが遅くなります。チェックポイントは、結果をファイルに書いてリネージを切り、そのあとのプランを「このファイルを読む」の1行にします。

ステップ

  1. /root/spk/cache/common.pyに、決済完了の注文に金額を付ける関数と、アクションを2回(count、sum(amount))呼ぶ関数を置き、/root/spk/cache/none.py(アプリspk-cache-none)で、キャッシュなしで2つのアクションの結果を、/root/spk/cache/out/none.jsonに書き込んでください。
  2. /root/spk/cache/cached.py(アプリspk-cache-on)で、同じDataFrameにcache()をかけ、2つのアクションの結果を、/root/spk/cache/out/cached.jsonに、ストレージレベル(str(df.storageLevel))を、/root/spk/cache/out/cached_level.txtに書き込んでください。
  3. /root/spk/cache/disk.py(アプリspk-cache-disk)で、persist(StorageLevel.DISK_ONLY)を使って同じことを行い、結果を、/root/spk/cache/out/disk.jsonに、ストレージレベルを、/root/spk/cache/out/level.txtに書き込んでください。
  4. /root/spk/cache/unpersist.py(アプリspk-cache-un)で、キャッシュ → アクション → unpersist() → アクションの順に実行し、キャッシュの有無(df.is_cached)の前後を、/root/spk/cache/out/unpersist.jsonに{"before": 참거짓, "after": 참거짓}の形式で書き込んでください(プレースホルダーは順に真偽値、真偽値です)。
  5. /root/spk/cache/ckpt.py(アプリspk-cache-ckpt)で、チェックポイントフォルダーを、/root/spk/cache/ckptに設定し、金額に1–10を順に足す10回のwithColumnのあとにcheckpoint()した結果で、チャネル別の金額の合計を、/root/spk/cache/out/ckpt_resultにParquet(channel・amount)で書き込んでください。
  6. /root/spk/cache/cache_sql.py(アプリspk-cache-sql)で、一時ビューpaidを作成し、CACHE TABLE paid_cached AS SELECT channel, amount FROM paidを実行してから、チャネル別の合計を求めて、/root/spk/cache/out/cache_table.jsonに{"is_cached": 참거짓, "by_channel": {채널: 합}}の形式で書き込んでください(プレースホルダーは順に真偽値、チャネル、合計です)。
  7. /root/spk/cache/size.py(アプリspk-cache-size、spark.eventLog.logBlockUpdates.enabled=true)でキャッシュを埋めたあと、そのログのrdd_ブロックのサイズを足して、/root/spk/cache/out/cache_size.jsonに{"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数、整数です)。
  8. /root/spk/cache/report.mdに、## 다시 계산하는 비용・## 저장 수준・## 캐시의 크기の3つの節を書いてください(見出しは韓国語で、順に「再計算のコスト」「ストレージレベル」「キャッシュのサイズ」を意味します)。3つ目の節には、ステップ7のブロック数とメモリバイト数を入れてください。

参考

キャッシュなしでアクション2回(元データを2回読む)

/root/spk/cache/common.pyに、決済完了の注文を商品と結合してamount(qty×price)を付ける関数と、DataFrameにcount()とsum(amount)の2つのアクションを呼んで{"count": 정수, "sum": 정수}をファイルに書き込む関数を置いてください(プレースホルダーは順に整数、整数です)。/root/spk/cache/none.pyを、アプリ名spk-cache-noneで作成し、キャッシュなしで結果を、/root/spk/cache/out/none.jsonに書き込んでください。

アクションが2つなら、SQL実行も2つで、どちらもプランの一番下にScan csvがあります。同じファイルを2回読み、同じ結合を2回行うという意味です。採点ツールは、キャッシュなしで元データを読んだ実行が2つ以上あるかを確認します。

cache()(最初のアクションが埋め、次のアクションが読む)

/root/spk/cache/cached.pyを、アプリ名spk-cache-onで作成し、ステップ1のDataFrameにcache()をかけて、同じ2つのアクションの結果を、/root/spk/cache/out/cached.jsonに、str(df.storageLevel)を、/root/spk/cache/out/cached_level.txtに書き込んでください。

2回目の実行のプランにInMemoryTableScanが現れます。DataFrameのデフォルトのストレージレベルは、RDDのcacheとは違うという点も見てください。メモリが足りなければ、ディスクに流します。結果の数値はステップ1と同じでなければなりません。

persist(DISK_ONLY)(ストレージレベルを選ぶ)

/root/spk/cache/disk.pyを、アプリ名spk-cache-diskで作成し、ステップ1のDataFrameをpersist(StorageLevel.DISK_ONLY)で保持して、同じ2つのアクションの結果を、/root/spk/cache/out/disk.jsonに、str(df.storageLevel)を、/root/spk/cache/out/level.txtに書き込んでください。

ディスクにだけ置くと、エグゼキューターのメモリは節約できますが、読むたびにデシリアライズが必要です。プランのInMemoryRelationの行に出力されたStorageLevelが変わったことを見てください。名前がInMemoryでも、ディスクにあることがあります。

unpersist(キャッシュを解放すると再計算する)

/root/spk/cache/unpersist.pyを、アプリ名spk-cache-unで作成し、ステップ1のDataFrameを、cache() → count() → unpersist() → sum(amount)の順に実行して、unpersist()の前後のdf.is_cachedを、/root/spk/cache/out/unpersist.jsonに{"before": 참거짓, "after": 참거짓}の形式で書き込んでください(プレースホルダーは順に真偽値、真偽値です)。

キャッシュは無料ではありません。エグゼキューターのメモリを占め、ほかの作業が使える場所を減らします。使い終わったキャッシュは解放してください。採点ツールは、キャッシュを読んだ実行のあとに、キャッシュを読まない実行が続くかをログで確認します。

チェックポイントでリネージを切る

/root/spk/cache/ckpt.pyを、アプリ名spk-cache-ckptで作成し、チェックポイントフォルダーを、/root/spk/cache/ckptに決め、ステップ1のDataFrameでamountに1, 2, …, 10を順に足すwithColumnを10回行ったあと、checkpoint()した結果で、チャネル別のamountの合計を、/root/spk/cache/out/ckpt_resultにParquet(channel・amount)で書き込んでください。

checkpoint()は、その場で結果をチェックポイントフォルダーに書き、返されたDataFrameのプランは、元データではなくそのファイルを読む1行(Scan ExistingRDD)になります。キャッシュはリネージを残し(失えば再計算)、チェックポイントはリネージを切ります(ファイルがそのまま元データ)。

SQLのCACHE TABLEは遅延評価ではない

/root/spk/cache/cache_sql.pyを、アプリ名spk-cache-sqlで作成し、ステップ1のDataFrameを一時ビューpaidとして登録して、CACHE TABLE paid_cached AS SELECT channel, amount FROM paidを実行したあと、paid_cachedからチャネル別のamountの合計を求めて、/root/spk/cache/out/cache_table.jsonに{"is_cached": spark.catalog.isCached("paid_cached"), "by_channel": {채널: 합}}の形式で書き込んでください(プレースホルダーは順にチャネル、合計です)。

DataFrameのcache()は印を付けるだけですが、CACHE TABLEはその場で埋めます(CACHE LAZY TABLEなら先送りします)。そのため、CACHE TABLE文自体がジョブを起動します。名前付きのキャッシュを読むプランには、Scan In-memory table paid_cachedが出力されます。採点ツールは、そのような実行があるかと、合計を確認します。

キャッシュが実際に使うメモリを測る

/root/spk/cache/size.pyを、アプリ名spk-cache-size、設定spark.eventLog.logBlockUpdates.enabled=trueで作成し、ステップ1のDataFrameをcache()してcount()で埋めてください。そのアプリのログで、SparkListenerBlockUpdatedのうちブロックIDがrdd_で始まるものを、ブロックごとに最後のものだけ残して、個数とMemory Size・Disk Sizeの合計を、/root/spk/cache/out/cache_size.jsonに{"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数、整数です)。

キャッシュブロックは、パーティションごとに1つ、rdd_<RDD 번호>_<파티션>の形式です(プレースホルダーは順にRDD番号、パーティションです)。メモリに展開した(デシリアライズした)キャッシュは、元のCSVより大きいことも小さいこともあります。選んだカラムと型によって違います。キャッシュを増やす前にこの数値を見る習慣が、エグゼキューターのメモリを守ります。

キャッシュを付ける場所と解放する場所

/root/spk/cache/report.mdに、## 다시 계산하는 비용・## 저장 수준・## 캐시의 크기の3つの節を書いてください(見出しは韓国語で、順に「再計算のコスト」「ストレージレベル」「キャッシュのサイズ」を意味します)。2つ目の節にはステップ2・3のストレージレベルの文字列を、3つ目の節にはステップ7のブロック数とメモリバイト数を入れてください。

最初の節には、キャッシュがないときに何を2回したか、2つ目の節には、2つのストレージレベルが何を引き換えにしているか、3つ目の節には、そのメモリがエグゼキューターのヒープ(1g)の何パーセントかを書けばよいです。