Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 파일로 쓰기 · 实验
파티션 레이크를 쓰고, 덮어쓰고, 작은 파일을 줄인다
목표
주문을 날짜로 나눈 Parquet 레이크로 쓰고, 저장 모드 넷(errorifexists·append·overwrite·ignore) 가운데 셋이 실제로 무엇을 하는지 확인한다. 정적 덮어쓰기와 동적 파티션 덮어쓰기의 차이를 디렉터리 수로 보고, 작은 파일을 줄이는 두 가지 손잡이(파티션 칼럼으로 모으기·파일당 행 수 상한)를 써 본다.
왜 중요한가
Spark 의 쓰기는 한 번에 끝나는 일이 아니다. 태스크마다 임시 경로에 쓰고, 잡이 성공하면 커밋 프로토콜이 결과를 제자리로 옮기고 _SUCCESS 표식을 남긴다. 그래서 반쯤 쓴 결과가 읽히지 않는다. 대신 저장 모드를 잘못 고르면 이 정직한 기계가 정직하게 사고를 낸다.
append 는 재시도 한 번에 같은 자료를 두 번 넣는다. overwrite 는 기본이 정적이라, 하루치만 고치려고 쓴 것이 대상 경로의 모든 날짜를 지운다. 동적 파티션 덮어쓰기(spark.sql.sources.partitionOverwriteMode=dynamic)는 쓰는 자료에 들어 있는 파티션만 바꾼다. 이 차이를 모르고 운영 레이크에 overwrite 를 한 번 쓰면 몇 달치가 한 번에 사라진다.
파일 수는 쓰는 태스크 수 × 그 태스크가 가진 파티션 값의 수다. 셔플 뒤 태스크가 모든 날짜를 조금씩 가지고 있으면 날짜 디렉터리마다 태스크 수만큼 작은 파일이 생긴다. 파티션 칼럼으로 먼저 모으거나(repartition("칼럼")), 파일당 행 수에 상한(maxRecordsPerFile)을 두면 크기를 다스릴 수 있다.
단계
1. /root/spk/write/common.py 에 주문을 읽어 order_date 를 붙이는 함수를 두고, /root/spk/write/lake.py(앱 spk-write-lake)로 /root/spk/write/lake/orders 에 order_date 로 나눈 Parquet 을 쓰세요.
2. /root/spk/write/exists.py(앱 spk-write-exists)로 같은 경로에 모드를 주지 않고 다시 쓰게 하고, 오류 조건 이름을 /root/spk/write/out/exists.txt 에 쓰세요.
3. /root/spk/write/append.py(앱 spk-write-append)로 1월 주문을 /root/spk/write/lake/append 에 overwrite 로 한 번, append 로 한 번 더 쓰고, 다시 읽은 행 수와 고유 order_id 수를 /root/spk/write/out/append.json 에 쓰세요.
4. /root/spk/write/static.py(앱 spk-write-static)로 전체를 /root/spk/write/lake/static 에 쓴 뒤 2026-02-10 하루치만 overwrite 로 다시 쓰고, 남은 파티션 디렉터리 수를 /root/spk/write/out/static.json 에 쓰세요.
5. /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 에 쓰세요.
6. /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 에 쓰세요.
7. /root/spk/write/maxrec.py(앱 spk-write-maxrec)로 주문을 repartition("channel") 한 뒤 maxRecordsPerFile=20000 으로 /root/spk/write/lake/by_channel 에 channel 로 나눠 쓰세요.
8. /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](https://spark.apache.org/docs/4.2.0/sql-data-sources-load-save-functions.html#save-modes) · [Bucketing, Sorting and Partitioning](https://spark.apache.org/docs/4.2.0/sql-data-sources-load-save-functions.html#bucketing-sorting-and-partitioning) · [Parquet Files](https://spark.apache.org/docs/4.2.0/sql-data-sources-parquet.html) · [Configuration — spark.sql.sources.partitionOverwriteMode](https://spark.apache.org/docs/4.2.0/configuration.html)
8个步骤
- 날짜로 나눈 레이크 쓰기
- 기본 모드는 멈춘다 — errorifexists
- append — 재시도가 두 배를 만든다
- 정적 덮어쓰기 — 하루를 고치려다 전부 지운다
- 동적 파티션 덮어쓰기 — 그날만 바꾸기
- 작은 파일 — 태스크 수 × 파티션 값
- 파일당 행 수에 상한 두기
- 쓰기 규칙을 팀 규칙으로