LabHub
배우기 러닝패스 코스

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

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

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

같은 DataFrame 에 행동을 두 번 부를 때 캐시가 없으면 원본부터 다시 계산하는 것을 계획으로 확인하고, cache·persist·unpersist·체크포인트·CACHE TABLE 이 계획과 이벤트 로그에 어떻게 나타나는지 본다. 캐시가 실제로 메모리를 얼마나 먹는지도 블록 갱신 이벤트로 잰다.

왜 중요한가

DataFrame 은 결과가 아니라 만드는 방법(계보)이다. 행동을 부를 때마다 Spark 는 그 방법을 처음부터 다시 실행한다. 같은 중간 결과를 두 번 쓰는 코드는 원본을 두 번 읽고 조인을 두 번 한다 — 오류도 경고도 없이. cache() 는 그 중간 결과를 실행기에 남겨 두라는 표시다. 표시일 뿐이라 첫 행동이 채우고, 그 뒤의 행동은 계획에서 원본을 읽는 대신 InMemoryTableScan 으로 캐시를 읽는다. 저장 수준(persist)은 메모리·디스크·직렬화 여부를 고르는 것이고, 메모리가 모자라면 무엇을 포기할지를 정한다. 다 쓴 캐시는 unpersist 로 내려놓아야 다른 작업의 메모리가 된다. 계보가 너무 길어지면(반복 계산) 계획 자체가 커져 드라이버가 느려진다. 체크포인트는 결과를 파일로 쓰고 계보를 끊어, 그 뒤의 계획을 '이 파일을 읽는다' 한 줄로 만든다.

단계

  1. /root/spk/cache/common.py 에 결제 완료 주문에 금액을 붙이는 함수와 행동 두 번(count, sum(amount))을 부르는 함수를 두고, /root/spk/cache/none.py(앱 spk-cache-none)로 캐시 없이 두 행동의 결과를 /root/spk/cache/out/none.json 에 쓰세요.
  2. /root/spk/cache/cached.py(앱 spk-cache-on)에서 같은 DataFrame 에 cache() 를 걸고 두 행동의 결과를 /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 을 차례로 더하는 열 번의 withColumncheckpoint() 한 결과로 채널별 금액 합을 /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## 다시 계산하는 비용 ## 저장 수준 ## 캐시의 크기 세 절을 쓰세요. 셋째 절에는 7단계의 블록 수와 메모리 바이트를 넣으세요.

참고

캐시 없이 행동 둘 — 원본을 두 번 읽는다

/root/spk/cache/common.py 에 결제 완료 주문을 상품과 조인해 amount(qty×price)를 붙이는 함수와, DataFrame 에 count()sum(amount) 두 행동을 부르고 {"count": 정수, "sum": 정수} 를 파일로 쓰는 함수를 두세요. /root/spk/cache/none.py 를 앱 이름 spk-cache-none 으로 만들어 캐시 없이 결과를 /root/spk/cache/out/none.json 에 쓰세요.

행동이 둘이면 SQL 실행도 둘이고, 둘 다 계획 맨 아래에 Scan csv 가 있습니다. 같은 파일을 두 번 읽고 같은 조인을 두 번 한다는 뜻입니다. 채점기는 캐시 없이 원본을 읽은 실행이 둘 이상인지 봅니다.

cache() — 첫 행동이 채우고 다음 행동이 읽는다

/root/spk/cache/cached.py 를 앱 이름 spk-cache-on 으로 만들어 1단계의 DataFrame 에 cache() 를 걸고 같은 두 행동의 결과를 /root/spk/cache/out/cached.json 에, str(df.storageLevel)/root/spk/cache/out/cached_level.txt 에 쓰세요.

두 번째 실행의 계획에 InMemoryTableScan 이 나타납니다. DataFrame 의 기본 저장 수준은 RDD 의 cache 와 다르다는 점도 보세요 — 메모리가 모자라면 디스크로 흘립니다. 결과 숫자는 1단계와 같아야 합니다.

persist(DISK_ONLY) — 저장 수준 고르기

/root/spk/cache/disk.py 를 앱 이름 spk-cache-disk 로 만들어 1단계의 DataFrame 을 persist(StorageLevel.DISK_ONLY) 로 두고 같은 두 행동의 결과를 /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 열 번 뒤 checkpoint() 한 결과로 채널별 amount 합을 /root/spk/cache/out/ckpt_result 에 Parquet(channel,amount)으로 쓰세요.

checkpoint() 는 그 자리에서 결과를 체크포인트 폴더에 쓰고, 돌려준 DataFrame 의 계획은 원본이 아니라 그 파일을 읽는 한 줄(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": 정수} 로 쓰세요.

캐시 블록은 파티션마다 하나(rdd_<RDD 번호>_<파티션>)입니다. 메모리에 풀어 둔(역직렬화) 캐시는 원본 CSV 보다 클 수도 작을 수도 있습니다 — 고른 칼럼과 형에 따라 다릅니다. 캐시를 늘리기 전에 이 숫자를 보는 습관이 실행기 메모리를 지킵니다.

캐시를 걸 자리와 내려놓을 자리

/root/spk/cache/report.md## 다시 계산하는 비용 ## 저장 수준 ## 캐시의 크기 세 절을 쓰세요. 둘째 절에는 2·3단계의 저장 수준 문자열을, 셋째 절에는 7단계의 블록 수와 메모리 바이트를 넣으세요.

첫 절에는 캐시가 없을 때 무엇을 두 번 했는지, 둘째 절에는 두 저장 수준이 무엇을 맞바꾸는지, 셋째 절에는 그 메모리가 실행기 힙(1g)의 몇 퍼센트인지를 적으면 됩니다.