Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 스키마와 깨진 줄 · 실습
협력사 CSV 의 깨진 줄을 세 가지 모드로 받아 낸다
목표
협력사가 보낸 주문 파일 6,000줄을 스키마를 정해 읽고, 깨진 줄을 PERMISSIVE·DROPMALFORMED·FAILFAST 세 모드가 각각 어떻게 다루는지 숫자로 확인한 뒤, 업무 규칙까지 더해 깨끗한 표와 격리 표로 나눈다.
왜 중요한가
CSV 에는 형이 없다. 칼럼 하나에 "two" 가 한 번 섞이면 inferSchema 는 그 칼럼 전체를 문자열로 물러서고, 그 사실은 오류가 아니라 조용한 형 변화로 나타난다. 어제까지 정수였던 칼럼이 오늘 문자열이 되면 뒤의 합계가 전부 깨진다. 그래서 파이프라인은 스키마를 추론하지 않고 정해서 읽는다.
스키마를 정하면 이번에는 스키마에 맞지 않는 줄을 어떻게 할지 정해야 한다. Spark 의 CSV 읽기에는 모드가 셋 있다. PERMISSIVE 는 깨진 줄을 남기고 원문을 손상 칼럼에 담고, DROPMALFORMED 는 버리고, FAILFAST 는 멈춘다. 어느 것이 옳은지는 업무가 정한다 — 돈을 세는 표라면 조용히 버리는 것이 가장 위험하다.
함정이 하나 있다. CSV 파서는 필요한 칼럼만 파싱한다. count() 처럼 칼럼이 하나도 필요 없는 행동은 줄을 파싱하지 않으므로, DROPMALFORMED 로 읽고 세면 깨진 줄까지 전부 세고 FAILFAST 도 멈추지 않는다. 검증을 count() 로 했다면 아무것도 검증하지 않은 것이다.
단계
1. 원본의 일곱 칼럼을 DDL 로 적어 /root/spk/schema/schema.ddl 에 저장하세요(qty 는 INT, order_ts 는 TIMESTAMP, 나머지는 STRING).
2. /root/spk/schema/infer.py 를 앱 이름 spk-schema-infer 로 만들어 inferSchema=True 로 읽고, 추론된 스키마의 simpleString() 을 /root/spk/schema/out/inferred.txt 에 쓰세요.
3. /root/spk/schema/load.py 를 앱 이름 spk-schema-load 로 만들어 1단계 스키마 끝에 _corrupt STRING 을 더해 PERMISSIVE 로 읽고, 전체를 /root/spk/schema/out/permissive 에 Parquet 으로 쓰세요.
4. /root/spk/schema/dropmal.py 를 앱 이름 spk-schema-drop 으로 만들어 DROPMALFORMED 로 읽고, count() 한 값과 /root/spk/schema/out/dropped 에 Parquet 으로 쓴 뒤 다시 센 값을 /root/spk/schema/out/drop_count.json 에 {"count_only": 정수, "written": 정수} 로 쓰세요.
5. /root/spk/schema/failfast.py 를 앱 이름 spk-schema-failfast 로 만들어 FAILFAST 로 읽은 것을 Parquet 으로 쓰게 하고, 실패의 원인이 된 파싱 오류 조건 이름을 /root/spk/schema/out/failfast.txt 첫 줄에 쓰세요.
6. /root/spk/schema/width.py 를 앱 이름 spk-schema-width 로 만들어 원본을 줄 단위로 읽고(머리줄 제외) 쉼표로 나눈 칸 수별 줄 수를 /root/spk/schema/out/width.json 에 {"칸 수": 줄 수} 로 쓰세요.
7. /root/spk/schema/clean.py 를 앱 이름 spk-schema-clean 으로 만들어 파싱이 된 줄 가운데 qty 가 1 이상인 것만 /root/spk/schema/out/clean 에, 나머지는 reason 칼럼(파싱 실패 parse, 수량 문제 qty)을 붙여 /root/spk/schema/out/quarantine 에 Parquet 으로 쓰세요.
8. /root/spk/schema/report.md 에 ## 추론이 틀린 곳 ## 세 가지 모드 ## 격리한 줄 세 절을 쓰세요. 둘째 절에는 손상 줄 수와 DROPMALFORMED 로 쓴 줄 수를, 셋째 절에는 깨끗한 줄 수와 격리한 줄 수를 숫자로 넣으세요.
참고
- 원본:
/data/shop/orders_dirty.csv(머리줄 1줄 + 6,000줄). 시각 형식은yyyy-MM-dd HH:mm:ss이므로timestampFormat으로 알려 주세요. - 손상 칼럼은 스키마에 직접 넣어야 합니다. 이름은
columnNameOfCorruptRecord옵션과 같아야 합니다. - 손상 칼럼만 골라 묻는 질의는 원본을 다시 파싱해야 해서 막혀 있을 수 있습니다. 한 번 읽은 것을
cache()해 두고 거르세요. - FAILFAST 의 예외는 파이썬에서 조건 이름을 바로 돌려주지 않을 수 있습니다. 메시지 안의 대괄호
[...]에 조건 이름이 들어 있습니다. - 흔한 실수:
count()로 모드의 효과를 확인하는 것, 손상 칼럼을 스키마에 넣지 않는 것, 빈 수량(null)을 깨끗한 줄로 보내는 것. - 공식 문서: [CSV Files](https://spark.apache.org/docs/4.2.0/sql-data-sources-csv.html) · [Data Types](https://spark.apache.org/docs/4.2.0/sql-ref-datatypes.html) · [Datetime Patterns](https://spark.apache.org/docs/4.2.0/sql-ref-datetime-pattern.html) · [Error Conditions](https://spark.apache.org/docs/4.2.0/sql-error-conditions.html)
8단계
- 스키마를 정해 두기
- 추론이 무엇을 놓치는지 보기
- PERMISSIVE — 버리지 않고 원문을 담기
- DROPMALFORMED — count() 는 거짓말을 한다
- FAILFAST — 멈추게 하고 이유 읽기
- 칸 수가 다른 줄 세기
- 업무 규칙까지 — 깨끗한 표와 격리 표
- 무엇을 버렸고 왜 버렸는지 남기기