第156回 更新後のCSVを機械的に突合する:Pythonで作る件数・キー・変更列の品質ゲート

CSV更新処理がエラーなく終了しても、想定した範囲だけが変更されたとは限りません。対象外の金額列が書き換わった、行が欠落した、承認した件数と実際の変更件数が違う、といった問題は終了コードだけでは検出できません。

第155回では、変更案を確認してから適用するPlan/Apply型の更新を扱いました。今回はその次の段階として、更新前のCSVと適用後の候補CSVを突合し、後続処理へ進めてよいかを機械的に判定します。

今回作る品質ゲート

例では、注文データのstatusreview_noteだけを更新できるものとします。Pythonスクリプトは次の条件を確認します。

  1. 更新前後で列名と列順が変わっていない
  2. 主キーの追加や欠落がない
  3. 主キーが空欄や重複になっていない
  4. 許可した列以外が変更されていない
  5. 更新前後の両方に存在するキーの変更行数が、承認済みの予定件数と一致する

これは行の追加や削除を行わない「既存行の更新専用」ルールです。追加・削除を含む業務では、この合格条件をそのまま使用できません。また、許可列に書かれた値の正しさまでは判定しません。

1.検証用CSVを用意する

Python 3.10以上を使用します。外部ライブラリは不要です。次の3ファイルを同じディレクトリへ保存してください。

更新前:before.csv

order_id,customer,amount,status,review_note
A001,青木商店,12000,pending,
A002,北川企画,8500,pending,
A003,南工業,5000,approved,確認済み

正しい更新例:after_ok.csv

order_id,customer,amount,status,review_note
A001,青木商店,12000,approved,担当者確認済み
A002,北川企画,8500,pending,
A003,南工業,5000,approved,確認済み

A001のstatusreview_noteだけが変わっています。変更されたセルは2つですが、変更行数は1行です。

誤った更新例:after_ng.csv

order_id,customer,amount,status,review_note
A001,青木商店,12000,approved,担当者確認済み
A002,北川企画,9500,pending,
A004,東商事,5000,approved,確認済み

こちらには、許可していないamountの変更と、A003の欠落、A004の追加が含まれています。

2.突合スクリプトを作る

次のコードをreconcile_csv.pyとして保存します。レポートには変更前後の値を含めず、追加キー、欠落キー、変更行、禁止列の変更を各一覧につき最大20件記録します。

from __future__ import annotations

import argparse
import csv
import json
import sys
from collections import Counter
from pathlib import Path

SAMPLE_LIMIT = 20


class CsvInputError(Exception):
    pass


def read_snapshot(
    path: Path, key: str
) -> tuple[list[str], dict[str, dict[str, str]]]:
    try:
        with path.open("r", encoding="utf-8-sig", newline="") as handle:
            reader = csv.DictReader(handle, strict=True)
            if reader.fieldnames is None:
                raise CsvInputError(f"{path}: ヘッダーがありません")

            headers = reader.fieldnames
            duplicate_names = [
                name
                for name, count in Counter(headers).items()
                if count > 1
            ]

            if any(name.strip() == "" for name in headers):
                raise CsvInputError(f"{path}: 空の列名があります")
            if duplicate_names:
                names = ", ".join(sorted(duplicate_names))
                raise CsvInputError(
                    f"{path}: 列名が重複しています: {names}"
                )
            if key not in headers:
                raise CsvInputError(
                    f"{path}: キー列 {key!r} がありません"
                )

            rows: dict[str, dict[str, str]] = {}
            for row in reader:
                line = reader.line_num

                if None in row:
                    raise CsvInputError(
                        f"{path}:{line}: ヘッダーより多い値があります"
                    )

                missing = [
                    name for name, value in row.items() if value is None
                ]
                if missing:
                    names = ", ".join(missing)
                    raise CsvInputError(
                        f"{path}:{line}: 値が不足している列: {names}"
                    )

                record_key = row[key]
                if record_key.strip() == "":
                    raise CsvInputError(
                        f"{path}:{line}: キー列 {key!r} が空です"
                    )
                if record_key != record_key.strip():
                    raise CsvInputError(
                        f"{path}:{line}: キーの前後に空白があります"
                    )
                if record_key in rows:
                    raise CsvInputError(
                        f"{path}:{line}: キー {record_key!r} が重複しています"
                    )

                rows[record_key] = row

    except CsvInputError:
        raise
    except (OSError, UnicodeError, csv.Error) as exc:
        raise CsvInputError(f"{path}: {exc}") from exc

    return headers, rows


def reconcile(
    before_headers: list[str],
    before_rows: dict[str, dict[str, str]],
    after_headers: list[str],
    after_rows: dict[str, dict[str, str]],
    key: str,
    allowed: set[str],
    expected_changed: int | None,
) -> dict[str, object]:
    before_columns = set(before_headers)
    after_columns = set(after_headers)
    header_added = sorted(after_columns - before_columns)
    header_removed = sorted(before_columns - after_columns)
    header_order_changed = (
        not header_added
        and not header_removed
        and before_headers != after_headers
    )

    before_keys = set(before_rows)
    after_keys = set(after_rows)
    added_keys = after_keys - before_keys
    removed_keys = before_keys - after_keys
    common_keys = sorted(before_keys & after_keys)

    change_counts: Counter[str] = Counter()
    changed_samples: list[dict[str, object]] = []
    forbidden_samples: list[dict[str, object]] = []
    changed_rows = 0
    forbidden_rows = 0
    forbidden_cells = 0

    for record_key in common_keys:
        changed = [
            column
            for column in before_headers
            if column != key
            and column in after_columns
            and before_rows[record_key][column]
            != after_rows[record_key][column]
        ]
        if not changed:
            continue

        changed_rows += 1
        change_counts.update(changed)

        if len(changed_samples) < SAMPLE_LIMIT:
            changed_samples.append(
                {"key": record_key, "columns": changed}
            )

        forbidden = [
            column for column in changed if column not in allowed
        ]
        if forbidden:
            forbidden_rows += 1
            forbidden_cells += len(forbidden)
            if len(forbidden_samples) < SAMPLE_LIMIT:
                forbidden_samples.append(
                    {"key": record_key, "columns": forbidden}
                )

    violations: list[dict[str, object]] = []

    if before_headers != after_headers:
        violations.append(
            {
                "rule": "same_header",
                "added": header_added,
                "removed": header_removed,
                "order_changed": header_order_changed,
            }
        )

    if added_keys or removed_keys:
        violations.append(
            {
                "rule": "same_key_set",
                "added_count": len(added_keys),
                "removed_count": len(removed_keys),
            }
        )

    if forbidden_cells:
        violations.append(
            {
                "rule": "allowed_columns_only",
                "row_count": forbidden_rows,
                "cell_count": forbidden_cells,
            }
        )

    if (
        expected_changed is not None
        and changed_rows != expected_changed
    ):
        violations.append(
            {
                "rule": "expected_changed_rows",
                "expected": expected_changed,
                "actual": changed_rows,
            }
        )

    return {
        "ok": not violations,
        "kind": "reconciliation",
        "policy": {
            "key": key,
            "allowed_columns": sorted(allowed),
            "expected_changed_rows": expected_changed,
        },
        "counts": {
            "before_rows": len(before_rows),
            "after_rows": len(after_rows),
            "common_rows": len(common_keys),
            "added_keys": len(added_keys),
            "removed_keys": len(removed_keys),
            "changed_rows": changed_rows,
            "unchanged_rows": len(common_keys) - changed_rows,
            "forbidden_changed_rows": forbidden_rows,
            "forbidden_changed_cells": forbidden_cells,
        },
        "change_counts_by_column": dict(sorted(change_counts.items())),
        "samples": {
            "added_keys": sorted(added_keys)[:SAMPLE_LIMIT],
            "removed_keys": sorted(removed_keys)[:SAMPLE_LIMIT],
            "changed_rows": changed_samples,
            "forbidden_changes": forbidden_samples,
        },
        "violations": violations,
    }


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser(
        description="更新前後のCSVを突合します"
    )
    parser.add_argument("--before", required=True, type=Path)
    parser.add_argument("--after", required=True, type=Path)
    parser.add_argument("--key", required=True)
    parser.add_argument("--allow", action="append", default=[])
    parser.add_argument("--expect-changed", type=int)
    parser.add_argument("--report", required=True, type=Path)
    args = parser.parse_args()

    if args.expect_changed is not None and args.expect_changed < 0:
        parser.error("--expect-changedには0以上を指定してください")

    before_path = args.before.resolve()
    after_path = args.after.resolve()
    report_path = args.report.resolve()

    if before_path == after_path:
        parser.error("--beforeと--afterには別のファイルを指定してください")
    if report_path in {before_path, after_path}:
        parser.error("--reportに入力CSVと同じパスは指定できません")

    return args


def write_report(path: Path, report: dict[str, object]) -> None:
    path.parent.mkdir(parents=True, exist_ok=True)
    payload = json.dumps(
        report, ensure_ascii=False, indent=2
    ) + "\n"
    path.write_text(payload, encoding="utf-8")


def main() -> int:
    args = parse_args()

    try:
        before_headers, before_rows = read_snapshot(
            args.before, args.key
        )
        after_headers, after_rows = read_snapshot(
            args.after, args.key
        )

        allowed = set(args.allow)
        if args.key in allowed:
            raise CsvInputError(
                "キー列は--allowに指定できません"
            )

        unknown = sorted(allowed - set(before_headers))
        if unknown:
            names = ", ".join(unknown)
            raise CsvInputError(
                f"--allowに存在しない列があります: {names}"
            )

        result = reconcile(
            before_headers,
            before_rows,
            after_headers,
            after_rows,
            args.key,
            allowed,
            args.expect_changed,
        )
        exit_code = 0 if result["ok"] else 1

    except CsvInputError as exc:
        result = {
            "ok": False,
            "kind": "input_error",
            "message": str(exc),
        }
        exit_code = 2

    try:
        write_report(args.report, result)
    except OSError as exc:
        print(
            f"ERROR: レポートを書き込めません: {exc}",
            file=sys.stderr,
        )
        return 2

    if exit_code == 2:
        message = result["message"]
        print(f"ERROR: {message}", file=sys.stderr)
    else:
        status = "PASS" if result["ok"] else "FAIL"
        counts = result["counts"]
        changed_rows = counts["changed_rows"]
        violation_count = len(result["violations"])
        print(
            f"{status}: changed_rows={changed_rows}, "
            f"violations={violation_count}"
        )

    return exit_code


if __name__ == "__main__":
    sys.exit(main())

3.正しい更新を検証する

--allowは変更を許可する列ごとに指定します。承認済みの変更対象が1行なら、--expect-changed 1を付けます。

python3 reconcile_csv.py --before before.csv --after after_ok.csv --key order_id --allow status --allow review_note --expect-changed 1 --report reconciliation_ok.json
echo $?

bashなどのシェルでは、標準出力と終了コードは次のようになります。

PASS: changed_rows=1, violations=0
0

reconciliation_ok.jsonには次のレポートが出力されます。

{
  "ok": true,
  "kind": "reconciliation",
  "policy": {
    "key": "order_id",
    "allowed_columns": [
      "review_note",
      "status"
    ],
    "expected_changed_rows": 1
  },
  "counts": {
    "before_rows": 3,
    "after_rows": 3,
    "common_rows": 3,
    "added_keys": 0,
    "removed_keys": 0,
    "changed_rows": 1,
    "unchanged_rows": 2,
    "forbidden_changed_rows": 0,
    "forbidden_changed_cells": 0
  },
  "change_counts_by_column": {
    "review_note": 1,
    "status": 1
  },
  "samples": {
    "added_keys": [],
    "removed_keys": [],
    "changed_rows": [
      {
        "key": "A001",
        "columns": [
          "status",
          "review_note"
        ]
      }
    ],
    "forbidden_changes": []
  },
  "violations": []
}

4.誤った更新を検出する

同じ条件でafter_ng.csvを検証します。

python3 reconcile_csv.py --before before.csv --after after_ng.csv --key order_id --allow status --allow review_note --expect-changed 1 --report reconciliation_ng.json
echo $?

結果は不合格です。

FAIL: changed_rows=2, violations=3
1

changed_rowsは更新前後の両方に存在するキーだけを比較した件数です。この例ではA001とA002が該当します。A004の追加とA003の欠落は、変更行数ではなく1件のキー集合違反として記録されます。これに、A002のamount変更と予定件数の不一致を加えた合計3件の違反です。

レポートには変更前後の値を記録しないため、金額の値そのものは含まれません。ただし、主キーにも機密性がある場合は、サンプルを記録しない実装へ変更するか、レポートをアクセス制御された場所へ保存してください。

終了コードを後続処理の条件にする

終了コード 意味 後続処理
0 このスクリプトで設定した突合条件に合格 必要な値検証や人の確認を行い、次工程へ進める
1 比較結果が業務ルールに違反 候補CSVを採用せず、レポートを確認する
2 引数、入力CSV、レポート書き込みなどのエラー 実行条件や入力ファイルを修正する

ファイル不足、文字コードエラー、CSV解析エラー、列数の不一致、重複キーなど、比較を続けられない問題は終了コード2になります。CSVを比較できたものの、許可されていない変更などが見つかった場合は終了コード1です。これにより、実行条件を修正すべきエラーと、担当者が差分を確認すべき不一致を分けられます。

入力処理中に検出したエラーでは、kindinput_errorのJSONレポートも出力します。ただし、必須引数の欠落などをargparseが検出した場合は、レポート作成前に終了コード2で停止します。レポートの書き込みに失敗した場合も終了コード2になりますが、指定先を直接上書きする実装のため、完全なJSONレポートが作成されたとは限りません。書き込み失敗の状況によっては、古いファイル、空のファイル、または途中まで書き込まれたファイルが残る可能性があります。後続処理はレポートの存在だけで成功を判断せず、必ず今回のプロセスの終了コードを確認してください。

AIを含む更新フローへの組み込み方

  1. AIまたは担当者が変更案を作る
  2. 元CSVを上書きせず、候補CSVへ適用する
  3. 承認時の予定件数と許可列を引数にして突合する
  4. 終了コードとJSONレポートを確認し、必要な値検証や人の確認を終えてから正式ファイルへ反映する

AIは分類や更新案の作成に使えますが、このスクリプトの合否は列、キー集合、変更範囲、件数という明示的な規則で判定します。終了コード0は、更新内容そのものが正しいという保証ではありません。

この実装の制約

  • 前後のCSVを辞書としてメモリへ読み込むため、利用可能メモリに収まるファイルが対象です。
  • 値は文字列として比較します。1200012000.0は異なる値です。数値や日付を意味で比較したい場合は、比較前の正規化が必要です。
  • 許可列に入っていれば、誤った値でも合格する可能性があります。状態値の候補や金額範囲などは、別の値検証ルールで補います。
  • --expect-changedが確認するのは、更新前後の両方に存在するキーの変更行数だけです。承認対象とは別の行が同じ件数だけ変更された場合までは検出できません。対象キーまで承認している運用では、承認済みキー集合との照合を追加する必要があります。
  • 行追加や削除を伴う処理では、キー集合一致のルールを業務仕様に合わせて変更する必要があります。
  • --expect-changedを省略すると変更行数は合否条件になりません。承認済み件数がある運用では明示的に指定します。
  • この実装は列順の変更も不合格にします。後続処理が列名だけを参照し、列順を問わない場合はルールを緩和できます。
  • JSONレポートは指定先へ直接上書きします。書き込み中の容量不足や強制終了などが起きると、古いファイル、空のファイル、または不完全なファイルが残る可能性があります。この影響を減らすには、一時ファイルへ書き込み、書き込み完了後に置き換える方式を検討してください。

参考にした公式資料

まとめ

安全なCSV更新は、処理が完了した時点では終わりません。更新前のスナップショットと候補CSVを突合し、列、キー集合、変更範囲、件数を確認することで、想定外の構造変更や更新を検出できます。まずは検証用CSVで終了コード0と1を確認し、入力ファイル名を誤らせるなどして終了コード2の扱いも確かめてください。そのうえで、値の妥当性確認や人の承認と組み合わせ、実際の更新処理へ品質ゲートとして組み込みます。