TT Lab
はじめる
学ぶ 学習パス コース

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

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

TT Labで続きを見る

目標

パートナーが送ってきた注文ファイルの6,000行をスキーマを決めて読み込み、壊れた行をPERMISSIVE・DROPMALFORMED・FAILFASTの3つのモードがそれぞれどう扱うかを数値で確認したうえで、業務ルールも加えて、きれいなテーブルと隔離テーブルに分けます。

なぜ重要なのか

CSVには型がありません。カラム1つに"two"が1回混ざると、inferSchemaはそのカラム全体を文字列にフォールバックし、そのことはエラーではなく静かな型の変化として現れます。昨日まで整数だったカラムが今日文字列になると、後ろの集計がすべて壊れます。そのため、パイプラインはスキーマを推論せず、決めて読み込みます。 スキーマを決めたら、今度はスキーマに合わない行をどうするか決める必要があります。SparkのCSV読み込みにはモードが3つあります。PERMISSIVEは壊れた行を残して元の文字列を破損レコード列に入れ、DROPMALFORMEDは捨て、FAILFASTは止まります。どれが正しいかは業務が決めます。お金を数えるテーブルなら、黙って捨てるのが最も危険です。 落とし穴が1つあります。CSVパーサーは必要なカラムだけをパースします。count()のようにカラムが1つも必要ないアクションは行をパースしないため、DROPMALFORMEDで読んで数えると壊れた行までそのまま数えてしまい、FAILFASTも止まりません。検証をcount()で行ったなら、何も検証していないのと同じです。

ステップ

  1. /root/spk/schema/schema.ddlに、元データの7つのカラムを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の1行目に書き込んでください。
  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に、## 추론이 틀린 곳・## 세 가지 모드・## 격리한 줄の3つの節を書いてください(見出しは韓国語で、順に「推論が間違っている箇所」「3つのモード」「隔離した行」を意味します)。2つ目の節には破損行数とDROPMALFORMEDで書き込んだ行数を、3つ目の節にはきれいな行数と隔離した行数を、数値で入れてください。

参考

スキーマを決めておく

元データ/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のきれいな行数と隔離した行数を、数値で入れてください。

数値は、自分の結果ファイルから数え直して書き写してください。パートナーに送るメールだと考えればかまいません。何行を受け取り、何行がなぜ壊れていて、何行を業務ルールで返却するのか、という内容です。