CSVを直接、本番tableへINSERTすると、途中の1行だけ型が違う、必須値が空、IDが重複するといった問題が、そのまま業務データへ入り込みます。逆に、最初の異常行で処理を止めるだけでは、利用者は修正すべき行を一度に確認できません。
安全な一括取込では、受信した文字列を一度staging tableへ保存し、構文、型、業務ルールを段階的に検証します。正常行だけを正規化して本番tableへ反映し、異常行は行番号と理由を残します。この記事では、CSV、SQLite、Pythonを使って再実行可能な取込契約を設計します。
staging tableが解決する問題
staging tableは、外部データと本番データの間に置く一時的な受け皿です。
CSV file
→ raw staging
→ 型変換・検証
→ valid / invalidに分類
→ transactionで本番へ反映
→ 取込結果を記録
この分離により、次の状態を区別できます。
- file自体を読めない
- headerや列数が契約と違う
- 行は読めるが値を変換できない
- 型は正しいが業務ルールに違反する
- 正常で本番へ反映できる
- 以前と同じfileまたはrecordが再送された
「失敗」という一語にまとめず、どの境界で何が失敗したかをデータとして残すことが重要です。
入力は文字列として保存する
CSVには、JSONのnumberやbooleanのような型情報がありません。次の00125、0.50、falseは、読み取った時点では文字列です。
source_record_id,part_number,name,mass_kg,active
00125,P-001,Bracket,0.50,true
00126,P-002,Shaft,abc,false
00125を先に整数へ変換すると先頭ゼロが失われます。abcを0へ置き換えると、入力エラーと実測値0を区別できません。stagingでは元のcellを変更せず保存し、正規化値を別fieldへ出力します。
{
"import_id": "imp-20261004-001",
"row_number": 3,
"raw": {
"source_record_id": "00126",
"part_number": "P-002",
"name": "Shaft",
"mass_kg": "abc",
"active": "false"
},
"validation_status": "invalid",
"errors": [
{
"field": "mass_kg",
"code": "invalid_decimal"
}
]
}
元値、正規化値、検証結果を分ければ、変換規則を変更した後でも再検証できます。
取込単位と行を分けて保存する
取込全体と各行を別tableへ保存します。
PRAGMA foreign_keys = ON;
CREATE TABLE import_runs (
import_id TEXT PRIMARY KEY NOT NULL,
source_system TEXT NOT NULL,
source_filename TEXT NOT NULL,
file_sha256 TEXT NOT NULL,
schema_version INTEGER NOT NULL CHECK (schema_version > 0),
status TEXT NOT NULL CHECK (
status IN ('received', 'validating', 'ready', 'applied', 'rejected')
),
received_at TEXT NOT NULL,
applied_at TEXT,
UNIQUE (source_system, file_sha256, schema_version)
);
CREATE TABLE part_staging (
import_id TEXT NOT NULL,
row_number INTEGER NOT NULL CHECK (row_number >= 2),
raw_source_record_id TEXT,
raw_part_number TEXT,
raw_name TEXT,
raw_mass_kg TEXT,
raw_active TEXT,
normalized_part_no TEXT,
normalized_name TEXT,
mass_grams INTEGER,
active_value INTEGER,
validation_status TEXT NOT NULL CHECK (
validation_status IN ('pending', 'valid', 'invalid')
),
error_codes_json TEXT,
PRIMARY KEY (import_id, row_number),
FOREIGN KEY (import_id) REFERENCES import_runs(import_id)
ON DELETE CASCADE
);
import_runsはfile単位、part_stagingはrecord単位です。row_numberはCSVの物理行を特定するために保存します。quoted field内に改行を含む場合、CSV record番号と物理行番号が一致しないことがあるため、利用しているparserの定義を確認します。
同じbyte列の再送を検出するため、連携元、SHA-256、schema versionの組み合わせを一意にしています。ただしfile hashは業務recordのIDではありません。同じ内容を別の目的で取り込む要件があるなら、job typeなども一意キーへ含めます。
本番tableは型と制約を強くする
stagingは外部表現を受け入れますが、本番tableには正規化した値だけを保存します。
CREATE TABLE parts (
source_system TEXT NOT NULL,
source_record_id TEXT NOT NULL,
part_number TEXT NOT NULL UNIQUE,
name TEXT NOT NULL CHECK (length(trim(name)) > 0),
mass_grams INTEGER NOT NULL CHECK (mass_grams >= 0),
active INTEGER NOT NULL CHECK (active IN (0, 1)),
updated_at TEXT NOT NULL,
PRIMARY KEY (source_system, source_record_id)
) STRICT;
SQLiteのSTRICT tableでは、列のdatatypeを明示し、losslessに変換できない値をdatatype制約違反として拒否できます。利用環境のSQLiteがSTRICT tableに対応しているか確認し、対応していない場合もNOT NULL、CHECK、UNIQUEなどで境界を守ります。
kgで受け取った質量は、例ではgの整数へ正規化しています。元の0.50はstagingに残し、本番には500を保存します。単位、倍率、丸め方法はschema versionごとに固定します。
検証を段階に分ける
一括取込では、次の順序で検証すると原因を特定しやすくなります。
file検証
- 許可した文字コードでdecodeできるか
- file sizeや行数が上限内か
- delimiter、quote、改行規則が契約どおりか
- header名、順序、重複、未知列をどう扱うか
- file hashがmanifestまたは再送履歴と一致するか
field検証
- 必須cellが空でないか
mass_kgを有限の十進数へ変換できるかactiveが許可したtrueまたはfalseか- IDの先頭ゼロや大文字・小文字を維持するか
- trimやUnicode正規化をどのfieldへ適用するか
record・集合検証
- 同一file内で
source_record_idが重複していないか - part numberが別recordと衝突していないか
- 参照先の分類や単位が存在するか
- 数値が業務範囲内か
- snapshot内に必要なrecordがそろっているか
null、空文字、0、falseは別の値です。CSVの空cellを自動的に0やfalseへ変換せず、fieldごとの欠損規則を定めます。
Pythonでraw行と検証結果を作る
次の例はheaderを固定し、質量をDecimalでgへ変換します。
import csv
from decimal import Decimal, InvalidOperation
from pathlib import Path
EXPECTED_FIELDS = [
"source_record_id",
"part_number",
"name",
"mass_kg",
"active",
]
def validate_row(row: dict[str, str | None]) -> tuple[dict, list[str]]:
errors: list[str] = []
source_record_id = (row.get("source_record_id") or "").strip()
part_number = (row.get("part_number") or "").strip()
name = (row.get("name") or "").strip()
if not source_record_id:
errors.append("source_record_id.required")
if not part_number:
errors.append("part_number.required")
if not name:
errors.append("name.required")
mass_grams = None
try:
mass_kg = Decimal(row.get("mass_kg") or "")
if not mass_kg.is_finite() or mass_kg < 0:
raise ValueError
grams = mass_kg * Decimal("1000")
if grams != grams.to_integral_value():
errors.append("mass_kg.too_precise")
else:
mass_grams = int(grams)
except (InvalidOperation, ValueError):
errors.append("mass_kg.invalid_decimal")
active_text = (row.get("active") or "").strip().lower()
active_value = {"true": 1, "false": 0}.get(active_text)
if active_value is None:
errors.append("active.invalid_boolean")
normalized = {
"source_record_id": source_record_id,
"part_number": part_number,
"name": name,
"mass_grams": mass_grams,
"active": active_value,
}
return normalized, errors
def read_csv_rows(path: Path):
with path.open("r", encoding="utf-8", newline="") as file:
reader = csv.DictReader(file)
if reader.fieldnames != EXPECTED_FIELDS:
raise ValueError("CSV headerが契約と一致しません")
for row in reader:
normalized, errors = validate_row(row)
yield reader.line_num, row, normalized, errors
Pythonのcsv公式資料は、CSVを開く際にnewline=""を使用する例を示しています。encodingも明示し、OSの既定値へ依存させません。DictReaderのkeyとheaderを照合し、未知列を黙って落とさない設計にします。
正常行をUPSERTする
検証済み行は、連携元と外部record IDを同一性のkeyとして本番へ反映します。
INSERT INTO parts (
source_system,
source_record_id,
part_number,
name,
mass_grams,
active,
updated_at
)
SELECT
r.source_system,
s.raw_source_record_id,
s.normalized_part_no,
s.normalized_name,
s.mass_grams,
s.active_value,
CURRENT_TIMESTAMP
FROM part_staging AS s
JOIN import_runs AS r ON r.import_id = s.import_id
WHERE s.import_id = :import_id
AND s.validation_status = 'valid'
ON CONFLICT (source_system, source_record_id)
DO UPDATE SET
part_number = excluded.part_number,
name = excluded.name,
mass_grams = excluded.mass_grams,
active = excluded.active,
updated_at = excluded.updated_at;
SQLiteのUPSERTは、PRIMARY KEYやUNIQUE制約との衝突時にDO UPDATEまたはDO NOTHINGを実行します。NOT NULL、CHECK、外部キー違反を修復する仕組みではないため、事前検証と本番tableの制約を併用します。
transactionで中途半端な反映を防ぐ
本番反映、取込status更新、集計値の記録は、一つのtransactionにします。
import sqlite3
def apply_import(conn: sqlite3.Connection, import_id: str) -> None:
invalid_count = conn.execute(
"""
SELECT count(*)
FROM part_staging
WHERE import_id = ? AND validation_status <> 'valid'
""",
(import_id,),
).fetchone()[0]
if invalid_count:
raise ValueError(f"不正行が{invalid_count}件あります")
with conn:
conn.execute(UPSERT_SQL, {"import_id": import_id})
conn.execute(
"""
UPDATE import_runs
SET status = 'applied',
applied_at = CURRENT_TIMESTAMP
WHERE import_id = ? AND status = 'ready'
""",
(import_id,),
)
Pythonのsqlite3.Connectionをcontext managerとして使うと、blockが正常終了した場合はcommitし、未処理の例外が発生した場合はrollbackします。値は文字列連結ではなくplaceholderで渡します。
この例のUPSERT_SQLには、直前に示したUPSERT文を格納します。実装ではSQLを定数または専用moduleへまとめ、検証処理と同じ版で管理します。
取込方針は次のどちらかを契約で選びます。
| 方針 | 動作 | 向く用途 |
|---|---|---|
| all-or-nothing | 1行でも不正なら本番反映しない | BOM、会計、設定snapshot |
| valid-rows-only | 正常行だけ反映し、不正行を隔離する | 独立した計測値、追記型log |
都合によって実行時に方針を変えると、同じfileから異なる結果が生まれます。import profileへ固定します。
冪等性・再実行・削除を設計する
同じfileを再送されたときは、新しい取込として二重計上せず、既存のimport_idと結果を返します。内容を修正したfileはhashが変わるため、新しいimportとして検証します。
record反映は安定した外部IDでUPSERTし、名称や全field一致を重複判定に使いません。変換規則へschema_versionまたはmapping_versionを持たせれば、過去のrawデータから再計算できます。
また、snapshotに存在しないrecordを即座に削除してはいけません。
- 全件snapshotか差分fileか
- 欠落が削除を意味するか
deleted=trueのtombstoneを使うか- 無効化と物理削除のどちらを行うか
- 参照中のrecordを削除できるか
削除伝播は通常のUPSERTと分離し、明示的な契約と監査記録を用意します。
エラー結果も再利用できる形にする
利用者へ「3行目で失敗」とだけ返すのではなく、機械処理できる結果を作ります。
{
"import_id": "imp-20261004-001",
"status": "rejected",
"total_rows": 2,
"valid_rows": 1,
"invalid_rows": 1,
"errors": [
{
"row_number": 3,
"field": "mass_kg",
"code": "invalid_decimal"
}
]
}
表示用の日本語messageと、処理分岐に使う安定したcodeを分けます。raw値に個人情報や原価が含まれる場合は、通常logへ出さず、staging tableの閲覧権限と保持期限を制限します。
よくある失敗
- CSVから本番へ直接INSERTする:異常行の調査と再処理が困難になります。
- 元値を正規化値で上書きする:変換ミスを後から検証できません。
- 空cellを0へ変える:欠損と数値0を区別できません。
- 行ごとにcommitする:途中失敗で半端な状態が残ります。
- UPSERTだけで検証できると考える:一意性以外の規則は別途必要です。
- 同じfileを毎回新規処理する:再送で二重計上されます。
- snapshotの欠落を自動削除する:抽出漏れが大量削除へ変わります。
- stagingを無期限保存する:機密情報と不要データが蓄積します。
CAD・BOM・API連携へ展開する
staging patternは、部品表、CAD属性、計測結果、API batch、表計算データに共通利用できます。
連携ごとの差は、raw schema、正規化関数、業務validatorへ閉じ込めます。import_runs、status、error code、file hash、transaction、監査項目は共通部品にできます。
この構造を保てば、将来DBやAPIを変更しても、外部入力を即座に本番modelへ結び付けず、変換境界で検証できます。RAGや検索へ渡すmetadataも、検証済みの本番tableから生成できます。
まとめ
安全なCSV一括取込では、読み込めたことと、本番へ反映できることを分けます。
- raw cellを文字列のままstagingへ保存する
- 元値、正規化値、検証結果を別fieldにする
- file、field、record、集合の順に検証する
- 本番tableを型とDB制約で守る
- 正常行だけを安定した外部IDでUPSERTする
- 本番反映とstatus更新を一つのtransactionにする
- all-or-nothingか部分反映かをprofileへ固定する
- file hashとschema versionで再送・再検証を管理する
- 削除伝播を通常更新から分離する
- エラーcode、行番号、保持期限、公開範囲を契約に含める
staging tableを単なる仮置き場ではなく、外部形式と内部modelを分離する検証境界として設計すれば、CSV、API、DB、CAD、BOMの一括連携を安全に再実行できます。

