Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
取引先 CSV の壊れた行を三つのモードで読み込む
目標
パートナーが送ってきた注文ファイルの6,000行をスキーマを決めて読み込み、壊れた行をPERMISSIVE・DROPMALFORMED・FAILFASTの3つのモードがそれぞれどう扱うかを数値で確認したうえで、業務ルールも加えて、きれいなテーブルと隔離テーブルに分けます。
なぜ重要なのか
CSVには型がありません。カラム1つに"two"が1回混ざると、inferSchemaはそのカラム全体を文字列にフォールバックし、そのことはエラーではなく静かな型の変化として現れます。昨日まで整数だったカラムが今日文字列になると、後ろの集計がすべて壊れます。そのため、パイプラインはスキーマを推論せず、決めて読み込みます。
スキーマを決めたら、今度はスキーマに合わない行をどうするか決める必要があります。SparkのCSV読み込みにはモードが3つあります。PERMISSIVEは壊れた行を残して元の文字列を破損レコード列に入れ、DROPMALFORMEDは捨て、FAILFASTは止まります。どれが正しいかは業務が決めます。お金を数えるテーブルなら、黙って捨てるのが最も危険です。
落とし穴が1つあります。CSVパーサーは必要なカラムだけをパースします。count()のようにカラムが1つも必要ないアクションは行をパースしないため、DROPMALFORMEDで読んで数えると壊れた行までそのまま数えてしまい、FAILFASTも止まりません。検証をcount()で行ったなら、何も検証していないのと同じです。
ステップ
- /root/spk/schema/schema.ddlに、元データの7つのカラムをDDLで書いて保存してください(
qtyはINT、order_tsはTIMESTAMP、残りはSTRINGです)。 - /root/spk/schema/infer.pyを、アプリ名
spk-schema-inferで作成し、inferSchema=Trueで読み込んで、推論されたスキーマのsimpleString()を、/root/spk/schema/out/inferred.txtに書き込んでください。 - /root/spk/schema/load.pyを、アプリ名
spk-schema-loadで作成し、ステップ1のスキーマの末尾に_corrupt STRINGを追加してPERMISSIVEで読み込み、全体を、/root/spk/schema/out/permissiveにParquetで書き込んでください。 - /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": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。 - /root/spk/schema/failfast.pyを、アプリ名
spk-schema-failfastで作成し、FAILFASTで読み込んだものをParquetで書き込ませてください。失敗の原因となったパースエラーの条件名を、/root/spk/schema/out/failfast.txtの1行目に書き込んでください。 - /root/spk/schema/width.pyを、アプリ名
spk-schema-widthで作成し、元データを行単位で読み込んで(ヘッダー行を除く)、カンマで区切ったフィールド数別の行数を、/root/spk/schema/out/width.jsonに{"칸 수": 줄 수}の形式で書き込んでください(プレースホルダーは順にフィールド数、行数です)。 - /root/spk/schema/clean.pyを、アプリ名
spk-schema-cleanで作成し、パースできた行のうちqtyが1以上のものだけを、/root/spk/schema/out/cleanに、残りはreasonカラム(パース失敗はparse、数量の問題はqty)を付けて、/root/spk/schema/out/quarantineにParquetで書き込んでください。 - /root/spk/schema/report.mdに、
## 추론이 틀린 곳・## 세 가지 모드・## 격리한 줄の3つの節を書いてください(見出しは韓国語で、順に「推論が間違っている箇所」「3つのモード」「隔離した行」を意味します)。2つ目の節には破損行数とDROPMALFORMEDで書き込んだ行数を、3つ目の節にはきれいな行数と隔離した行数を、数値で入れてください。
参考
- 元データ:
/data/shop/orders_dirty.csv(ヘッダー行1行と6,000行)。時刻の形式はyyyy-MM-dd HH:mm:ssなので、timestampFormatで指定してください。 - 破損レコード列は、スキーマに直接入れる必要があります。名前は
columnNameOfCorruptRecordオプションと同じでなければなりません。 - 破損レコード列だけを選んで問い合わせるクエリは、元データを再パースする必要があるため、禁止されている場合があります。一度読んだものを
cache()しておいてから絞り込んでください。 - FAILFASTの例外は、Pythonで条件名を直接返さないことがあります。メッセージの中の角括弧
[...]に条件名が入っています。 - よくあるミス:
count()でモードの効果を確認すること、破損レコード列をスキーマに入れないこと、空の数量(null)をきれいな行に送ること。 - 公式ドキュメント: CSV Files・Data Types・Datetime Patterns・Error Conditions
スキーマを決めておく
元データ/data/shop/orders_dirty.csvのヘッダー行の順に7つのカラムをDDL文字列1行で書き、/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に書き込んでください。
推論はデータをもう一度走査するジョブです。イベントログでジョブ数を見ると、スキーマを決めたときより多くなっています。結果でqtyとorder_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": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。
2つの数値が違う値になるのが正常です。count()にはカラムが1つも必要ないため、CSVパーサーが行をパースせず、パースしなければ壊れているかどうかもわかりません。すべてのカラムを使うアクションを行って初めて、捨てられます。
FAILFAST(止めて理由を読む)
/root/spk/schema/failfast.pyを、アプリ名spk-schema-failfastで作成し、ステップ1のスキーマとmode=FAILFASTで読み込んだものをParquetで書き込ませ(パスは自由です)、失敗を捕捉して、パースエラーの条件名(MALFORMED_…で始まるもの)を、/root/spk/schema/out/failfast.txtの1行目に書き込んでください。
書き込みはすべてのカラムをパースするので、最初に壊れた行でジョブが失敗します。例外は何重にも包まれていて、外側はファイル読み込みの失敗であり、原因側にパースの失敗があります。メッセージから角括弧内の名前をすべて取り出してみてください。採点ツールは、そのアプリで失敗したジョブが実際にあったかも確認します。
フィールド数が違う行を数える
/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つの節を書いてください(見出しは韓国語で、順に「推論が間違っている箇所」「3つのモード」「隔離した行」を意味します)。2つ目の節にはステップ3の破損行数とステップ4で書き込んで数えた行数を、3つ目の節にはステップ7のきれいな行数と隔離した行数を、数値で入れてください。
数値は、自分の結果ファイルから数え直して書き写してください。パートナーに送るメールだと考えればかまいません。何行を受け取り、何行がなぜ壊れていて、何行を業務ルールで返却するのか、という内容です。