現場で受け取るCSVやExcelがそのままではAIモデルや集計に使えず、作業が止まってしまう──そんな経験は多いはずです。本記事では「受け取り側の契約(スキーマ)」を明確にし、検証→自動修復→記録→隔離という実務ワークフローをPythonで実装します。CSVとExcel(.xlsx)の読み込み、検証失敗時の例外処理、修復ログと隔離ファイルの出力までを一つの最小構成にまとめます。第134回・第135回のファイルI/O/クリーニング記事の知識がある前提です。
狙いと前提
目的は、現場で受け取る表データを「そのまま処理できる状態」に自動で近づけることです。完全自動で正解にすることを約束するのではなく、許容できる自動修復ポリシーを定義して運用する点を重視します。
この記事の実装例では、入力ファイルを変更せず、処理結果を次の3種類に分けて出力します。
accepted.csv:検証を通過したレコードrejected.csv:修復または検証に失敗したレコードaudit.jsonl:行番号、処理結果、修復前後、適用ルールを記録した監査ログ
ステップ1:スキーマ設計の実務手順
まずは業務要件から逆算して、スキーマ(契約)を設計します。下表は実務で使いやすいテンプレートです。
| カラム名 | 型 | 必須 | 許容値/正規表現 | 欠損ポリシー | ユニーク | 外部参照 |
|---|---|---|---|---|---|---|
| customer_id | string | 必須 | ^[A-Z0-9\-]{8,}$ | 欠損は隔離 | 真 | 顧客マスタ参照 |
| string | 条件付き | 業務で合意した形式 | 空文字を許容するか明示 | 業務次第 | なし | |
| order_date | date | 必須 | 出力はISO日付 | パース不可は隔離 | false | なし |
| amount | number | 必須 | >=0 | 欠損は隔離、または業務ポリシーを別途定義 | false | 通貨変換ルール |
設計時のチェックリスト:
- 「必須」か「許容する欠損」かを明確にする
- 許容されるデータ型とフォーマットを具体的に例示する
- 自動修復ルールを優先順位付きで決める
- ユニーク制約や外部参照を行単位で確認するか、バッチ後に確認するかを決める
- スキーマのバージョン管理と互換性方針を定める
ステップ2:環境構築と入力ファイル
以下の例はPython 3.10以降を想定します。CSVは標準ライブラリのcsv、Excel(.xlsx)はpandasとopenpyxlで読み込みます。
mkdir table-validator
cd table-validator
python -m venv .venv
# macOS / Linux
source .venv/bin/activate
# Windows PowerShellの場合
# .venv\Scripts\Activate.ps1
python -m pip install --upgrade pip
python -m pip install pandas openpyxl jsonschema "pydantic>=2" pandera pytest
動作確認用のinput.csvを作成します。
customer_id,email,order_date,amount
ABCD-1234, MAILTO:USER@EXAMPLE.COM ,31/01/2025,"1,200"
BAD,user@example.com,not-a-date,-10
EFGH-5678,,2025-02-01,300
1行目は修復後に検証を通過し、2行目は日付のパースに失敗して隔離されます。3行目はそのまま、または軽微な正規化後に通過します。
ステップ3:dataclassからJSON Schemaを生成する軽量実装
次のコードをworkflow.pyとして保存してください。dataclassの型情報とメタデータからJSON Schemaの辞書を生成し、jsonschema.validateで行単位に検証します。
from __future__ import annotations
import argparse
import csv
import json
import unicodedata
from dataclasses import MISSING, asdict, dataclass, field, fields
from datetime import datetime
from pathlib import Path
from typing import Any, get_type_hints
import pandas as pd
from jsonschema import FormatChecker, ValidationError, validate
INPUT_FIELDS = ["customer_id", "email", "order_date", "amount"]
@dataclass
class OrderRow:
customer_id: str = field(
metadata={"pattern": r"^[A-Z0-9\-]{8,}$"}
)
order_date: str = field(
metadata={"format": "date"}
)
amount: float = field(
metadata={"minimum": 0}
)
email: str = ""
def dataclass_to_json_schema(model_class: type) -> dict[str, Any]:
"""この例で使うstrとfloatをJSON Schemaへ変換する。"""
type_map = {
str: "string",
int: "integer",
float: "number",
bool: "boolean",
}
hints = get_type_hints(model_class)
properties: dict[str, Any] = {}
required: list[str] = []
for item in fields(model_class):
python_type = hints[item.name]
if python_type not in type_map:
raise TypeError(f"未対応の型です: {item.name}={python_type}")
property_schema = {"type": type_map[python_type]}
for key in ("pattern", "format", "minimum", "maximum"):
if key in item.metadata:
property_schema[key] = item.metadata[key]
properties[item.name] = property_schema
if item.default is MISSING and item.default_factory is MISSING:
required.append(item.name)
return {
"$schema": "https://json-schema.org/draft/2020-12/schema",
"type": "object",
"properties": properties,
"required": required,
"additionalProperties": False,
}
ORDER_SCHEMA = dataclass_to_json_schema(OrderRow)
def normalize_text(value: Any) -> str:
text = "" if value is None else str(value)
return unicodedata.normalize("NFKC", text).strip()
def repair_row(row: dict[str, Any]) -> tuple[dict[str, Any], list[dict[str, Any]]]:
"""許可した修復だけを適用し、修復内容も返す。"""
repaired = dict(row)
changes: list[dict[str, Any]] = []
def set_value(field_name: str, new_value: Any, rule: str) -> None:
old_value = repaired.get(field_name)
repaired[field_name] = new_value
if old_value != new_value:
changes.append(
{
"field": field_name,
"rule": rule,
"before": old_value,
"after": new_value,
}
)
customer_id = normalize_text(repaired.get("customer_id"))
set_value("customer_id", customer_id, "NFKC変換と前後空白削除")
email = normalize_text(repaired.get("email")).lower()
if email.startswith("mailto:"):
email = email.removeprefix("mailto:").strip()
set_value("email", email, "小文字化、前後空白削除、mailto:削除")
amount_text = normalize_text(repaired.get("amount")).replace(",", "")
try:
amount = float(amount_text)
except ValueError as exc:
raise ValueError(f"amountを数値に変換できません: {amount_text!r}") from exc
set_value("amount", amount, "NFKC変換、カンマ削除、float変換")
date_text = normalize_text(repaired.get("order_date"))
parsed_date = None
selected_format = None
for date_format in ("%Y-%m-%d", "%d/%m/%Y"):
try:
parsed_date = datetime.strptime(date_text, date_format).date()
selected_format = date_format
break
except ValueError:
continue
if parsed_date is None:
raise ValueError(f"order_dateを日付に変換できません: {date_text!r}")
set_value(
"order_date",
parsed_date.isoformat(),
f"{selected_format}で解析しISO日付へ変換",
)
return repaired, changes
def read_rows(path: Path) -> list[dict[str, Any]]:
"""CSVまたは.xlsxを辞書のリストとして読み込む。"""
suffix = path.suffix.lower()
if suffix == ".csv":
with path.open("r", encoding="utf-8-sig", newline="") as handle:
return [dict(row) for row in csv.DictReader(handle)]
if suffix == ".xlsx":
dataframe = pd.read_excel(
path,
engine="openpyxl",
dtype=str,
keep_default_na=False,
)
return dataframe.to_dict(orient="records")
raise ValueError("対応形式は.csvまたは.xlsxです")
def write_csv(
path: Path,
rows: list[dict[str, Any]],
fieldnames: list[str],
) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("w", encoding="utf-8-sig", newline="") as handle:
writer = csv.DictWriter(handle, fieldnames=fieldnames)
writer.writeheader()
writer.writerows(rows)
def process_file(input_path: Path, output_dir: Path) -> None:
accepted: list[dict[str, Any]] = []
rejected: list[dict[str, Any]] = []
audit_records: list[dict[str, Any]] = []
for line_number, original in enumerate(read_rows(input_path), start=2):
changes: list[dict[str, Any]] = []
try:
repaired, changes = repair_row(original)
validate(
instance=repaired,
schema=ORDER_SCHEMA,
format_checker=FormatChecker(),
)
# 検証後のキーと型を固定する。
normalized = asdict(
OrderRow(
customer_id=repaired["customer_id"],
email=repaired.get("email", ""),
order_date=repaired["order_date"],
amount=repaired["amount"],
)
)
accepted.append(normalized)
status = "repaired" if changes else "accepted"
audit_records.append(
{
"line": line_number,
"status": status,
"original": original,
"result": normalized,
"changes": changes,
}
)
except (ValueError, ValidationError) as exc:
rejected_row = dict(original)
rejected_row["_line"] = line_number
rejected_row["_error"] = str(exc)
rejected.append(rejected_row)
audit_records.append(
{
"line": line_number,
"status": "rejected",
"original": original,
"changes": changes,
"error_type": type(exc).__name__,
"error": str(exc),
}
)
output_dir.mkdir(parents=True, exist_ok=True)
write_csv(
output_dir / "accepted.csv",
accepted,
["customer_id", "order_date", "amount", "email"],
)
extra_reject_fields = sorted(
{
key
for row in rejected
for key in row
if key not in INPUT_FIELDS and key not in {"_line", "_error"}
}
)
reject_fields = INPUT_FIELDS + extra_reject_fields + ["_line", "_error"]
write_csv(output_dir / "rejected.csv", rejected, reject_fields)
with (output_dir / "audit.jsonl").open("w", encoding="utf-8") as handle:
for record in audit_records:
handle.write(json.dumps(record, ensure_ascii=False) + "\n")
print(f"accepted={len(accepted)} rejected={len(rejected)}")
print(f"output={output_dir.resolve()}")
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("input", type=Path)
parser.add_argument("--output-dir", type=Path, default=Path("output"))
args = parser.parse_args()
process_file(args.input, args.output_dir)
if __name__ == "__main__":
main()
保存後、次の順番で実行します。
python workflow.py input.csv --output-dir output
Excelを入力する場合は、同じ列を持つinput.xlsxを用意して次のように実行します。出力形式はCSVとJSON Linesです。
python workflow.py input.xlsx --output-dir output
repair_row内の変換に失敗した場合はValueError、JSON Schemaに違反した場合はjsonschema.ValidationErrorを捕捉します。いずれも処理全体を停止させず、その行をrejected.csvへ隔離します。入力に想定外の列がある場合も、additionalProperties: falseにより検証エラーになります。
ステップ4:実装パターンの比較
| パターン | 概要 | メリット | デメリット |
|---|---|---|---|
| 軽量(dataclass + jsonschema) | dataclassからJSON Schemaを生成し、行単位で検証する | 処理順と修復ルールを明示しやすい | 型変換や複雑な修復は手実装になる |
| Pydantic | モデル単位で型変換と検証を行う | 詳細な検証エラーを取得しやすい | 独自の修復ルールにはvalidatorの実装が必要 |
| Pandera | pandas DataFrame全体の列型や制約を検証する | ユニーク制約などの表単位の検査を記述しやすい | 行単位の修復・隔離は別途設計が必要 |
Pydanticによる行単位の検証例
Pydantic v2を使う場合は、model_validateで入力辞書を検証し、ValidationErrorを捕捉します。次は日付、金額、メールをモデル内で正規化する最小例です。
import unicodedata
from datetime import date, datetime
from typing import Any
from pydantic import (
BaseModel,
ConfigDict,
Field,
ValidationError as PydanticValidationError,
field_validator,
)
class OrderModel(BaseModel):
model_config = ConfigDict(extra="forbid")
customer_id: str = Field(pattern=r"^[A-Z0-9\-]{8,}$")
order_date: date
amount: float = Field(ge=0)
email: str = ""
@field_validator("customer_id", mode="before")
@classmethod
def normalize_customer_id(cls, value: Any) -> str:
return unicodedata.normalize("NFKC", str(value)).strip()
@field_validator("amount", mode="before")
@classmethod
def normalize_amount(cls, value: Any) -> float:
text = unicodedata.normalize("NFKC", str(value)).strip()
return float(text.replace(",", ""))
@field_validator("order_date", mode="before")
@classmethod
def normalize_order_date(cls, value: Any) -> date:
if isinstance(value, date):
return value
text = str(value).strip()
for date_format in ("%Y-%m-%d", "%d/%m/%Y"):
try:
return datetime.strptime(text, date_format).date()
except ValueError:
continue
raise ValueError(f"日付に変換できません: {text!r}")
@field_validator("email", mode="before")
@classmethod
def normalize_email(cls, value: Any) -> str:
text = "" if value is None else str(value)
text = unicodedata.normalize("NFKC", text).strip().lower()
return text.removeprefix("mailto:").strip()
raw = {
"customer_id": "ABCD-1234",
"email": " MAILTO:USER@EXAMPLE.COM ",
"order_date": "31/01/2025",
"amount": "1,200",
}
try:
order = OrderModel.model_validate(raw)
print(order.model_dump(mode="json"))
except PydanticValidationError as exc:
for error in exc.errors():
print(error)
実運用では、成功したmodel_dump()の結果を正常系へ送り、PydanticValidationErrorを捕捉したレコードを隔離します。修復前の辞書は上書きせず、監査ログへ保存してください。
PanderaによるDataFrame全体の検証例
行単位の修復後に、ユニーク制約や列全体の型を確認する場合はPanderaを利用できます。次の例は、先ほど生成したoutput/accepted.csvを再検証します。
import pandas as pd
import pandera as pa
dataframe = pd.read_csv(
"output/accepted.csv",
keep_default_na=False,
)
order_schema = pa.DataFrameSchema(
{
"customer_id": pa.Column(
str,
checks=pa.Check.str_matches(r"^[A-Z0-9\-]{8,}$"),
unique=True,
),
"order_date": pa.Column("datetime64[ns]", nullable=False),
"amount": pa.Column(float, checks=pa.Check.ge(0)),
"email": pa.Column(str, nullable=False),
},
strict=True,
coerce=True,
)
try:
validated_dataframe = order_schema.validate(dataframe, lazy=True)
print(validated_dataframe)
except pa.errors.SchemaErrors as exc:
print(exc.failure_cases)
exc.failure_cases.to_csv(
"output/pandera_failure_cases.csv",
index=False,
encoding="utf-8-sig",
)
lazy=Trueでは、検出した複数の違反をSchemaErrorsとしてまとめて確認できます。failure_casesは違反内容の調査に使えますが、元レコードを隔離する処理は、利用するインデックスや業務上の主キーに合わせて別途実装してください。
ステップ5:自動修復ポリシー
| フィールド | 修復ルール | 優先度 | ログ/差分 |
|---|---|---|---|
| 文字列 | NFKC正規化、前後空白の削除 | 中 | 修復前後と適用ルールを保存 |
| amount | 全角数字の正規化、カンマ削除、float変換 | 高 | 変換前後を保存 |
| order_date | %Y-%m-%d、%d/%m/%Yの順で試し、ISO日付へ統一 | 高 | 採用した日付形式を保存 |
| 小文字化、前後空白削除、mailto:削除 | 低〜中 | 修復前後を保存 |
金額のクリッピングや欠損値の0埋めは、正しい値を推測する変更になり得ます。そのため上記の実装には含めていません。実施する場合は、業務上の閾値、承認条件、HITL(人の確認)を明文化してください。
- 修復前の入力ファイルを上書きしない
- 修復理由を構造化ログとして残す
- 重大な修復は自動確定せず、担当者の確認対象にする
- 日付のように解釈が分かれる形式は、送信元と優先順位を合意する
ステップ6:ワークフロー統合と運用フロー
バリデーションはETLの初段に組み込みます。定期実行にはcron、ワークフロー管理ツール、CIなど、運用環境に合う方法を選びます。重要なのは、検証失敗時に処理全体を停止するのか、行単位で隔離して継続するのかを事前に決めることです。
| ケース | 自動対応 | 運用アクション |
|---|---|---|
| 軽微な修復で処理可能 | 修復後に再検証して正常系へ送る | 修復ログを記録 |
| 自動修復では確定できない | 隔離ファイルへ出力 | 担当者が原本とエラーを確認 |
| 重大な整合性違反 | 本番投入を止めるか対象行を隔離 | 送信元へのフィードバックとスキーマ再合意 |
SLOの例:
- 受信ファイルの90%を自動修復で本番投入可能にする(週次評価)
- 重大な隔離レコードは24時間以内に原因調査を開始する
テストとCI
次のコードをtest_workflow.pyとして保存します。
import pytest
from jsonschema import FormatChecker, ValidationError, validate
from workflow import ORDER_SCHEMA, repair_row
def test_repair_and_validate():
raw = {
"customer_id": " ABCD-1234 ",
"email": " MAILTO:USER@EXAMPLE.COM ",
"order_date": "31/01/2025",
"amount": "1,200",
}
repaired, changes = repair_row(raw)
validate(
instance=repaired,
schema=ORDER_SCHEMA,
format_checker=FormatChecker(),
)
assert repaired["order_date"] == "2025-01-31"
assert repaired["amount"] == 1200.0
assert repaired["email"] == "user@example.com"
assert changes
def test_negative_amount_is_rejected():
raw = {
"customer_id": "ABCD-1234",
"email": "",
"order_date": "2025-01-31",
"amount": "-10",
}
repaired, _ = repair_row(raw)
with pytest.raises(ValidationError):
validate(
instance=repaired,
schema=ORDER_SCHEMA,
format_checker=FormatChecker(),
)
def test_invalid_date_is_rejected():
raw = {
"customer_id": "ABCD-1234",
"email": "",
"order_date": "not-a-date",
"amount": "100",
}
with pytest.raises(ValueError):
repair_row(raw)
テストを実行します。
pytest -q
| CI段階 | チェック内容 |
|---|---|
| PR時 | pytestで検証と自動修復のユニットテストを実行 |
| マージ後 | 匿名化したサンプルデータで統合テストを実行 |
スキーマ設計用YAMLテンプレート
次のYAMLは送信元との合意やレビューに使う設計テンプレートです。上記のPythonコードがこのYAMLを直接読み込むわけではありません。
fields:
- name: customer_id
type: string
required: true
pattern: '^[A-Z0-9\-]{8,}$'
- name: order_date
type: date
required: true
formats:
- '%Y-%m-%d'
- '%d/%m/%Y'
- name: amount
type: number
required: true
minimum: 0
- name: email
type: string
required: false
運用上の注意点と落とし穴
- 過度な厳格化は業務停止を招くため、最低限の必須ルールから始める
- 入力側とのスキーマ合意とサンプル交換を定期的に行う
- スキーマをバージョン管理し、破壊的変更は事前に通知する
- ログや監査情報を保存し、どの入力にどの修復を行ったか追跡できるようにする
- CSVの文字コード、区切り文字、Excelの対象シートも受け渡し契約に含める
次の一歩
- 主要テーブル1つから導入し、隔離理由の多い項目を確認する
- ユニーク制約や外部参照をバッチ検証へ追加する
- 生成したスキーマを使ったデータ生成テストで破壊的変更を検知する
まとめ
- スキーマは受け取り側と渡し手の「契約」であり、型、必須、形式、修復範囲を合意することが重要です。
- dataclassとjsonschemaを使うと、行単位の修復、検証、例外処理、隔離を小さく実装できます。
- Pydanticはモデル単位の型変換と詳細なエラー、PanderaはDataFrame全体の制約確認に利用できます。
- 自動修復では入力を上書きせず、修復前後と適用ルールを監査ログへ残してください。
- pytestとCIで正常系、逸脱例、境界値を継続的に確認してください。