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

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

スキーマは推論に任せず、壊れた行は捨てずに選り分ける

TT Labで続きを見る

一言でいうと

CSVには型がないので、誰かがスキーマを決める必要があります。推論はデータをもう一度読み、値が1つあるだけでカラム全体の型が揺らぎます。スキーマを直接指定し、壊れた行はPERMISSIVEモードと破損レコード列で捨てずに選り分けて、何行がなぜ壊れたかを数値で残すのが基本です。

なぜスキーマが問題になるのか

Parquetのような形式は、ファイルの中にスキーマを持っています。CSVはカンマで区切られた文字にすぎません。2が整数なのか、2026-01-03 10:00:00がタイムスタンプなのかは、ファイルは教えてくれません。そのため、SparkがCSVを読むときは、次のどちらかをする必要があります。人間がスキーマを指定するか、Sparkがデータを見て推測するかです。

推測にはコストがかかります。CSVデータソースのドキュメントは、inferSchemaのデフォルトがfalseで、オンにすると「データをもう一度走査する必要がある」と書いています。推論に使う割合samplingRatioのデフォルトは1.0、つまり全部です。パートナーが毎日送ってくる10GBのファイルなら、毎日10GBをもう一度読むことになります。

さらに大きな問題は、結果がデータによって揺らぐことです。推論は、カラムのすべての値を収められる型を選びます。数量のフィールドにある日、twoという文字列が1行混ざると、そのカラムは整数ではなく文字列と推論されます。このラボのパートナーのファイルがまさにそうです。推論に任せるとqtyとorder_tsがどちらもstringになります。昨日まで動いていたsum("qty")が今日は変な値を返すのに、コードは1文字も変わっていません。

どう動くのか

PySparkのcsvドキュメントは、全体をもう一度読むのを避けたいなら推論をオフにするかスキーマを直接指定するよう勧めており、スキーマはStructTypeだけでなくDDL文字列でも受け付けます。

ddl = ("order_id STRING, customer_id STRING, product_id STRING, qty INT, "
       "order_ts TIMESTAMP, status STRING, channel STRING, _corrupt_record STRING")
df = (spark.read
      .option("header", True)
      .option("mode", "PERMISSIVE")             # 기본값이지만 적어 둔다
      .schema(ddl)
      .csv("/data/partner/orders.csv"))
bad = df.where(F.col("_corrupt_record").isNotNull())

スキーマを指定すると、Sparkはそのスキーマどおりに値を変換しようとし、変換できない行に出会います。その行をどう扱うかを決めるのがmodeオプションです。ドキュメントが定義している3つは次のとおりです。

破損レコード列の名前はcolumnNameOfCorruptRecordオプションで変更でき、指定しなければ設定のspark.sql.columnNameOfCorruptRecordのデフォルト値_corrupt_recordを使います。

同じCSVの4行を3つのモードで読んだ図です。数量のフィールドにtwoが入った2行目を、PERMISSIVEは数量だけnullにして元の文字列を破損レコード列に入れ、4行すべてを出力し、DROPMALFORMEDはその行を黙って除外して3行を出力し、FAILFASTはその行でMALFORMED_RECORD_IN_PARSINGとなって止まります。下には、countだけを行うとパースが省略されて3つのモードすべてで4行が出るという落とし穴が書かれています

FAILFASTで失敗すると、外側の例外はファイルの読み込みに失敗したというFAILED_READ_FILEで、その原因としてMALFORMED_RECORD_IN_PARSINGが付きます。このエラーの説明文は、壊れたレコードをnullとして処理したいならmodeをPERMISSIVEにするよう案内しています。

countが嘘をつく理由

ここが、このモジュールで最もよく間違える箇所です。同じドキュメントのmodeの説明には、短い警告が付いています。カラムプルーニングのもとでは、CSVは必要なカラムだけをパースしようとするため、どの行が破損と判定されるかは、要求されたカラムの集合によって変わります。この動作はspark.sql.csv.parser.columnPruning.enabledで調整でき、デフォルトでオンです。

count()にはカラムが1つも必要ありません。そのためパーサーはqtyを整数に変換しようともせず、行数だけを数えます。DROPMALFORMEDで読んで数えると壊れた行を除いていない数値が出て、FAILFASTで読んで数えると何事もなく終わります。ラボの6,000行のファイルでも、どちらのモードも6,000になります。すべてのカラムを使う書き込みをして初めて、壊れた行が除かれた5,560行になります。「FAILFASTで実行して通ったからファイルはきれいだ」という結論は、何を要求したかとセットでなければ成り立ちません。

ドキュメントと動作が食い違う箇所

ドキュメントのPERMISSIVEの説明は、フィールド数がスキーマより少ない行や多い行を、CSVでは破損レコードとして扱わないと書いています。少なければ足りないフィールドをnullで埋め、多ければ余ったトークンを捨てるというのです。ところが、このPodのSpark 4.2で破損レコード列を持つスキーマで読んでみると、フィールド数が合わない行も元の文字列が破損レコード列に入って出てきます。ラボのファイルでそのように選り分けられる行は、全部で440行です。バージョンが変わると、こうした細部は動きます。ドキュメントは出発点であり、判定は自分で数えた数値で行います。

現場での姿

黙って捨てるモードが最も危険です。DROPMALFORMEDは便利そうに見えますが、何行をなぜ捨てたのか、何の記録も残しません。パートナーが形式を変えて半分が壊れても、パイプラインは緑のままです。現場では、PERMISSIVEで読み、破損レコード列が埋まった行を隔離テーブルとして別に書き出し、その数をメトリクスとして残します。

パースが通っても業務ルールは別です。数量-3は、整数として問題なく変換されます。パーサーは形式しか見ず、意味は知りません。負の数量や未来の日付のように業務上ありえない値は、パースのあとに別の条件で選り分け、同じ隔離テーブルに送ります。きれいな結果とは、2つのゲートをどちらも通過した行です。

ヘッダー行はデフォルトでは無視されます。ドキュメントでenforceSchemaのデフォルトはtrueで、このとき指定または推論したスキーマを強制的に適用するため、CSVのヘッダー行は無視されます。パートナーがカラムの順序を変えて送ると、名前ではなく位置で入り、黙って間違った値になります。ドキュメントも、誤った結果を避けるにはこのオプションをオフにするよう勧めています。

実務で本当に大切なこと

次のラボですること

パートナーが送ってきた注文CSVの6,000行を、DDLスキーマで読み込みます。まず推論に任せたときにqtyとorder_tsが文字列と推論されることを確認し、PERMISSIVEと破損レコード列で壊れた行を選り分けて何行あるか数えます。同じファイルをDROPMALFORMEDで読んで数えただけの値と、書き込んでから数え直した値がどう違うかを見て、FAILFASTで書き込んで失敗した例外からパースエラーの条件名を受け取ります。行をカンマで分けたフィールド数別に数えて、壊れた行のうちフィールド数が合わないものがどれだけあるかを比べ、最後にパース失敗と、空の数量・負の数量のような業務ルール違反を、理由カラムとともに別に隔離して、きれいな結果を作ります。