スキーマ — 検証しなければ後で倍払う
一言でいうと
外部から入ってくるデータは、すべて文字列だと仮定して始め、明示的に検証して変換したあとにだけ、型のある表に入れる必要があります。
なぜ必要なのか
実務の元データは、例外なく汚れています。同じ日付が4つの形式で入ってきて、金額には通貨記号と千の位のカンマが付いており、ステータス値は大文字小文字と前後の空白がバラバラで、必須項目が空になっています。
これをそのまま型のある表に入れようとすると、ロードが失敗します。すると、よくある2つの誤った対応が出てきます。失敗した行をそのままスキップしたり(誰にも気づかれないままデータが消えます)、すべてのカラムをテキストにしてしまったりする(問題を下流に先送りします)ことです。
どう動くのか
標準的な構造は、3つの層です。
- 元データ(staging): すべてのカラムが文字列。そのままの形で保管します。絶対に直しません。
- クリーン(clean): 検証を通過した行だけが、型を備えた状態で入ります。
- リジェクト(reject): 通過できなかった行と、その理由が入ります。
この構造の核心は、保存則です。クリーンの件数とリジェクトの件数の合計が、元データの件数とちょうど同じである必要があります。どちらにもない行があれば、それは黙って消えたデータであり、パイプラインで最も危険な事故です。
リジェクトの理由を残すことも、妥協できません。理由なしに捨てられた行は、あとで誰も復旧できません。「amountなし」「emailなし」のように、短くてもよいので、必ず一緒に書きます。
変換ルールも明示的である必要があります。日付形式が複数あるなら、それぞれの形式を正規表現で判別して、それぞれ別のパースルールを適用します。自動推論に任せると、03/04/2025が3月4日なのか4月3日なのかによって、静かに間違います。
金額で空文字列を0に変換するのも、よくあるミスです。値がないことと0は違います。決済金額が空なら、それは0ウォンの決済ではなく、情報の欠落であり、0で埋めた瞬間に、その事実が消えます。
現場での姿
スキーマの契約をコードとして残すと役に立ちます。クリーンテーブルのカラムと型の一覧を別の表に記録しておけば、あとで誰かがカラムの型を変えたときに、下流のパイプラインがすぐに検知できます。スキーマ変更は、もともと静かに起き、数日後に変な数字として現れる種類の事故です。
フォーマットの選択も、触れておく価値があります。CSVはどこでも読めますが、型情報がなく、区切り文字のエスケープが脆弱です。JSONは入れ子を表現できますが、容量が大きいです。Parquetのような列指向フォーマットは、型と統計を一緒に持ち、圧縮率が良いため、分析ワークロードに有利です。元データの保管はCSVやJSONで、分析用の再ロードは列フォーマットで分ける構成が、よく見られます。
スキーマ変更を安全に行うルール
パイプラインのスキーマは、生産者とコンシューマーの間の契約です。片方だけを変えると、反対側が壊れるため、どの方向に互換性があるかを先に決めます。
| 互換性の方向 | 意味 | 許可される変更 |
|---|---|---|
| 後方互換(backward) | 新しいコンシューマーが古いデータを読む | フィールドの削除、デフォルト値のあるフィールドの追加 |
| 前方互換(forward) | 古いコンシューマーが新しいデータを読む | フィールドの追加、オプションフィールドの削除 |
| 完全互換(full) | 両方 | デフォルト値のあるオプションフィールドの追加・削除だけ |
ストリーミングでは、後方互換を基本として使います。コンシューマーを先に上げて、生産者をあとで上げればよいからです。順序を逆にすると、新しいデータを古いコンシューマーが受け取って、壊れます。
必須フィールドの追加は、常に破壊的変更です。デフォルト値を与えてオプションとして入れたあと、すべての生産者が埋めるようになったら、そのときに必須に上げます。2段階に分けるのが定石です。
ファイル形式が性能を決める
| 形式 | 構造 | 向いている場所 | 注意 |
|---|---|---|---|
| CSV | 行 | 人が見る少量のデータ | 型がない。エンコーディング・区切り文字の地獄 |
| JSON Lines | 行 | スキーマが流動的なロード | 大きくて遅い |
| Avro | 行 | ストリーミング、イベント | スキーマレジストリと一緒に |
| Parquet | 列 | 分析クエリ | 書き込みが重い。小さなファイルに不向き |
列指向(Parquet)が分析で速い理由は、必要な列だけを読むからです。カラム50個のうち3個だけを使うクエリがよくありますが、行指向は50個をすべて読みます。さらに、列単位の圧縮がよく効き、サイズも小さくなります。
ただし、Parquetは小さなファイルが多いと、かえって遅くなります。ファイルごとにメタデータを読む必要があり、それが実際のデータより大きくなることがあります。128MB–1GBを目標にまとめます。
パーティションはクエリのパターンに従う
s3://lake/events/dt=2026-09-06/hour=14/part-0001.parquet
└─ 날짜로 자르는 질의가 대부분이면 이렇게
パーティションキーを誤って選ぶと、すべてのクエリが全体をスキャンします。逆に、細かく分けすぎると、小さなファイルの問題が生じます。1日にファイルが何個できるかを計算してみて、決めます。
日付をdt=2026-09-06のように1つの文字列として置くほうが、year=/month=/day=に分けるより、たいてい便利です。範囲クエリを使いやすく、ディレクトリの深さも浅くなります。
次のラボですること
文字列だけでできたstaging.orders_rawをプロファイリングして欠陥を数え、日付と金額とステータスを正規化したあと、クリーンテーブルとリジェクトテーブルに分けてロードし、保存則が成り立つかを確認します。