Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
パーティション化したレイクを書き、上書きし、小さなファイルを減らす
한국어 원문으로 표시합니다.
목표
주문을 날짜로 나눈 Parquet 레이크로 쓰고, 저장 모드 넷(errorifexists·append·overwrite·ignore) 가운데 셋이 실제로 무엇을 하는지 확인한다. 정적 덮어쓰기와 동적 파티션 덮어쓰기의 차이를 디렉터리 수로 보고, 작은 파일을 줄이는 두 가지 손잡이(파티션 칼럼으로 모으기·파일당 행 수 상한)를 써 본다.
왜 중요한가
Spark 의 쓰기는 한 번에 끝나는 일이 아니다. 태스크마다 임시 경로에 쓰고, 잡이 성공하면 커밋 프로토콜이 결과를 제자리로 옮기고 _SUCCESS 표식을 남긴다. 그래서 반쯤 쓴 결과가 읽히지 않는다. 대신 저장 모드를 잘못 고르면 이 정직한 기계가 정직하게 사고를 낸다.
append 는 재시도 한 번에 같은 자료를 두 번 넣는다. overwrite 는 기본이 정적이라, 하루치만 고치려고 쓴 것이 대상 경로의 모든 날짜를 지운다. 동적 파티션 덮어쓰기(spark.sql.sources.partitionOverwriteMode=dynamic)는 쓰는 자료에 들어 있는 파티션만 바꾼다. 이 차이를 모르고 운영 레이크에 overwrite 를 한 번 쓰면 몇 달치가 한 번에 사라진다.
파일 수는 쓰는 태스크 수 × 그 태스크가 가진 파티션 값의 수다. 셔플 뒤 태스크가 모든 날짜를 조금씩 가지고 있으면 날짜 디렉터리마다 태스크 수만큼 작은 파일이 생긴다. 파티션 칼럼으로 먼저 모으거나(repartition("칼럼")), 파일당 행 수에 상한(maxRecordsPerFile)을 두면 크기를 다스릴 수 있다.
단계
- /root/spk/write/common.py 에 주문을 읽어
order_date를 붙이는 함수를 두고, /root/spk/write/lake.py(앱spk-write-lake)로 /root/spk/write/lake/orders 에order_date로 나눈 Parquet 을 쓰세요. - /root/spk/write/exists.py(앱
spk-write-exists)로 같은 경로에 모드를 주지 않고 다시 쓰게 하고, 오류 조건 이름을 /root/spk/write/out/exists.txt 에 쓰세요. - /root/spk/write/append.py(앱
spk-write-append)로 1월 주문을 /root/spk/write/lake/append 에 overwrite 로 한 번, append 로 한 번 더 쓰고, 다시 읽은 행 수와 고유order_id수를 /root/spk/write/out/append.json 에 쓰세요. - /root/spk/write/static.py(앱
spk-write-static)로 전체를 /root/spk/write/lake/static 에 쓴 뒤 2026-02-10 하루치만 overwrite 로 다시 쓰고, 남은 파티션 디렉터리 수를 /root/spk/write/out/static.json 에 쓰세요. - /root/spk/write/dynamic.py(앱
spk-write-dynamic,partitionOverwriteMode=dynamic)로 전체를 /root/spk/write/lake/dynamic 에 쓴 뒤 2026-02-10 에서 취소 주문을 뺀 것만 overwrite 로 쓰고, 남은 파티션 디렉터리 수를 /root/spk/write/out/dynamic.json 에 쓰세요. - /root/spk/write/small.py(앱
spk-write-small)로 클릭을repartition(40)해서 /root/spk/write/lake/clicks_many 에,repartition("page")해서 /root/spk/write/lake/clicks_few 에 각각page로 나눠 쓰고, 두 레이크의 데이터 파일 수를 /root/spk/write/out/files.json 에 쓰세요. - /root/spk/write/maxrec.py(앱
spk-write-maxrec)로 주문을repartition("channel")한 뒤maxRecordsPerFile=20000으로 /root/spk/write/lake/by_channel 에channel로 나눠 쓰세요. - /root/spk/write/report.md 에
## 저장 모드## 정적과 동적 덮어쓰기## 파일 수세 절을 쓰세요. 둘째 절에는 4·5단계의 파티션 수를, 셋째 절에는 6단계의 두 파일 수를 넣으세요.
참고
- 스크립트는
/root/spk/write에 두고 거기서 돌리세요(from common import …). - 레이크의 데이터 파일은
part-로 시작합니다._SUCCESS는 커밋 표식이고.crc는 체크섬입니다(세지 않습니다). - 동적 덮어쓰기는 세션 설정(
spark.sql.sources.partitionOverwriteMode)이나 쓰기 옵션(.option("partitionOverwriteMode", "dynamic"))으로 켭니다. - 흔한 실수: 한 경로를 여러 단계에서 고쳐 써 앞 단계 결과를 지우는 것(단계마다 경로가 다릅니다), 정적 덮어쓰기로 하루만 바뀔 거라 믿는 것, maxRecordsPerFile 이 파일 크기(바이트)를 정한다고 여기는 것.
- 공식 문서: Generic Load/Save Functions — Save Modes · Bucketing, Sorting and Partitioning · Parquet Files · Configuration — spark.sql.sources.partitionOverwriteMode
날짜로 나눈 레이크 쓰기
/root/spk/write/common.py 에 주문을 스키마로 읽고 order_date = to_date(order_ts) 를 붙이는 함수를 두고, /root/spk/write/lake.py 를 앱 이름 spk-write-lake 로 만들어 /root/spk/write/lake/orders 에 partitionBy("order_date") 로 Parquet 을 쓰세요(mode("overwrite")).
날짜 90일이면 디렉터리 90개입니다. 잡이 성공해야 _SUCCESS 가 생깁니다 — 이 표식이 없는 디렉터리는 반쯤 쓴 것일 수 있으니 읽는 쪽이 먼저 확인합니다. 채점기는 디렉터리 수·행 수·표식을 봅니다.
기본 모드는 멈춘다 — errorifexists
/root/spk/write/exists.py 를 앱 이름 spk-write-exists 로 만들어 1단계와 같은 경로에 모드를 주지 않고 다시 쓰게 하고, 예외의 getCondition() 을 /root/spk/write/out/exists.txt 첫 줄에 쓰세요. 1단계의 레이크는 그대로 남아 있어야 합니다.
기본 저장 모드는 '이미 있으면 오류' 입니다. 실수로 같은 경로에 쓰는 것을 막는 안전장치이고, 운영 잡에서 모드를 명시하는 습관이 필요한 이유이기도 합니다.
append — 재시도가 두 배를 만든다
/root/spk/write/append.py 를 앱 이름 spk-write-append 로 만들어 2026년 1월 주문을 /root/spk/write/lake/append 에 overwrite 로 한 번, append 로 한 번 더 쓰고, 다시 읽은 행 수와 고유 order_id 수를 /root/spk/write/out/append.json 에 {"rows": 정수, "distinct_orders": 정수} 로 쓰세요.
append 는 이미 있는 파일을 보지 않고 새 파일을 더합니다. 같은 입력으로 잡을 재시도하면 행이 두 배가 되고, 파일 이름이 달라 겉으로는 멀쩡해 보입니다. 멱등하게 쓰려면 덮어쓰기(가능하면 동적)나 키로 합치는 방식을 씁니다.
정적 덮어쓰기 — 하루를 고치려다 전부 지운다
/root/spk/write/static.py 를 앱 이름 spk-write-static 으로 만들어 전체 주문을 /root/spk/write/lake/static 에 order_date 로 나눠 쓴 뒤, 2026-02-10 하루치만 같은 경로에 overwrite(설정 기본값)로 다시 쓰세요. 남은 order_date= 디렉터리 수를 /root/spk/write/out/static.json 에 {"partitions_after": 정수} 로 쓰세요.
partitionOverwriteMode 의 기본은 static 입니다. overwrite 는 쓰기 전에 대상 경로 아래를 통째로 지웁니다 — 쓰는 자료에 무슨 날짜가 들었는지는 보지 않습니다. 이 단계의 결과는 사고입니다. 그것을 눈으로 보는 것이 목적입니다.
동적 파티션 덮어쓰기 — 그날만 바꾸기
/root/spk/write/dynamic.py 를 앱 이름 spk-write-dynamic, 설정 spark.sql.sources.partitionOverwriteMode=dynamic 으로 만들어 전체 주문을 /root/spk/write/lake/dynamic 에 order_date 로 나눠 쓴 뒤, 2026-02-10 주문에서 취소(cancelled)를 뺀 것만 같은 경로에 overwrite 로 쓰세요. 남은 order_date= 디렉터리 수를 /root/spk/write/out/dynamic.json 에 {"partitions_after": 정수} 로 쓰세요.
동적 모드에서는 쓰는 자료에 들어 있는 파티션(여기서는 하루)만 갈아 끼웁니다. 나머지 89일은 그대로입니다. 채점기는 디렉터리 수, 그날에 취소 주문이 없는지, 다른 날은 원본 그대로인지를 봅니다.
작은 파일 — 태스크 수 × 파티션 값
/root/spk/write/small.py 를 앱 이름 spk-write-small 로 만들어 /data/clicks/clicks.jsonl 을 repartition(40) 해서 /root/spk/write/lake/clicks_many 에, repartition("page") 해서 /root/spk/write/lake/clicks_few 에 각각 partitionBy("page") 로 쓰고, 두 레이크의 데이터 파일(page=*/part-*) 수를 /root/spk/write/out/files.json 에 {"many": 정수, "few": 정수} 로 쓰세요.
쓰는 태스크 하나는 자기가 가진 파티션 값마다 파일을 하나씩 엽니다. 태스크 40개가 페이지 일곱 개를 다 가지고 있으면 최대 280개입니다. 파티션 칼럼으로 먼저 모으면 페이지 하나가 태스크 하나에 들어가 페이지당 파일 하나가 됩니다(대신 한 페이지가 크면 그 태스크가 무거워집니다).
파일당 행 수에 상한 두기
/root/spk/write/maxrec.py 를 앱 이름 spk-write-maxrec 으로 만들어 주문을 repartition("channel") 한 뒤 option("maxRecordsPerFile", 20000) 으로 /root/spk/write/lake/by_channel 에 partitionBy("channel") 로 쓰세요.
채널로 모으면 채널 하나가 태스크 하나에 들어가 파일 하나가 되는데, 가장 큰 채널은 17만 행이 넘습니다. 상한을 두면 태스크가 2만 행마다 파일을 새로 엽니다. 채점기는 채널마다 파일 수가 올림(행 수 ÷ 20000)인지와 파일마다 행이 2만을 넘지 않는지를 Parquet 메타데이터로 봅니다.
쓰기 규칙을 팀 규칙으로
/root/spk/write/report.md 에 ## 저장 모드 ## 정적과 동적 덮어쓰기 ## 파일 수 세 절을 쓰세요. 둘째 절에는 4·5단계의 남은 파티션 수를, 셋째 절에는 6단계의 두 파일 수를 넣으세요.
운영 레이크에 쓰는 잡이라면 어떤 모드를 기본으로 하고 무엇을 금지할지를 규칙으로 적어 보세요. 숫자는 여러분의 결과 파일에서 옮기세요.