LabHub
배우기 러닝패스 코스

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

取引先 CSV の壊れた行を三つのモードで読み込む

LabHub 에서 이어서 보기

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

목표

협력사가 보낸 주문 파일 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 의 머리줄 순서대로 일곱 칼럼을 DDL 문자열 한 줄로 적어 /root/spk/schema/schema.ddl 에 저장하세요. qty 는 INT, order_ts 는 TIMESTAMP, 나머지는 STRING 입니다.

DDL 은 이름 형, 이름 형, ... 모양입니다. head -1 로 머리줄을 보고 순서를 그대로 따르세요 — CSV 는 이름이 아니라 위치로 칼럼을 맞춥니다.

추론이 무엇을 놓치는지 보기

/root/spk/schema/infer.py 를 앱 이름 spk-schema-infer 로 만들어 inferSchema=True·header=True 로 원본을 읽고, df.schema.simpleString()/root/spk/schema/out/inferred.txt 에 쓰세요.

추론은 자료를 한 번 더 훑는 잡입니다. 이벤트 로그에서 잡 수를 보면 스키마를 정했을 때보다 많습니다. 결과에서 qtyorder_ts 가 어떤 형으로 잡혔는지, 왜 그런지 원본을 grep 해 보세요.

PERMISSIVE — 버리지 않고 원문을 담기

/root/spk/schema/load.py 를 앱 이름 spk-schema-load 로 만들어 1단계 DDL 끝에 _corrupt STRING 을 더한 스키마로, mode=PERMISSIVE·columnNameOfCorruptRecord=_corrupt·timestampFormat=yyyy-MM-dd HH:mm:ss 로 읽어 전체를 /root/spk/schema/out/permissive 에 Parquet 으로 쓰세요.

PERMISSIVE 는 줄 수를 바꾸지 않습니다. 6,000줄이 그대로 남고, 형에 맞지 않거나 칸 수가 다른 줄은 원문이 _corrupt 에 들어가며 해당 칸은 null 이 됩니다. 채점기는 원본을 다시 읽어 깨진 줄을 세고 여러분의 _corrupt 와 견줍니다.

DROPMALFORMED — count() 는 거짓말을 한다

/root/spk/schema/dropmal.py 를 앱 이름 spk-schema-drop 으로 만들어 1단계 스키마와 mode=DROPMALFORMED 로 읽고, (가) count() 한 값, (나) /root/spk/schema/out/dropped 에 Parquet 으로 쓴 뒤 다시 읽어 센 값을 /root/spk/schema/out/drop_count.json{"count_only": 정수, "written": 정수} 로 쓰세요.

두 숫자가 다르게 나와야 정상입니다. count() 에는 칼럼이 하나도 필요 없어서 CSV 파서가 줄을 파싱하지 않고, 파싱하지 않으면 깨졌는지도 모릅니다. 모든 칼럼을 써야 하는 행동을 해야 비로소 버립니다.

FAILFAST — 멈추게 하고 이유 읽기

/root/spk/schema/failfast.py 를 앱 이름 spk-schema-failfast 로 만들어 1단계 스키마와 mode=FAILFAST 로 읽은 것을 Parquet 으로 쓰게 하고(경로는 자유), 실패를 잡아 파싱 실패의 오류 조건 이름(MALFORMED_… 로 시작하는 것)을 /root/spk/schema/out/failfast.txt 첫 줄에 쓰세요.

쓰기는 모든 칼럼을 파싱하므로 첫 깨진 줄에서 잡이 실패합니다. 예외는 겹겹이 싸여 있어서 바깥은 파일 읽기 실패이고, 원인 쪽에 파싱 실패가 있습니다. 메시지에서 대괄호 안의 이름을 모두 뽑아 보세요. 채점기는 그 앱에 실패한 잡이 실제로 있었는지도 봅니다.

칸 수가 다른 줄 세기

/root/spk/schema/width.py 를 앱 이름 spk-schema-width 로 만들어 spark.read.text 로 원본을 줄 단위로 읽고, 머리줄을 뺀 줄을 쉼표로 나눈 칸 수별로 세어 /root/spk/schema/out/width.json{"칸 수": 줄 수} 로 쓰세요(열쇠는 문자열).

F.size(F.split("value", ",")) 가 칸 수입니다. 이 원본은 따옴표 안에 쉼표가 없어서 단순히 나눠도 됩니다(따옴표가 있는 CSV 라면 이 방법은 틀립니다). 7칸이 아닌 줄이 3단계의 손상 줄 가운데 얼마를 차지하는지 견줘 보세요.

업무 규칙까지 — 깨끗한 표와 격리 표

/root/spk/schema/clean.py 를 앱 이름 spk-schema-clean 으로 만들어 PERMISSIVE 로 읽은 뒤, 파싱이 되었고(_corrupt 가 null) qty 가 1 이상인 줄만 /root/spk/schema/out/clean 에(손상 칼럼 없이), 나머지는 reason 칼럼(parse 또는 qty)을 붙여 /root/spk/schema/out/quarantine 에 Parquet 으로 쓰세요.

형식은 맞지만 업무상 틀린 줄이 있습니다 — 음수 수량, 빈 수량. 파서는 이것을 깨진 줄로 보지 않으므로 규칙을 직접 적어야 합니다. 이유를 칼럼으로 남겨 두면 협력사에 돌려보낼 때 무엇을 고치라고 할지 바로 나옵니다.

무엇을 버렸고 왜 버렸는지 남기기

/root/spk/schema/report.md## 추론이 틀린 곳 ## 세 가지 모드 ## 격리한 줄 세 절을 쓰세요. 둘째 절에는 3단계의 손상 줄 수와 4단계에서 써서 센 줄 수를, 셋째 절에는 7단계의 깨끗한 줄 수와 격리한 줄 수를 숫자로 넣으세요.

숫자는 여러분의 결과 파일에서 다시 세어 옮기세요. 협력사에 보내는 메일이라고 생각하면 됩니다 — 몇 줄을 받았고, 몇 줄이 왜 깨졌고, 몇 줄을 업무 규칙으로 돌려보내는지.