第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の扱いも確かめてください。そのうえで、値の妥当性確認や人の承認と組み合わせ、実際の更新処理へ品質ゲートとして組み込みます。

第155回 AIの提案を直接反映しない:Pythonで作るドライランとPlan/Apply型CSV更新

AIに顧客情報や申請データを分類させる場合、結果をそのまま元データへ書き込むと、誤分類や入力の取り違えを後から確認しにくくなります。そこで今回は、AIや担当者が作った更新案をいったん「変更計画」に変換し、人が確認してから別のCSVへ適用する仕組みを作ります。

前回扱った増分処理と冪等性が「同じ処理を安全に再実行する設計」だとすれば、今回の主題は「提案と適用の間に確認可能な境界を置く設計」です。生成AIのAPI自体は呼び出さないため、APIキーなしで再現できます。

今回作るワークフロー

  1. 現在のデータをrecords.csvに用意する。
  2. AIまたは担当者の更新案をsuggestions.csvとして受け取る。
  3. planコマンドで入力を検証し、plan.jsonを作る。この段階では元データを変更しない。
  4. 担当者が変更内容を確認し、計画ファイルのSHA-256ハッシュを取得する。
  5. applyコマンドでハッシュと元データの一致を再確認し、更新結果を新しいCSVへ出力する。

適用時に元CSVが計画作成時と一致しない場合、または計画ファイルが承認ハッシュ取得時と一致しない場合、適用処理は失敗します。元CSVは上書きせず、出力先がすでに存在する場合も停止します。

前提と対象範囲

  • Python 3.10以上を使用します。外部ライブラリは不要です。
  • 例ではpendingapprovedrejectedの3状態だけを許可します。
  • 入力をメモリへ読み込むため、担当者がレビューできる規模のCSVを対象とします。
  • 同じ出力先へ複数のプロセスが同時に書き込まない運用を前提とします。

入力CSVを用意する

records.csvは現在の正本です。

record_id,name,status
A001,青木商店,pending
A002,北川企画,pending
A003,南製作所,approved

suggestions.csvはAIや担当者が作った提案です。理由も必須にして、レビュー時に判断根拠を確認できるようにします。

record_id,proposed_status,reason
A001,approved,必要項目を確認済み
A002,rejected,連絡先が不足
A003,approved,現状維持

この例ではA001とA002が変更対象です。A003は現在値と提案値が同じなので、計画には変更操作として追加しません。

Plan/Applyを実装する

次のコードをworkflow.pyとして保存します。CSVの重複ID、未知のID、列不足、不正な状態、空の理由を計画作成前に検出します。

from __future__ import annotations

import argparse
import csv
import hashlib
import io
import json
import os
import sys
import tempfile
from pathlib import Path
from typing import Any

ALLOWED_STATUSES = {'pending', 'approved', 'rejected'}
PLAN_VERSION = 1


class WorkflowError(Exception):
    pass


def sha256_bytes(data: bytes) -> str:
    return hashlib.sha256(data).hexdigest()


def parse_csv(
    data: bytes,
    label: str,
    required: set[str],
) -> tuple[list[str], list[dict[str, str]]]:
    try:
        text = data.decode('utf-8-sig')
    except UnicodeDecodeError as exc:
        raise WorkflowError(f'{label}: UTF-8として読み込めません') from exc

    reader = csv.DictReader(io.StringIO(text, newline=''), strict=True)
    if not reader.fieldnames:
        raise WorkflowError(f'{label}: ヘッダーがありません')

    fields = list(reader.fieldnames)
    if any(name == '' for name in fields):
        raise WorkflowError(f'{label}: 空の列名があります')
    if len(fields) != len(set(fields)):
        raise WorkflowError(f'{label}: 列名が重複しています')

    missing = sorted(required - set(fields))
    if missing:
        missing_text = ', '.join(missing)
        raise WorkflowError(f'{label}: 必須列がありません: {missing_text}')

    rows: list[dict[str, str]] = []
    for data_no, row in enumerate(reader, start=1):
        if None in row or any(value is None for value in row.values()):
            raise WorkflowError(f'{label}: データ行{data_no}の列数が不正です')
        rows.append(dict(row))
    return fields, rows


def load_records(
    data: bytes,
    label: str,
) -> tuple[list[str], list[dict[str, str]], dict[str, dict[str, str]]]:
    fields, rows = parse_csv(data, label, {'record_id', 'status'})
    index: dict[str, dict[str, str]] = {}

    for data_no, row in enumerate(rows, start=1):
        record_id = row['record_id']
        status = row['status']

        if not record_id or record_id != record_id.strip():
            raise WorkflowError(f'{label}: データ行{data_no}のrecord_idが不正です')
        if record_id in index:
            raise WorkflowError(f'{label}: record_idが重複しています: {record_id}')
        if status not in ALLOWED_STATUSES:
            raise WorkflowError(f'{label}: {record_id}のstatusが不正です: {status}')

        index[record_id] = row
    return fields, rows, index


def atomic_write_text(
    path: Path,
    text: str,
    replace: bool = False,
) -> None:
    if not path.parent.is_dir():
        raise WorkflowError(f'出力先ディレクトリがありません: {path.parent}')
    if path.exists() and not replace:
        raise WorkflowError(f'出力先がすでに存在します: {path}')

    temp_path: Path | None = None
    try:
        with tempfile.NamedTemporaryFile(
            mode='w',
            encoding='utf-8',
            newline='',
            dir=path.parent,
            prefix=f'.{path.name}.',
            suffix='.tmp',
            delete=False,
        ) as temporary:
            temp_path = Path(temporary.name)
            temporary.write(text)
            temporary.flush()
            os.fsync(temporary.fileno())

        if path.exists() and not replace:
            raise WorkflowError(f'出力先がすでに存在します: {path}')
        os.replace(temp_path, path)
    except BaseException:
        if temp_path is not None:
            try:
                temp_path.unlink(missing_ok=True)
            except OSError:
                pass
        raise


def make_plan(
    records_path: Path,
    suggestions_path: Path,
    plan_path: Path,
    replace: bool,
) -> dict[str, Any]:
    if plan_path.resolve() in {
        records_path.resolve(),
        suggestions_path.resolve(),
    }:
        raise WorkflowError('計画ファイルを入力CSVと同じ場所にはできません')

    records_bytes = records_path.read_bytes()
    suggestions_bytes = suggestions_path.read_bytes()
    _, record_rows, records = load_records(records_bytes, str(records_path))
    _, suggestions = parse_csv(
        suggestions_bytes,
        str(suggestions_path),
        {'record_id', 'proposed_status', 'reason'},
    )

    seen: set[str] = set()
    operations: list[dict[str, str]] = []
    unchanged = 0

    for data_no, suggestion in enumerate(suggestions, start=1):
        record_id = suggestion['record_id']
        proposed = suggestion['proposed_status']
        reason = suggestion['reason'].strip()

        if not record_id or record_id != record_id.strip():
            raise WorkflowError(
                f'{suggestions_path}: データ行{data_no}のrecord_idが不正です'
            )
        if record_id in seen:
            raise WorkflowError(f'提案のrecord_idが重複しています: {record_id}')
        seen.add(record_id)

        if record_id not in records:
            raise WorkflowError(f'元データに存在しないrecord_idです: {record_id}')
        if proposed not in ALLOWED_STATUSES:
            raise WorkflowError(f'{record_id}の提案状態が不正です: {proposed}')
        if not reason:
            raise WorkflowError(f'{record_id}のreasonが空です')

        current = records[record_id]['status']
        if current == proposed:
            unchanged += 1
            continue

        operations.append(
            {
                'record_id': record_id,
                'before': current,
                'after': proposed,
                'reason': reason,
            }
        )

    operations.sort(key=lambda operation: operation['record_id'])
    plan: dict[str, Any] = {
        'format_version': PLAN_VERSION,
        'records_sha256': sha256_bytes(records_bytes),
        'suggestions_sha256': sha256_bytes(suggestions_bytes),
        'record_count': len(record_rows),
        'suggestion_count': len(suggestions),
        'unchanged_count': unchanged,
        'operations': operations,
    }

    plan_text = json.dumps(
        plan,
        ensure_ascii=False,
        sort_keys=True,
        indent=2,
    ) + '\n'
    atomic_write_text(plan_path, plan_text, replace=replace)
    return plan


def load_plan(data: bytes, label: str) -> dict[str, Any]:
    try:
        parsed = json.loads(data.decode('utf-8-sig'))
    except (UnicodeDecodeError, json.JSONDecodeError) as exc:
        raise WorkflowError(f'{label}: 有効なJSONではありません') from exc

    if not isinstance(parsed, dict):
        raise WorkflowError(f'{label}: JSONオブジェクトではありません')
    if type(parsed.get('format_version')) is not int:
        raise WorkflowError(f'{label}: format_versionが不正です')
    if parsed['format_version'] != PLAN_VERSION:
        raise WorkflowError(f'{label}: 未対応のformat_versionです')

    source_hash = parsed.get('records_sha256')
    if (
        not isinstance(source_hash, str)
        or len(source_hash) != 64
        or any(character not in '0123456789abcdef' for character in source_hash.lower())
    ):
        raise WorkflowError(f'{label}: records_sha256が不正です')
    if not isinstance(parsed.get('operations'), list):
        raise WorkflowError(f'{label}: operationsが配列ではありません')
    return parsed


def apply_plan(
    records_path: Path,
    plan_path: Path,
    output_path: Path,
    approved_hash: str,
) -> dict[str, Any]:
    if output_path.resolve() in {records_path.resolve(), plan_path.resolve()}:
        raise WorkflowError('出力CSVを入力CSVまたは計画ファイルと同じ場所にはできません')

    approved = approved_hash.strip().lower()
    if (
        len(approved) != 64
        or any(character not in '0123456789abcdef' for character in approved)
    ):
        raise WorkflowError('承認ハッシュは64桁の16進数で指定してください')

    plan_bytes = plan_path.read_bytes()
    actual_plan_hash = sha256_bytes(plan_bytes)
    if actual_plan_hash != approved:
        raise WorkflowError('計画ファイルが承認時の内容と一致しません')
    plan = load_plan(plan_bytes, str(plan_path))

    records_bytes = records_path.read_bytes()
    if sha256_bytes(records_bytes) != plan['records_sha256'].lower():
        raise WorkflowError(
            '入力CSVが計画作成後に変化しています。計画を作り直してください'
        )

    fields, rows, records = load_records(records_bytes, str(records_path))
    seen: set[str] = set()

    for operation_no, operation in enumerate(plan['operations'], start=1):
        if not isinstance(operation, dict):
            raise WorkflowError(f'操作{operation_no}がオブジェクトではありません')

        required = {'record_id', 'before', 'after', 'reason'}
        if not required.issubset(operation):
            raise WorkflowError(f'操作{operation_no}の項目が不足しています')
        if not all(isinstance(operation[key], str) for key in required):
            raise WorkflowError(f'操作{operation_no}に文字列以外の値があります')

        record_id = operation['record_id']
        before = operation['before']
        after = operation['after']

        if record_id in seen:
            raise WorkflowError(f'計画内のrecord_idが重複しています: {record_id}')
        seen.add(record_id)
        if record_id not in records:
            raise WorkflowError(f'計画対象が元データにありません: {record_id}')
        if before not in ALLOWED_STATUSES or after not in ALLOWED_STATUSES:
            raise WorkflowError(f'{record_id}の状態が不正です')
        if before == after:
            raise WorkflowError(f'{record_id}は変更前後が同じです')
        if not operation['reason'].strip():
            raise WorkflowError(f'{record_id}のreasonが空です')
        if records[record_id]['status'] != before:
            raise WorkflowError(f'{record_id}の変更前状態が一致しません')

        records[record_id]['status'] = after

    buffer = io.StringIO(newline='')
    writer = csv.DictWriter(buffer, fieldnames=fields, lineterminator='\n')
    writer.writeheader()
    writer.writerows(rows)
    atomic_write_text(output_path, buffer.getvalue())

    return {
        'applied': len(plan['operations']),
        'output': str(output_path),
        'plan_sha256': actual_plan_hash,
    }


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        description='CSVの状態変更を計画し、承認後に別ファイルへ適用します'
    )
    commands = parser.add_subparsers(dest='command', required=True)

    plan_command = commands.add_parser('plan', help='変更計画を作成します')
    plan_command.add_argument('--records', type=Path, required=True)
    plan_command.add_argument('--suggestions', type=Path, required=True)
    plan_command.add_argument('--plan', type=Path, required=True)
    plan_command.add_argument('--replace-plan', action='store_true')

    apply_command = commands.add_parser('apply', help='承認済み計画を適用します')
    apply_command.add_argument('--records', type=Path, required=True)
    apply_command.add_argument('--plan', type=Path, required=True)
    apply_command.add_argument('--output', type=Path, required=True)
    apply_command.add_argument('--approve-hash', required=True)
    return parser


def main(argv: list[str] | None = None) -> int:
    args = build_parser().parse_args(argv)
    try:
        if args.command == 'plan':
            plan = make_plan(
                args.records,
                args.suggestions,
                args.plan,
                replace=args.replace_plan,
            )
            result = {
                'mode': 'plan',
                'plan': str(args.plan),
                'plan_sha256': sha256_bytes(args.plan.read_bytes()),
                'changes': len(plan['operations']),
                'unchanged': plan['unchanged_count'],
            }
        else:
            result = {
                'mode': 'apply',
                **apply_plan(
                    args.records,
                    args.plan,
                    args.output,
                    args.approve_hash,
                ),
            }
    except (WorkflowError, OSError, csv.Error) as exc:
        print(f'error: {exc}', file=sys.stderr)
        return 1

    print(json.dumps(result, ensure_ascii=False, sort_keys=True))
    return 0


if __name__ == '__main__':
    raise SystemExit(main())

実行順と入出力

1. 変更計画を作る

python workflow.py plan --records records.csv --suggestions suggestions.csv --plan plan.json

標準出力には、変更件数、変更なしの件数、計画ファイルのSHA-256ハッシュがJSONで1行出力されます。この例では変更2件、変更なし1件です。元のrecords.csvは変更されません。

計画内容は次のコマンドで読みやすく表示できます。

python -m json.tool plan.json

operationsにあるrecord_idbeforeafterreasonを確認します。既存の計画を作り直す場合だけ--replace-planを指定します。置き換えると以前のハッシュは一致しなくなるため、再確認が必要です。

2. 確認した計画のハッシュを取得する

レビューを終えた後、次のコマンドで計画ファイルの現在のハッシュを表示します。

python -c "from pathlib import Path; import hashlib; print(hashlib.sha256(Path('plan.json').read_bytes()).hexdigest())"

表示された64桁の値を、確認済み計画の識別子として控えます。ハッシュを取得した後に計画ファイルを編集した場合は、内容を再確認してハッシュも取り直します。

3. 計画を適用する

次のHASHを、直前に取得した実際の値へ置き換えて実行します。

python workflow.py apply --records records.csv --plan plan.json --output updated_records.csv --approve-hash HASH

成功時は適用件数、出力先、計画ハッシュが標準出力へJSONで出ます。生成されるupdated_records.csvの論理的な内容は次のとおりです。

record_id,name,status
A001,青木商店,approved
A002,北川企画,rejected
A003,南製作所,approved

元のrecords.csvは残ります。CSVは再生成されるため、引用符、BOM、改行コードなどのバイト表現は元ファイルと異なる場合があります。

失敗時の動作

  • 利用方法の誤りはargparseが標準エラーへ案内を出し、終了コード2で終了します。
  • 入力不備、ハッシュ不一致、既存の出力先など、想定内の業務エラーは標準エラーへ理由を出し、終了コード1で停止します。
  • 全操作の検証が終わるまで出力先へ切り替えないため、途中の操作だけを反映したCSVは残りません。
  • 一時ファイルは出力先と同じディレクトリへ作り、書き込み完了後にos.replaceで出力先へ切り替えます。

たとえば計画作成後にrecords.csvを編集すると、適用時に「入力CSVが計画作成後に変化しています」と表示され、出力CSVは作成されません。改行コードやBOMだけの変更もハッシュ差として検出します。これは、論理的に同じに見える別ファイルを誤って適用しないための保守的な動作です。

最小テストを追加する

次をtest_workflow.pyとして保存します。正常適用、計画作成後の入力変更、承認後の計画変更を検査します。

import csv
import tempfile
import unittest
from pathlib import Path

from workflow import WorkflowError, apply_plan, make_plan, sha256_bytes


RECORDS = '''record_id,name,status
A001,青木商店,pending
A002,北川企画,pending
A003,南製作所,approved
'''

SUGGESTIONS = '''record_id,proposed_status,reason
A001,approved,必要項目を確認済み
A002,rejected,連絡先が不足
A003,approved,現状維持
'''


class WorkflowTest(unittest.TestCase):
    def prepare(self, root: Path) -> tuple[Path, Path, Path]:
        records = root / 'records.csv'
        suggestions = root / 'suggestions.csv'
        plan = root / 'plan.json'
        records.write_text(RECORDS, encoding='utf-8')
        suggestions.write_text(SUGGESTIONS, encoding='utf-8')
        return records, suggestions, plan

    def test_apply_changes_only_planned_rows(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            plan_data = make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            result = apply_plan(records, plan, output, approval)

            self.assertEqual(len(plan_data['operations']), 2)
            self.assertEqual(result['applied'], 2)
            self.assertEqual(
                sha256_bytes(records.read_bytes()),
                plan_data['records_sha256'],
            )

            with output.open(encoding='utf-8', newline='') as file:
                rows = list(csv.DictReader(file))
            statuses = {row['record_id']: row['status'] for row in rows}
            self.assertEqual(
                statuses,
                {
                    'A001': 'approved',
                    'A002': 'rejected',
                    'A003': 'approved',
                },
            )

    def test_rejects_changed_input(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            changed = RECORDS.replace(
                'A001,青木商店,pending',
                'A001,青木商店,approved',
            )
            records.write_text(changed, encoding='utf-8')

            with self.assertRaisesRegex(
                WorkflowError,
                '入力CSVが計画作成後に変化',
            ):
                apply_plan(records, plan, output, approval)
            self.assertFalse(output.exists())

    def test_rejects_changed_plan(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            plan.write_text(
                plan.read_text(encoding='utf-8') + '\n',
                encoding='utf-8',
            )

            with self.assertRaisesRegex(
                WorkflowError,
                '計画ファイルが承認時の内容と一致しません',
            ):
                apply_plan(records, plan, output, approval)
            self.assertFalse(output.exists())


if __name__ == '__main__':
    unittest.main()

テストは次のコマンドで実行します。

python -m unittest -v test_workflow

3つのテストがokになれば、サンプルの正常系、元データ変更の拒否、計画ファイル変更の拒否を確認できます。

この実装で保証しないこと

SHA-256ハッシュはファイル内容の一致確認には使えますが、誰が承認したかを証明するものではありません。ハッシュを渡すだけで承認者の認証や電子署名になるわけでもありません。チーム運用では、計画ファイルとハッシュをアクセス制御されたレビューやCIの承認手順へ渡し、提案者、確認者、適用日時を別途記録します。

また、この例は同時実行ロック、大容量CSVのストリーミング、状態遷移ごとの権限制御を実装していません。複数プロセスが同じ出力先を扱う場合は排他制御を追加し、大容量データでは変更対象をデータベースや不変スナップショットへ移す設計を検討してください。表計算ソフトで開くCSVに外部由来の文字列を含める場合は、数式として解釈される文字列への対策も別途必要です。

導入時のチェックリスト

  • AIの出力を正本へ直接書き込まず、提案ファイルとして分離する。
  • 許可する値、必須列、重複ID、未知のIDを機械的に検証する。
  • 変更前後の値と理由を計画ファイルで確認する。
  • 確認後に取得した計画ファイルのハッシュを適用時に照合する。
  • 元データを上書きせず、出力結果を比較してから次工程へ渡す。

最初は実際の業務CSVを少量だけ複製し、許可状態と必須列を自社のルールへ置き換えて試してください。AIの判断精度とは別に、誤った提案を適用前に止められる工程を持つことが重要です。

第154回 実務で使える増分処理と冪等性の設計:Pythonで作る差分検出・重複除外・安全な書き込みパターン

はじめに — つまずきに寄り添う

増分処理や冪等性は、現場で繰り返し失敗すると面倒さが身にしみるテーマです。処理が重複してコストが増えたり、途中で止まったときに再実行してデータが壊れたりすると、日常運用が不安定になります。本記事では「現場ですぐ使える」Pythonパターンに絞り、差分検出から重複除外、冪等な書き込み、再実行・回復の設計まで段階的に解説します。シリーズ: AIとPythonの実務

なぜ増分処理と冪等性が重要か

  • コスト削減:フル処理を避け、対象の差分だけ処理することで時間・APIコール・計算資源を節約できます。
  • 再実行の安全性:途中失敗後のリトライで重複や不整合を避けられます。
  • データ一貫性:一貫した状態での保存・集計が可能になり、監査性も向上します。

増分検出の主要戦略と選び方

代表的な戦略を比較表で示します。運用コストや導入しやすさを見て選んでください。

戦略 メリット 欠点 運用コスト
mtime / ファイル更新時刻 簡単、ファイルシステム依存で高速 時刻の変更で誤検出、タイムゾーン/NFS の問題
シーケンスID/タイムスタンプ 順序が安定していれば増分を追いやすく、API向け ソース側のサポートが必要。同一タイムスタンプや遅延到着への対策も必要
ハッシュ比較(コンテンツハッシュ) 正確に差分を検出、内容ベース 計算コストが高い(大ファイル)、部分差分の検出が難しい 中〜高
差分ログ(change log / CDC) 完全な変更履歴、正確な再現性 導入コストが高い、ログ保持の管理が必要

Pythonでの差分検出パターン(ハンズオン指針)

1) mtimeチェック(ファイル単位)

シンプルで初期導入が容易。例:

from pathlib import Path

def changed_files(src_dir, last_run_time):
    for p in Path(src_dir).rglob('*'):
        try:
            if not p.is_file():
                continue
            mtime = p.stat().st_mtime
        except OSError:
            continue
        if mtime > last_run_time:
            yield p

注意点:rglob('*')はディレクトリも列挙するため、is_file()で通常ファイルだけに絞ります。また、タイムスタンプの精度やNFSタイムスタンプの不一致を考慮してください。

2) 行単位ハッシュ(内容ベース)

CSVやログの差分検出に有効。行ごとにハッシュを取り、既存ハッシュセットと照合します。

import hashlib

def line_hash(line: str) -> str:
    return hashlib.sha256(line.encode('utf-8')).hexdigest()

seen = set()  # 永続化するのが実運用
new_rows = []
for line in open('data.csv'):
    h = line_hash(line)
    if h not in seen:
        seen.add(h)
        new_rows.append(line)

3) ウォーターマーク管理(SQLiteを使う例)

処理位置を小さなDBで管理します。タイムスタンプだけを保存すると、同じタイムスタンプを持つ複数レコードや遅れて到着したデータを取りこぼす可能性があります。そのため、タイムスタンプと一意IDの複合カーソルを使い、(timestamp, unique_id)の順で処理位置を管理します。

import sqlite3

def init_db(path='meta.sqlite'):
    conn = sqlite3.connect(path)
    conn.execute('''
        CREATE TABLE IF NOT EXISTS watermark (
            name TEXT PRIMARY KEY,
            timestamp TEXT NOT NULL,
            unique_id TEXT NOT NULL
        )
    ''')
    conn.commit()
    return conn

def get_watermark(conn, name):
    cur = conn.execute(
        'SELECT timestamp, unique_id FROM watermark WHERE name=?',
        (name,)
    )
    row = cur.fetchone()
    return (row[0], row[1]) if row else None

def set_watermark(conn, name, timestamp, unique_id):
    conn.execute('''
        INSERT INTO watermark(name, timestamp, unique_id)
        VALUES (?, ?, ?)
        ON CONFLICT(name) DO UPDATE SET
            timestamp=excluded.timestamp,
            unique_id=excluded.unique_id
    ''', (name, timestamp, unique_id))
    conn.commit()

取得側では、日時を同一形式・同一タイムゾーンに正規化したうえで、概念的にはtimestamp > 前回timestamp OR (timestamp = 前回timestamp AND unique_id > 前回unique_id)という条件を使い、同じ順序で並べます。一意IDには、順序比較が安定する値を使用してください。

遅延到着があり得る場合は、前回位置より少し前から重複区間を再取得し、一意キーによるupsertや処理済みIDで重複を除外します。ウォーターマークは対象データの書き込みが成功した後にだけ更新し、可能であればデータ更新と同じトランザクションに含めます。

重複除外と冪等な書き込みの実装パターン

  • upsert設計:データベース側の一意キーを用いてINSERT … ON CONFLICTまたはREPLACEを利用する。
  • トランザクション的書き込み:一連の更新をトランザクションでまとめ、途中失敗時には全体をロールバック。
  • 一時ファイル→原子置換:出力はまずtempに書き込み、最終的にos.replaceで置き換える。これで中途半端なファイル状態の露出を防ぐ。
  • ファイルロック・ポジション保存:長時間処理や並列処理ではロックで同時書き込みを回避。位置情報を保存して部分再開も可能。

原子置換の最小例(atomic write):

import os
from tempfile import NamedTemporaryFile

def atomic_write(path, data_bytes):
    dirpath = os.path.dirname(path) or '.'
    with NamedTemporaryFile(dir=dirpath, delete=False) as tmp:
        tmp.write(data_bytes)
        tmp.flush()
        tmp_name = tmp.name
    os.replace(tmp_name, path)  # 原子的に置換

SQLiteでの簡単なupsert例:

# テーブル: items(id PRIMARY KEY, data TEXT)
conn.execute('INSERT INTO items(id,data) VALUES (?,?) ON CONFLICT(id) DO UPDATE SET data=excluded.data', (id, data))
conn.commit()

API連携での増分処理

  • バッチ化:API呼び出しはまとめて行い、オフセットや複合カーソルをウォーターマークで管理します。
  • 並列呼び出し時の重複リスク:位置管理に加え、APIが対応している場合はidempotency-keyを付与してAPI側で重複排除できる設計にします。
  • de-dupキーの設計:ソースの一意ID+バッチIDやコンテンツハッシュを候補にします。

idempotency-keyは試行ごとではなく、同じ論理処理に対して一度だけ生成します。送信前に処理状態とともに永続化し、成功が確定するまでは再試行でも同じキーを再利用します。以下はSQLiteにキーを保存する例です。

import sqlite3
import uuid
import requests

def init_request_db(path='meta.sqlite'):
    conn = sqlite3.connect(path)
    conn.execute('''
        CREATE TABLE IF NOT EXISTS api_requests (
            operation_id TEXT PRIMARY KEY,
            idempotency_key TEXT NOT NULL,
            status TEXT NOT NULL
        )
    ''')
    conn.commit()
    return conn

def get_or_create_key(conn, operation_id):
    row = conn.execute(
        'SELECT idempotency_key FROM api_requests WHERE operation_id=?',
        (operation_id,)
    ).fetchone()
    if row:
        return row[0]

    key = str(uuid.uuid4())
    conn.execute(
        'INSERT INTO api_requests(operation_id, idempotency_key, status) VALUES (?, ?, ?)',
        (operation_id, key, 'pending')
    )
    conn.commit()  # API送信前に保存する
    return key

def post_with_idempotency(conn, operation_id, url, payload):
    key = get_or_create_key(conn, operation_id)
    headers = {'Idempotency-Key': key}

    # 例外時もpendingのレコードとキーが残るため、
    # 同じoperation_idでの再試行では同じキーが使われる
    resp = requests.post(
        url,
        json=payload,
        headers=headers,
        timeout=30
    )
    resp.raise_for_status()

    conn.execute(
        'UPDATE api_requests SET status=? WHERE operation_id=?',
        ('succeeded', operation_id)
    )
    conn.commit()
    return resp

operation_idにはバッチIDなど、同じ論理処理を再試行しても変わらない識別子を指定します。この方法は、連携先APIがIdempotency-Keyと同一キーの再送をサポートしていることが前提です。

再実行・失敗からの回復設計

  • チェックポイントの粒度:ファイル単位、バッチ単位、レコード単位など。処理コストと回復の細かさで決める。
  • 部分再実行:複合カーソルや処理済みIDのリストで未処理部分のみ再実行する。遅延到着がある場合は重複区間を再取得する。
  • 冪等なリトライ戦略:リトライ前提なら、同じ論理処理で同じidempotencyキーを再利用するか、upsertを利用して重複を許容しない設計にする。
  • 監査ログ:処理開始時/終了時/失敗時にキー情報(ウォーターマーク、バッチID、idempotencyキー、件数、エラー)を残す。

テストとCIで検証する項目

差分・重複・削除などのケースをユニットと統合で検証します。モックを使って外部APIやファイルシステムの失敗を再現しましょう。

  • 差分ケース:追加・更新・削除・同一データの重複・同一タイムスタンプ・遅延到着を網羅。
  • モックAPI:レスポンス遅延や部分エラーを再現し、同じ論理処理のリトライでidempotencyキーが変わらないことを確認。
  • CIでの回帰確認ポイント:ウォーターマークの更新、重複率閾値、処理スループット。

運用チェックリストとモニタリング指標

項目 説明 / 計測方法 推奨閾値例
スループット レコード/秒、バッチあたり処理件数 業務要件に依存
重複率 同一IDやハッシュが処理済みで再度処理された割合 <1%(目安)
遅延時間 データ生成から処理完了までの時間 SLAに応じて設定
ウォーターマーク整合性 ウォーターマークと実処理の不整合チェック 不整合は即アラート

実務テンプレートと次の一歩

まずは小さなCLIテンプレートで試してください。argparseベースの最小構成例:

import argparse

def main():
    p = argparse.ArgumentParser()
    p.add_argument('--src', required=True)
    p.add_argument('--meta', default='meta.sqlite')
    args = p.parse_args()
    # 差分検出 → 処理 → ウォーターマーク更新 の流れを書く

if __name__ == '__main__':
    main()

配布案内:この記事で示した小さなリファレンス実装(CLIテンプレート、SQLiteメタ管理、atomic writeユーティリティ)をManage AIのリポジトリで配布予定です(公開予定日付近に案内します)。次回は「運用の自動化と監査強化」を取り上げる予定です。

まとめ

増分処理と冪等性は、現場の安定稼働に直結する重要な設計要素です。まずはmtimeや行ハッシュ、タイムスタンプと一意IDによる複合カーソルから始め、遅延到着がある場合は重複区間の再取得と重複除外を組み合わせてください。重要なのは「完璧を目指すより、失敗から安全に回復できる設計」を段階的に導入することです。

チェックリスト要約:

  • 差分戦略を選ぶ(mtime / ID / ハッシュ / CDC)
  • 複合カーソルまたは永続的なハッシュDBで状態管理し、必要に応じて重複区間を再取得する
  • 書き込みは原子的/トランザクション的に行う(upsert, atomic write)
  • APIでは同じ論理処理のidempotencyキーを保存し、成功確定まで再利用する
  • CI・テストで差分ケースを網羅し、運用監視を設定する

Manage AI(https://manageai.online)では、実務に使える小さなツール群とテンプレートを継続的に紹介します。まずはこの記事のサンプルパターンを試し、貴社の運用要件に合わせてカスタマイズしてください。

第153回 堅牢な中断・再開対応パイプライン設計:Pythonで作るチェックポイントと状態管理の実務手順

長時間処理や外部APIを伴うワークフローでは、途中で止まってしまったときに「どこから再開するか」「重複処理をどう防ぐか」で悩むことが多いはずです。本稿は、第152回(並行処理と再実行設計)を読んだ前提で、Python(標準ライブラリ+軽量外部依存)で現場に導入しやすいチェックポイント(checkpointing)と状態管理のパターン、実装例、運用上の注意点を実務目線で整理します。まずは、現場でよくあるつまずきに寄り添いながら進めます。

現場でよくある失敗パターンとチェックポイントでの解決狙い

実際の現場で見かける代表的な問題と、チェックポイント導入で期待できる効果を整理します。

よくある問題 影響 チェックポイントでの解決策
処理途中で失敗し、全件を再実行する 時間とコストの浪費、APIレート超過 進捗を保存し、最後に成功した境界から再開
重複処理(重複課金・二重登録) 外部サービスでのコストやデータ不整合 冪等キー・トランザクション境界でスキップ判定
部分成功データの中途半端な出力 解析や後工程でのエラーの温床 原子的コミットや一時ファイル+renameで整合性を担保

基本設計パターン(実務の観点)

チェックポイント設計で押さえておくべき基本パターンを短く整理します。

  • 冪等性 (idempotency):同じ操作を複数回実行しても結果が変わらないよう、外部APIには冪等キーを渡したり、ローカルで既処理レコードを参照してスキップする。
  • トランザクション境界:「ここまで成功したら永続化する」という明確な境界を決める。境界は小さくするほど回復が細かくできるが管理コストが増える。
  • 原子的コミット:ローカルファイルは tmp → atomic rename(os.replace)で最終化する。S3ではrename相当の操作がないため、条件付き書き込みや実行ごとに一意なキーを使い、同じ最終キーへの並行書き込みを制御する。
  • 部分出力の扱い:中間出力は明示的に状態として残し、後工程はそのスキーマを前提に再実行できるようにする。
  • ロールフォワードとロールバック:失敗時はロールバックよりロールフォワード(欠損部分を補う)が現場では現実的な場合が多い。重要データはトランザクションで保護。

実装パターンと具体例(利点・注意点付き)

ここでは現場で採用しやすい4つのパターンを示します。コードはそのまま貼れるPythonスニペットです。

1) ファイルベース(tmpファイル → atomic rename)

長時間処理の中間結果をローカルファイルに書くときの基本パターン。小規模バッチに最適。

利点 注意点
依存が少なく、実装が簡単 分散環境やクラッシュでの復旧ロジックを追加する必要あり
原子的操作(os.replace)で整合性確保 ファイル破損やディスク容量に注意
import os
import json

def atomic_write(path, data):
    tmp = path + ".tmp"
    with open(tmp, 'w', encoding='utf-8') as f:
        json.dump(data, f)
        f.flush()
        os.fsync(f.fileno())
    os.replace(tmp, path)  # atomic on most OS

2) JSONL(追記ログ)

処理ごとに追記していくログ形式。イベントソーシング的に使える。

利点 注意点
追記のみなので衝突が少ない サイズ管理とGCが必要(古いログの圧縮)
追記順で再生(ロールフォワード)可能 途中での不整合行は復旧が難しい場合あり
import json
import os
from datetime import datetime

def append_event(path, event):
    line = json.dumps({"ts": datetime.utcnow().isoformat(), **event}, ensure_ascii=False)
    with open(path, 'a', encoding='utf-8') as f:
        f.write(line + "\n")
        f.flush()
        os.fsync(f.fileno())

3) SQLite(軽量状態DB)

状態管理(チェックポイント、未処理リスト、成功フラグ等)にはSQLiteが便利。トランザクションで整合性を保てます。

利点 注意点
ACID対応で安全 高並列書き込みは苦手(適切な排他が必要)
ファイルで持てるのでバックアップ容易 スキーマ変更時のマイグレーションが必要
import sqlite3

def init_db(path):
    con = sqlite3.connect(path, isolation_level=None)
    cur = con.cursor()
    cur.execute('''
    CREATE TABLE IF NOT EXISTS checkpoints (
        key TEXT PRIMARY KEY,
        state JSON,
        updated_at TEXT
    )
    ''')
    return con

def save_checkpoint(con, key, state_json):
    cur = con.cursor()
    cur.execute('REPLACE INTO checkpoints (key, state, updated_at) VALUES (?, ?, CURRENT_TIMESTAMP)', (key, state_json))

4) クラウドストレージ(S3)でのアップロード

S3ではrename操作がないため、複数の実行が固定の一時キーを共有する方式は避けます。最終キーがまだ存在しない場合だけ公開する用途では、条件付きPUTを使うことで、同じキーに対する並行実行の上書きを防げます。上書きが必要なワークフローでは、実行ごとに一意な最終キーを使うか、別途排他制御を設けてください。

利点 注意点
耐久性が高く、共有が容易 rename操作がないため、条件付き書き込みや一意なキーの設計が必要
分散環境で扱いやすい リージョンや認証まわりの運用設定に注意
# boto3 を使い、最終キーが未作成の場合だけアップロードする例
import boto3
from botocore.exceptions import ClientError

s3 = boto3.client('s3')

def upload_if_absent(bucket, key, data_bytes):
    try:
        s3.put_object(
            Bucket=bucket,
            Key=key,
            Body=data_bytes,
            IfNoneMatch='*'
        )
    except ClientError as exc:
        status = exc.response.get('ResponseMetadata', {}).get('HTTPStatusCode')
        if status == 412:
            return False  # すでに同じ最終キーが存在する
        raise
    return True

再開(復旧)ロジックの実装ポイント

処理の再実行判定やスキップロジックはシンプルに保つことが重要です。代表的なパターンを示します。

  • チェックサム(入力のハッシュ)を保存しておき、入力が変わっていなければスキップできる。
  • ステータスマーク(pending / running / success / failed)を保存し、runningが長時間続く場合は調査フラグにする。
  • 部分再処理は「最小の再実行単位(バケット)」を設計しておき、そこ単位で再実行する。

ただし、チェックポイントだけではexactly-once(厳密に1回だけの実行)は保証できません。外部処理の成功後、成功チェックポイントの保存前にプロセスが停止すると、再開時に外部処理が再実行されます。二重登録や二重課金を防ぐには、再試行しても変わらない冪等キーを外部処理へ渡し、外部サービス側でも同じキーの重複実行を抑止する必要があります。

以下の擬似例では、process()を、受け取った冪等キーを外部APIへ転送するアプリケーション側の関数として扱います。同じ項目と同じ入力には、再開後も同じ冪等キーを使用します。

import hashlib

def checksum_of_bytes(b):
    return hashlib.sha256(b).hexdigest()

# 再開ロジック(擬似)
item_key = 'item-123'
input_checksum = checksum_of_bytes(input_bytes)
idempotency_key = f'{item_key}:{input_checksum}'

state = load_checkpoint(item_key)  # state は dict
if state and state.get('status') == 'success' and state.get('input_checksum') == input_checksum:
    print('skip: already processed')
else:
    try:
        mark_running(item_key)
        result = process(
            input_bytes,
            idempotency_key=idempotency_key
        )
        save_checkpoint(item_key, {
            'status': 'success',
            'input_checksum': input_checksum
        })
    except Exception:
        save_checkpoint(item_key, {'status': 'failed'})
        raise

テストと検証

チェックポイント機能はテスト設計が重要です。単体テスト、統合テスト、障害注入テストのチェックリストを示します。

レイヤ 検証項目 方法
単体テスト ファイル書き込みが原子的に見えるか/SQLiteのトランザクション 一時ディレクトリやin-memory SQLiteで検証、mockでIOを代替
統合テスト 実際に中断→再開して期待通り結果が得られるか CIで一時S3バケットやテストDBを使って流し込み
障害注入 書き込み時のクラッシュ、ネットワーク断を再現 プロセス強制終了、モックで例外を発生させる

CIに組み込むべき検証ポイント例:

  • チェックポイントを書いた後の再起動で処理がスキップされること
  • 中間ファイルが残らない(.tmp等の掃除)こと
  • 外部APIに対する冪等キーによる二重請求防止の確認

運用と監視の実務チェックリスト

導入後に確認・監視すべきポイントを表で整理します。

監視項目 閾値/ルール 対応アクション
未完了ジョブ数 常時0に近いこと。上昇トレンドは問題 原因調査・自動再実行の検討
最終成功時刻 想定頻度を超える遅延がある場合アラート 処理負荷や外部APIの障害確認
再実行頻度 一定以上なら根本原因の調査が必要 コードや外部依存の堅牢化

ゴミデータ(古いチェックポイントや一時ファイル)は定期クリーンアップをスケジュールし、最小保持ポリシーを決めて自動化してください。

移行・互換性とメンテナンス

チェックポイントデータのスキーマ進化は実務で必ず起きます。対処方針の例:

  • checkpoint に version フィールドを追加し、読み込み時に古いバージョンを検出してマイグレーション関数を呼ぶ。
  • 重大変更時はバックフィル計画を作り、テスト環境でまず適用する。
  • 古いチェックポイントは段階的に削除するポリシーを定める(例:90日経過でアーカイブ→削除)。
# 簡単なバージョニング読み込み例
def load_state_with_migration(raw):
    version = raw.get('version', 1)
    if version == 1:
        return migrate_v1_to_v2(raw)
    return raw

まとめと次の一歩

本稿では、チェックポイント設計の目的とトレードオフ、冪等性や原子的コミットなどの基本パターン、ファイル/JSONL/SQLite/S3それぞれの実装例、再開ロジック、テスト・運用・移行のポイントを実務的に整理しました。まずは小さなバッチで「tmp→atomic rename」や「SQLiteでの状態管理」から導入し、問題が起きたら障害注入テストで洗い出すのが現場で現実的な進め方です。

テンプレートコード(抜粋)を下に示します。WordPressにそのまま貼れる形式にしてあるので、プロジェクトにコピーして試してください。

# チェックポイント簡易テンプレート
import os, json, hashlib

def atomic_write(path, data):
    tmp = path + '.tmp'
    with open(tmp, 'w', encoding='utf-8') as f:
        json.dump(data, f, ensure_ascii=False)
        f.flush(); os.fsync(f.fileno())
    os.replace(tmp, path)

def checksum(b):
    return hashlib.sha256(b).hexdigest()

def process_item(item_bytes, state_path):
    try:
        with open(state_path, 'r', encoding='utf-8') as f:
            state = json.load(f)
    except FileNotFoundError:
        state = {}

    cs = checksum(item_bytes)
    if state.get('input_checksum') == cs and state.get('status') == 'success':
        return 'skipped'

    # ... 実処理ここから
    # 外部APIを呼ぶ場合は、csなどから作った同一の冪等キーを再試行時にも渡す
    # 処理が成功したらチェックポイントを原子的に保存
    new_state = {'status': 'success', 'input_checksum': cs}
    atomic_write(state_path, new_state)
    return 'done'

次回候補としては、作ったチェックポイントをAirflowやPrefectのようなオーケストレーションツールとどう連携するか、監査ログやガバナンスとどうつなぐかを扱う予定です。本シリーズ「AIとPythonの実務」の流れで、現場で使える実装を続けて紹介します。

抜粋:長時間処理や部分失敗が発生する現場向けに、処理の中断・再開(checkpointing)と状態管理をPythonで実装する手順を解説します。第152回(並行処理と再実行設計)の続きとして、再実行可能・冪等・安全な状態保存の実務的パターンに焦点を当てます。

第152回 実務で使える並行処理と再実行設計 — Pythonで作る安全で効率的なバッチ&API呼び出しワークフロー

はじめに — 現場でよくあるつまずきに寄り添う

大量のCSVを処理したり、埋め込みや推論APIを大量に呼び出したりすると、処理が遅い・途中で失敗してやり直しが大変・コストが膨らむ、という問題に直面します。本記事は「どの並列化を選ぶか」「失敗時にどう回復するか」を、Pythonコードと実務向けのチェックリストで示します。第151回(logging/CLI/config)との連携を想定し、運用につながる形で説明します。

1) いつ並列化すべきか(事前計測の手順)

まずは並列化する前にボトルネックを把握します。小さなサンプルで計測し、I/O待ちかCPU負荷か、多接続が原因かを見極めます。

  • 測定項目:1件あたりの平均処理時間、外部APIの待ち時間、CPU使用率、メモリ増加
  • 測定方法の例:代表的な100件でwall-clock時間を測る。外部通信時間が総時間の大部分ならI/Oバウンド。
状況 並列化の候補 注目点
I/Oバウンド(ネットワーク/API) ThreadPoolExecutor / asyncio 同時接続数の制御、タイムアウト設定
CPUバウンド(大きな変換/圧縮) ProcessPoolExecutor GIL回避、pickle可能性
多数の同時接続(高並列API呼び出し) asyncio + aiohttp 等 バックプレッシャー、レート制限対応

2) パターン紹介と短いコード例

ここでは代表的な3パターンの最小構成を示します。実務では利用するAPI、ライブラリ、ネットワーク条件に合わせて、タイムアウトや例外処理を追加してください。

ThreadPoolExecutor(簡単に並列化)

外部APIを呼ぶI/Oバウンド処理で使いやすい方法です。次の例は標準ライブラリで複数URLを並行取得します。スレッド数を増やしすぎず、接続先の制限に合わせて調整してください。

from concurrent.futures import ThreadPoolExecutor
from urllib.request import urlopen


def fetch_one(url):
    with urlopen(url, timeout=10) as response:
        return response.read()


urls = [
    "https://example.com/",
    "https://example.com/",
]

with ThreadPoolExecutor(max_workers=4) as executor:
    bodies = list(executor.map(fetch_one, urls))

executor.map()は入力順に結果を返します。個別の完了順で処理したい場合や、失敗した項目を特定したい場合は、後述の実践例のようにsubmit()wait()を使います。

ProcessPoolExecutor(CPUバウンド向け)

CPU負荷の高い処理でGILの影響を避けたい場合に使います。ワーカーへ渡す関数や引数、戻り値はpickle可能である必要があります。ワーカー関数はモジュール直下に定義します。

from concurrent.futures import ProcessPoolExecutor


def cpu_task(value):
    return sum(i * i for i in range(value))


if __name__ == "__main__":
    values = [100_000, 120_000, 140_000]
    with ProcessPoolExecutor(max_workers=4) as executor:
        results = list(executor.map(cpu_task, values))
    print(results)

asyncio + aiohttp(高並列のAPI呼び出し)

多数の短時間リクエストを扱う場合に有効です。次の例では、1つのClientSessionを再利用し、asyncio.Semaphoreで同時実行数を制限します。aiohttpは標準ライブラリではないため、利用環境へ別途導入する必要があります。

import asyncio

import aiohttp


async def fetch_json(session, url, semaphore):
    async with semaphore:
        timeout = aiohttp.ClientTimeout(total=10)
        async with session.get(url, timeout=timeout) as response:
            response.raise_for_status()
            return await response.json()


async def main(urls):
    semaphore = asyncio.Semaphore(10)
    async with aiohttp.ClientSession() as session:
        tasks = [fetch_json(session, url, semaphore) for url in urls]
        return await asyncio.gather(*tasks)


# results = asyncio.run(main(urls))

asyncio.gather()は、既定ではいずれかのタスクが例外を送出すると、その例外を呼び出し元へ伝えます。項目ごとに失敗を記録する場合は、各タスク内で対象となる例外を処理する設計にします。

3) レート制御・同時接続制限

APIにはレート制限やコストがあるため、同時接続数や時間当たりの呼び出し数を制御します。同時実行数だけを抑えるならSemaphore、時間当たりの許容量まで管理するならトークンバケットなどを検討します。

方法 説明 利点/注意点
Semaphore 同時実行数を単純に制限する 実装が簡単。時間当たりのリクエスト数は直接制御しない。
トークンバケット 時間当たりの許容量を制御する レート制限に適合させやすいが、実装と状態管理がやや複雑。

4) 再試行戦略とbackoff実装

一時的なネットワーク障害に対しては、回数に上限を設けた再試行を使います。すべての例外を一律に再試行するのではなく、失敗理由とAPI仕様に応じて再試行可否を判断してください。

  • 再試行を検討する例:一時的なタイムアウト、接続断、502/503など。
  • 原則として再試行しない例:認証エラー(401)、入力不備(400など)、権限不足。
  • 429の扱い:一時的なレート制限ならRetry-AfterやAPI仕様に従って待機し、上限回数内で再試行する。契約上のクォータ枯渇など、待機しても解消しない429は再試行を打ち切る。

次のコードは、再試行可能と判定済みの失敗だけをRetryableErrorとして受け取り、exponential backoffにjitterを加えて待機します。max_attemptsは初回実行を含む総試行回数です。

import random
import time


class RetryableError(Exception):
    def __init__(self, message, retry_after=None):
        super().__init__(message)
        self.retry_after = retry_after


def backoff_delay(attempt, base=1.0, cap=30.0):
    upper = min(cap, base * (2 ** attempt))
    return random.uniform(0, upper)


def call_with_retry(operation, max_attempts=5):
    for attempt in range(max_attempts):
        try:
            return operation()
        except RetryableError as exc:
            if attempt == max_attempts - 1:
                raise

            if exc.retry_after is not None:
                delay = max(0.0, exc.retry_after)
            else:
                delay = backoff_delay(attempt)

            time.sleep(delay)

実際のAPIアダプター側では、タイムアウトや一時的な502/503、再試行可能な429だけをRetryableErrorへ変換します。恒久的なエラーは変換せず、そのまま呼び出し元へ返します。429のRetry-Afterの形式や単位は利用するAPIの仕様に従って解釈してください。

外部ライブラリを利用できる環境では、tenacityなどで停止条件、待機条件、再試行対象の例外を設定する方法もあります。導入時は利用中のバージョンの公式ドキュメントに合わせて設定してください。

5) 再実行に強い設計:チェックポイント・idempotency・トランザクション的コミット

長時間バッチは途中で中断されることを想定します。チェックポイントと冪等性(idempotency)を組み合わせ、中断後に同じ処理が一部重複しても結果が壊れないようにします。

項目 実務上の設計例
入力の同一性確認 入力ファイルのハッシュ、サイズ、更新日時などを保存し、再開時に同じ入力か確認する。ハッシュだけでは処理済み位置は分からない。
進捗チェックポイント 処理済みの行ID、確定済みレコードID、または安全に再開できるオフセットを保存する。
idempotency APIやDBが対応している場合は、行IDなどの安定した一意キーを使い、重複実行でも二重登録しない仕様にする。
コミット 一時ファイルからの置換、DBトランザクション、確定状態の記録などにより、不完全な出力を完成済みとして扱わない。

並行処理では完了順が入力順と一致しないため、単純に「最後に完了した行番号」だけを保存すると、未完了行を飛ばす可能性があります。処理済みIDを保存するか、順序どおりに確定した連続位置だけをチェックポイントとして記録してください。

6) 優雅な終了と再開(SIGTERM対応・中断保存)

実行中に終了要求を受けた場合は、新しい処理の投入を止め、実行中の処理を可能な範囲で完了させてからチェックポイントを保存します。Unix系環境ではSIGTERMを捕捉する方法が一般的です。

import signal
import threading


stop_requested = threading.Event()


def request_stop(signum, frame):
    stop_requested.set()


if hasattr(signal, "SIGTERM"):
    signal.signal(signal.SIGTERM, request_stop)
signal.signal(signal.SIGINT, request_stop)

シグナルハンドラー内では複雑なI/Oを行わず、終了要求のフラグだけを設定します。メイン処理がそのフラグを確認し、安全な位置でチェックポイントを書き出します。保存頻度はコストとのトレードオフであり、数秒ごと、N件ごと、または成功ごとに保存する方法があります。

7) 監視とログ(運用につなげる)

ログは詳細すぎても不足でも困ります。エラー原因の特定に必要な情報として、リクエストID、行ID、エラーコード、再試行回数、最終的な成否を記録します。認証情報や入力本文などの機密情報をそのまま出力しないよう注意してください。第151回の設定と組み合わせ、CLIで並列度やタイムアウトを外から調整できるようにしておくと便利です。

8) 実践例:CSV→埋め込みAPI→ベクトルDB(最小構成)

ここでは、CSVを逐次読み込み、ThreadPoolExecutorで処理し、成功したIDをチェックポイントへ保存する最小構成を示します。特定事業者のAPIやベクトルDBには依存せず、process_rowへ「API呼び出しと保存」を行う関数を渡す形です。

ステップ ポイント
1. ストリーミング読込 CSVを1行ずつ読み、全件を一度にメモリへ載せない。
2. 同時実行数の制限 保留中のFuture数に上限を設け、読み込み側へバックプレッシャーをかける。
3. API呼び出し タイムアウトを設定し、再試行可能な失敗だけをRetryableErrorへ変換する。
4. 成功時処理 ベクトルDBなどへの保存が成功した後に、処理済みIDをチェックポイントへ追加する。
5. 失敗時処理 最大試行回数を超えた項目をエラーファイルへ記録する。
6. 終了/再開 入力ハッシュを確認し、処理済みIDを除外して再開する。

以下のコードは、前節のcall_with_retrystop_requestedが定義済みであることを前提にしています。CSVには重複しないid列が必要です。

import csv
import hashlib
import json
import os
from concurrent.futures import FIRST_COMPLETED, ThreadPoolExecutor, wait


def file_sha256(path):
    digest = hashlib.sha256()
    with open(path, "rb") as source:
        for chunk in iter(lambda: source.read(1024 * 1024), b""):
            digest.update(chunk)
    return digest.hexdigest()


def load_checkpoint(path, input_sha256):
    if not os.path.exists(path):
        return set()

    with open(path, "r", encoding="utf-8") as source:
        data = json.load(source)

    if data["input_sha256"] != input_sha256:
        raise ValueError("入力ファイルがチェックポイント作成時と異なります")

    return set(data.get("completed_ids", []))


def save_checkpoint(path, input_sha256, completed_ids):
    temporary_path = path + ".tmp"
    data = {
        "input_sha256": input_sha256,
        "completed_ids": sorted(completed_ids),
    }

    with open(temporary_path, "w", encoding="utf-8") as destination:
        json.dump(data, destination, ensure_ascii=False)
        destination.flush()
        os.fsync(destination.fileno())

    os.replace(temporary_path, path)


def append_error(path, row_id, exc):
    record = {
        "id": row_id,
        "error_type": type(exc).__name__,
        "message": str(exc),
    }
    with open(path, "a", encoding="utf-8") as destination:
        destination.write(json.dumps(record, ensure_ascii=False) + "\n")


def run_csv(input_path, process_row, max_workers=8, max_attempts=5):
    checkpoint_path = input_path + ".checkpoint.json"
    error_path = input_path + ".errors.jsonl"
    input_sha256 = file_sha256(input_path)
    completed_ids = load_checkpoint(checkpoint_path, input_sha256)
    max_pending = max_workers * 2

    def submit_row(executor, row):
        return executor.submit(
            call_with_retry,
            lambda: process_row(row),
            max_attempts,
        )

    def collect_finished(pending, block_until_one=False):
        if not pending:
            return

        timeout = None if block_until_one else 0
        done, _ = wait(
            pending,
            timeout=timeout,
            return_when=FIRST_COMPLETED,
        )

        for future in done:
            row_id = pending.pop(future)
            try:
                future.result()
            except Exception as exc:
                append_error(error_path, row_id, exc)
            else:
                completed_ids.add(row_id)
                save_checkpoint(
                    checkpoint_path,
                    input_sha256,
                    completed_ids,
                )

    pending = {}
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        with open(input_path, newline="", encoding="utf-8") as source:
            reader = csv.DictReader(source)

            if not reader.fieldnames or "id" not in reader.fieldnames:
                raise ValueError("CSVにはid列が必要です")

            for row in reader:
                if stop_requested.is_set():
                    break

                row_id = row["id"]
                if not row_id:
                    raise ValueError("idが空の行があります")
                if row_id in completed_ids:
                    continue

                while len(pending) >= max_pending:
                    collect_finished(pending, block_until_one=True)

                future = submit_row(executor, row)
                pending[future] = row_id
                collect_finished(pending)

        while pending:
            collect_finished(pending, block_until_one=True)

process_row(row)は、対象APIを呼び出して埋め込みを取得し、ベクトルDBなどへの保存が完了した時点で正常終了する関数として実装します。一時障害はRetryableErrorとして送出し、認証エラーや入力不備などは再試行対象にしません。

APIまたは保存先がidempotency keyや一意制約を提供している場合は、CSVのidを安定したキーとして利用します。保存成功直後、チェックポイント更新前にプロセスが強制終了すると、再開時に同じ行が再実行される可能性があるためです。外部APIが冪等性を保証しない場合は、事前照会、DBの一意制約、処理状態テーブルなど、保存先に合った重複防止策が必要です。

この最小構成は処理済みIDをJSONへ保存するため、件数が非常に多い場合はチェックポイントファイルも大きくなります。その場合は、DBの状態テーブルや確定済みの連続オフセットなどへ置き換えてください。

9) 注意点と落とし穴

  • GILの存在:CPUバウンドはProcessPoolを検討する。
  • 外部APIのレート制限とコスト:過度な並列は高額請求やAPIブロックを招く。
  • pickleできないオブジェクト:ProcessPoolへ渡す関数、引数、戻り値に注意する。
  • 再試行による重複:書き込みを伴う処理ではidempotencyや一意制約を設計する。
  • チェックポイントの誤用:入力ハッシュは同一性確認用であり、進捗情報の代わりにはならない。
  • ログの過不足:問題追跡に必要な情報を確保しつつ、機密情報を記録しない。

まとめ — 実務で取り入れるための短い手順

まずは小さなサンプルでボトルネックを計測し、I/O、CPU、多接続のどれが原因かを見極めます。I/OバウンドならThreadPoolかasyncio、CPUバウンドならProcessPoolを選択します。再試行は対象となる失敗を限定し、exponential backoffとjitter、APIが指定する待機時間を組み合わせます。入力の同一性確認と進捗チェックポイントを分け、idempotencyによって途中中断から安全に再開できるようにします。運用設定(並列度・タイムアウト・再試行回数)はCLIやconfigで外部化し、監視とログで挙動を把握してください。

読後のアクションリスト(短期)

  • 現行バッチのボトルネック計測を行う
  • 小さなサンプルでThreadPoolまたはasyncioを試す
  • 再試行対象と再試行しないエラーを明文化する
  • 入力の同一性情報と進捗チェックポイントを分けて保存する
  • 本番前に低負荷で運用テストを行い、ログと再開動作を確認する

次回は、実際のCLIテンプレートと組み合わせた簡易フレームワークを提示し、Cronやスケジューラとの安全な連携方法を紹介します。シリーズ「AIとPythonの実務」の流れで、実務で使える小さな改善を積み重ねていきましょう。

参考リンク:Manage AI 第151回(logging/CLI/config)の設定を先に確認すると運用へスムーズにつなげられます。

第151回 実務で使えるPython基礎:logging・argparse・configparserで作る現場向けCLIツールと運用ログ設計

ちょっとしたスクリプトは「動けば良い」となりがちで、運用に回すとログが足りなかったり、設定がハードコーディングされて苦労したりします。本記事では、小規模なAI/データ処理を現場で安全に運用するためのCLI設計とログ設計を、Python 3.9以降の標準ライブラリだけで動かせるテンプレートとチェックリストでまとめます。コードはWordPressに貼り付けられる<pre><code>形式で掲載します。

1)なぜ小さなCLIツールと良いログ設計が必要か(現場の失敗例)

現場でよくあるつまずきと影響を整理します。

よくある問題 現場での影響
引数が曖昧でHelpが無い 使う人が誤操作しやすく、自動化に組み込みにくい
設定がコード内にハードコーディング 環境ごとに編集が必要で漏洩リスクがある
ログがプレーンテキストで機械処理しにくい 障害検出や集約が難しく、原因追跡が遅れる
エラーが標準出力に散らばる 例外時の再現性が低く、対応が困難になる

2)要件整理:引数・設定・出力・監視

設計を始める前に満たすべき要件を明確にします。

領域 要件(実務的観点)
引数 必須・任意を明確にし、サブコマンドで処理を分離する。-h/--helpを充実させる
設定 設定ファイル(configparser)を環境変数で上書きする。秘密情報は設定ファイルに保存しない
出力/ログ コンソールは人間向け、ファイルはJSON Lines形式で保存し、ローテーションを設定する
監視 相関ID、ログレベル、処理時間、エラーを検索できるようにする

3)argparseで作る使いやすいCLI

完成テンプレートでは、入力ファイルの空行を除いた行数を数え、結果をJSONファイルへ保存するprocessサブコマンドを用意します。--input--outputは必須です。

グローバルオプションの--config--correlation-idは、次のようにサブコマンドより前へ指定します。

python cli_tool.py --config config.ini --correlation-id demo-001 process --input input.txt --output result.json

ヘルプは以下のコマンドで確認できます。

python cli_tool.py -h
python cli_tool.py process -h

必須引数が不足している場合など、引数解析のエラーはargparseが標準エラー出力へ表示し、終了コード2で終了します。

4)configparserと環境変数で安全に設定を読み込む手順

ログファイルの保存先とログレベルをINIファイルで管理します。たとえば、config.iniを次の内容で作成します。

[app]
log_file = logs/cli.jsonl
log_level = INFO

完成コードでは、APP_LOG_FILEAPP_LOG_LEVELが設定されていれば、INIファイルの値より環境変数を優先します。設定ファイルを指定しない場合は、コード内の非機密な初期値を使用します。

APP_LOG_LEVEL=DEBUG python cli_tool.py --config config.ini process --input input.txt --output result.json

APIキーなどの秘密情報はINIファイルへ書かず、実際にAPI処理を追加する段階でos.environまたはos.getenv()から直接取得してください。このテンプレート自体は外部APIを呼び出さないため、APIキーは要求しません。

5)logging入門:ハンドラとフォーマッタ(プレーン/JSON)

ポイントは「コンソールは人間向け、ファイルは機械処理向け」の使い分けです。完成コードでは、ルートロガーに次の2つのハンドラを設定します。

出力先 設定
コンソール StreamHandlerを使用し、INFO以上を読みやすいプレーンテキストで出力
ログファイル RotatingFileHandlerと独自のJSONフォーマッタを使用し、1行1JSONで出力

ログファイル側のレベルはlog_levelで変更できます。初期値はINFOです。DEBUGを指定した場合でも、コンソールはINFO以上、ファイルはDEBUG以上という使い分けになります。

6)実践:ログ回転、例外の一元キャッチ、相関IDの付与

完成コードでは、ログファイルが5MBに達したらローテーションし、過去5ファイルを保持します。maxBytesbackupCountは運用条件に合わせて調整してください。

相関IDはcontextvars.ContextVarに保存します。--correlation-idが指定されていない場合は、CLIの実行ごとにUUIDを生成します。同じ実行中に出力されたログを、共通のcorrelation_idで検索できます。

コマンド処理中の例外はmain()で一元的に捕捉し、logger.exception()でスタックトレースを記録して終了コード1を返します。Ctrl+CなどのKeyboardInterruptは終了コード130で終了します。

7)入出力とJSONログの確認例

input.txtを次の内容で作成します。

alpha

beta

設定ファイルを使って実行します。

python cli_tool.py --config config.ini --correlation-id demo-001 process --input input.txt --output result.json

result.jsonには次の結果が保存されます。

{
  "input": "input.txt",
  "non_empty_lines": 2,
  "correlation_id": "demo-001"
}

logs/cli.jsonlには、開始、処理完了、コマンド成功のログが1行1JSONで記録されます。時刻や処理時間は実行ごとに変わります。

{"timestamp":"2025-01-01T00:00:00+00:00","level":"INFO","logger":"cli_tool","message":"processing_finished","correlation_id":"demo-001","duration_ms":1.25,"input_path":"input.txt","output_path":"result.json","non_empty_lines":2}

入力ファイルが存在しない場合は、例外情報を含むERRORログが出力され、コマンドは終了コード1で終了します。

8)モニタリング連携の入口(ログからSLO/アラートにつなげる方法)

まずはログ構造を整え、次のチェックリストに従ってください。

目的 短いチェックリスト
エラー検知 JSONログのlevelcorrelation_idexceptionを検索できるか確認する
パフォーマンス監視 processing_finishedduration_msを集計し、平均値や上位値を確認する
アラート設計 短期のエラーバーストと、長期のエラーレート上昇を別の条件として扱う

9)デプロイと運用チェックリスト

CIから本番運用までの最低限の確認項目です。

カテゴリ チェック項目
CI 必要に応じてflake8やBlackによるチェックを行い、主要フローと異常系をテストする
設定 秘密情報は環境変数などで注入し、設定ファイルには保存しない。本番の設定ファイルは不用意に変更できない権限にする
実行 終了コードを監視し、systemdやコンテナを利用する場合は再起動条件とログ転送先を設定する
ログ ローテーション後のファイル数と容量を確認し、必要な保存期間に合わせて調整する

10)完成コード(そのまま貼れる短縮テンプレート)

以下をcli_tool.pyとして保存してください。CLI、設定読み込み、環境変数による上書き、JSONログ、ログ回転、相関ID、例外処理、入出力処理を含むテンプレートです。

#!/usr/bin/env python3
import argparse
import json
import logging
import os
import sys
import time
import uuid
from configparser import ConfigParser
from contextvars import ContextVar
from datetime import datetime, timezone
from logging.handlers import RotatingFileHandler
from pathlib import Path


correlation_id_var = ContextVar("correlation_id", default="-")
logger = logging.getLogger("cli_tool")


class JsonFormatter(logging.Formatter):
    """logging.LogRecordを1行のJSONへ変換する。"""

    def format(self, record):
        payload = {
            "timestamp": datetime.fromtimestamp(
                record.created, tz=timezone.utc
            ).isoformat(),
            "level": record.levelname,
            "logger": record.name,
            "message": record.getMessage(),
            "correlation_id": correlation_id_var.get(),
        }

        for key in (
            "duration_ms",
            "input_path",
            "output_path",
            "non_empty_lines",
        ):
            value = getattr(record, key, None)
            if value is not None:
                payload[key] = value

        if record.exc_info:
            payload["exception"] = self.formatException(record.exc_info)

        return json.dumps(payload, ensure_ascii=False)


def parse_args():
    parser = argparse.ArgumentParser(
        description="入力ファイルを処理し、結果と運用ログを出力します。"
    )
    parser.add_argument(
        "--config",
        type=Path,
        help="INI形式の設定ファイル",
    )
    parser.add_argument(
        "--correlation-id",
        help="ログを関連付ける相関ID。省略時はUUIDを生成します。",
    )

    subparsers = parser.add_subparsers(
        dest="command",
        required=True,
    )

    process_parser = subparsers.add_parser(
        "process",
        help="入力ファイルの空行を除いた行数を集計します。",
    )
    process_parser.add_argument(
        "--input",
        type=Path,
        required=True,
        help="UTF-8の入力テキストファイル",
    )
    process_parser.add_argument(
        "--output",
        type=Path,
        required=True,
        help="結果を書き込むJSONファイル",
    )
    process_parser.set_defaults(handler=run_process)

    return parser.parse_args()


def load_config(config_path):
    parser = ConfigParser()
    parser.read_dict(
        {
            "app": {
                "log_file": "logs/cli.jsonl",
                "log_level": "INFO",
            }
        }
    )

    if config_path is not None:
        if not config_path.is_file():
            raise FileNotFoundError(
                f"設定ファイルが見つかりません: {config_path}"
            )
        with config_path.open("r", encoding="utf-8") as file:
            parser.read_file(file)

    return {
        "log_file": os.getenv(
            "APP_LOG_FILE",
            parser.get("app", "log_file"),
        ),
        "log_level": os.getenv(
            "APP_LOG_LEVEL",
            parser.get("app", "log_level"),
        ),
    }


def get_log_level(level_name):
    level = getattr(logging, level_name.upper(), None)
    if not isinstance(level, int):
        raise ValueError(f"不正なログレベルです: {level_name}")
    return level


def setup_logging(log_file, log_level):
    file_level = get_log_level(log_level)
    log_path = Path(log_file)
    log_path.parent.mkdir(parents=True, exist_ok=True)

    root_logger = logging.getLogger()
    root_logger.setLevel(logging.DEBUG)

    for handler in root_logger.handlers[:]:
        root_logger.removeHandler(handler)
        handler.close()

    console_handler = logging.StreamHandler(sys.stderr)
    console_handler.setLevel(logging.INFO)
    console_handler.setFormatter(
        logging.Formatter(
            "%(asctime)s %(levelname)s "
            "correlation_id=%(correlation_id)s %(message)s"
        )
    )
    console_handler.addFilter(CorrelationIdFilter())

    file_handler = RotatingFileHandler(
        log_path,
        maxBytes=5 * 1024 * 1024,
        backupCount=5,
        encoding="utf-8",
    )
    file_handler.setLevel(file_level)
    file_handler.setFormatter(JsonFormatter())

    root_logger.addHandler(console_handler)
    root_logger.addHandler(file_handler)


class CorrelationIdFilter(logging.Filter):
    def filter(self, record):
        record.correlation_id = correlation_id_var.get()
        return True


def setup_fallback_logging():
    """設定読み込み前の例外を標準エラー出力へ記録する。"""
    root_logger = logging.getLogger()
    if root_logger.handlers:
        return

    root_logger.setLevel(logging.INFO)
    handler = logging.StreamHandler(sys.stderr)
    handler.setFormatter(
        logging.Formatter(
            "%(asctime)s %(levelname)s "
            "correlation_id=%(correlation_id)s %(message)s"
        )
    )
    handler.addFilter(CorrelationIdFilter())
    root_logger.addHandler(handler)


def run_process(args):
    started_at = time.perf_counter()

    logger.info(
        "processing_started",
        extra={
            "input_path": str(args.input),
            "output_path": str(args.output),
        },
    )

    text = args.input.read_text(encoding="utf-8")
    non_empty_lines = sum(
        1 for line in text.splitlines() if line.strip()
    )

    result = {
        "input": str(args.input),
        "non_empty_lines": non_empty_lines,
        "correlation_id": correlation_id_var.get(),
    }

    args.output.parent.mkdir(parents=True, exist_ok=True)
    args.output.write_text(
        json.dumps(result, ensure_ascii=False, indent=2) + "\n",
        encoding="utf-8",
    )

    duration_ms = round(
        (time.perf_counter() - started_at) * 1000,
        2,
    )
    logger.info(
        "processing_finished",
        extra={
            "duration_ms": duration_ms,
            "input_path": str(args.input),
            "output_path": str(args.output),
            "non_empty_lines": non_empty_lines,
        },
    )


def main():
    args = parse_args()
    correlation_id = args.correlation_id or str(uuid.uuid4())
    token = correlation_id_var.set(correlation_id)

    try:
        config = load_config(args.config)
        setup_logging(
            config["log_file"],
            config["log_level"],
        )

        logger.info("command_started")
        args.handler(args)
        logger.info("command_succeeded")
        return 0
    except KeyboardInterrupt:
        setup_fallback_logging()
        logger.warning("command_interrupted")
        return 130
    except Exception:
        setup_fallback_logging()
        logger.exception("command_failed")
        return 1
    finally:
        correlation_id_var.reset(token)


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

11)次の一歩

まずは上記コード、config.iniinput.txtを同じディレクトリに保存し、記載したコマンドで実行してください。その後、正常系と入力ファイルが存在しない異常系の両方を試し、result.json、コンソール出力、logs/cli.jsonl、終了コードを確認します。

ログを一定期間保存して検索できることを確認したら、必要に応じてログ集約基盤やメトリクス監視へつなげます。外部ライブラリを導入する前に、まずは標準ライブラリで引数、設定、ログ、例外、終了コードの一貫した設計を作ることが重要です。

まとめ

  • CLIはargparseで設計し、サブコマンド、必須引数、-h/--helpを用意する。
  • 非機密設定はconfigparserで管理し、環境変数による上書きを可能にする。
  • ログはコンソールとJSONファイルを使い分け、ローテーションと相関IDを設定する。
  • 例外を一元的に記録し、自動実行側が判定できる終了コードを返す。
  • 処理時間やエラーをJSONログから集計し、監視やアラートへつなげる。

小さな自動化スクリプトほど、運用方法を最初に整えることで効果が持続します。まずはテンプレートを実行し、現場の保存期間や監視条件に合わせて調整してください。

第150回 実務で使えるPython基礎:dataclassと型注釈で作る安全な表データスキーマとシリアライズワークフロー

はじめに — 「行データ」がいつも曲者に感じるあなたへ

CSVやJSONLで扱う「一行」の揺らぎ(欠損、型のばらつき、日時表記の差異)は、現場で繰り返し問題になります。辞書で取り回すと可読性が下がり、変換ルールが散らばって保守が難しくなります。この記事では、dataclass と型注釈を使って行データを型付きオブジェクトに落とし込み、変換・検証・シリアライズを一貫して扱う実務的な手順を示します。第149回のストリーミング埋め込みの続きとして、パイプラインに組み込みやすい層を作ることが狙いです。

1) なぜ dataclass と型注釈か — 現場でのメリットと適用範囲

短く言うと「可読性」「静的解析」「移行しやすさ」です。主な利点を表にまとめます。

項目 辞書ベース dataclass + 型注釈
可読性 キー文字列に依存しがちで散らばる フィールド名でまとまる。型が明示される
変換位置 処理ごとに分散しやすい 一箇所で変換・正規化できる
型チェック 実行時まで不明 mypy 等で静的チェックしやすい
互換性対策 キーの存在/欠損の追跡が難しい バージョンやデフォルトで一元管理可能

適用範囲は「テーブル状の行データ」を扱う処理、ETL の前段、埋め込み生成前の正規化、ログ行の正規化などです。大量行の超高頻度処理では dataclass のオーバーヘッドを考慮する必要があります(後述)。

2) 基本パターン:単純な行 → dataclass の定義と変換

型注釈を用いた基本例を示します。ここでは型には typing を使います。

例:行スキーマ
フィールド: id (str), created_at (datetime|None), score (Optional[float]), tags (List[str])

変換の雛形(読み取り → dataclass)は次のように整理できます。

雛形コード(概念)
from dataclasses import dataclass, field
from typing import Any, List, Mapping, Optional
from datetime import datetime

@dataclass
class Row:
id: str
created_at: Optional[datetime]
score: Optional[float]
tags: List[str] = field(default_factory=list)

# CSV/JSON の dict を Row に変換するユーティリティ
def dict_to_row(d: Mapping[str, Any]) -> Row:
# 日付文字列→datetime、空文字→None、tags を split などの正規化を行う
row_id = d[‘id’]
if not isinstance(row_id, str):
raise TypeError(‘id must be a string’)
created = parse_date_or_none(d.get(‘created_at’))
score = parse_float_or_none(d.get(‘score’))
tags = parse_tags(d.get(‘tags’))
return Row(id=row_id, created_at=created, score=score, tags=tags)

実務では parse_* 関数を小さく分けておくと再利用しやすく、単体テストも書きやすくなります。

3) 欠損値・デフォルト・型変換の実務ハンドリング

欠損値や空文字の扱いは現場で最も差が出る部分です。型が確定した値の内部整形には __post_init__ を利用できますが、CSVやJSONから来る文字列などの外部入力は、入力型を広く取る専用ファクトリ関数で受けて一元的に変換すると、dataclass の型注釈との矛盾を避けられます。

パターン 説明
__post_init__ 型に適合した値でインスタンス化した直後に、内部整形やフィールド間の整合性確認を行う。
ファクトリ関数 外部データ源(CSV/JSON)に特化した入力型と変換を分離。テストしやすい。
ユーティリティ関数 日付パーサ、数値パーサ、リスト正規化などを小さく作る。

実例(外部入力を専用ファクトリで受ける構成):

コードスニペット
from dataclasses import dataclass, field
from typing import Any, List, Mapping, Optional
from datetime import datetime

@dataclass
class Row:
id: str
created_at: Optional[datetime]
score: Optional[float]
tags: List[str] = field(default_factory=list)

@classmethod
def from_dict(cls, d: Mapping[str, Any]) -> “Row”:
row_id = d[‘id’]
if not isinstance(row_id, str):
raise TypeError(‘id must be a string’)

return cls(
id=row_id,
created_at=parse_date_or_none(d.get(‘created_at’)),
score=parse_float_or_none(d.get(‘score’)),
tags=parse_tags(d.get(‘tags’)),
)

# 必要に応じて既存の変換関数からファクトリを呼び出す
def dict_to_row(d: Mapping[str, Any]) -> Row:
return Row.from_dict(d)

4) ネスト・リスト・可変長フィールドの扱い

ネストした構造は dataclass をネストして表現します。可変長フィールドは List 型で表し、デフォルトは default_factory を使います。

from dataclasses import dataclass, field
from typing import List, Optional

@dataclass
class Item:
name: str
qty: int

@dataclass
class Row:
id: str
items: List[Item] = field(default_factory=list)

ネストの変換は再帰的にファクトリを呼び出すか、専用の parse_item_list 関数を用意して一元化します。

5) CSV/JSONL ⇄ dataclass シリアライズ/逆シリアライズ(ストリーミング対応)

大量行を扱う場合、メモリに全部読まないストリーミング処理が実務では重要です。読み取り側は行を順次 yield するジェネレータ、書き出し側は iterable を順次処理してファイルへ直接書き込む関数として構成できます。

CSV → dataclass(ジェネレータ)
import csv

def stream_rows_from_csv(fp):
reader = csv.DictReader(fp)
for d in reader:
try:
yield dict_to_row(d)
except Exception as e:
# ログに残してスキップか再試行のルールをここで適用
handle_conversion_error(d, e)

JSONL の場合は一行ずつ json.loads して同様に yield します。逆方向(dataclass → CSV/JSONL)は、受け取った行を一件ずつストリームへ書き出します。

dataclass → JSONL(ストリームへの逐次書き込み)
import json

def stream_write_jsonl(rows, fp):
for row in rows:
obj = asdict_for_serialization(row) # 日付は ISO 化など
fp.write(json.dumps(obj, ensure_ascii=False) + “\n”)

ポイント:

  • 日付は統一フォーマット(例: ISO 8601)で保存する
  • 列順を固定したい場合は列名リストを管理して出力順を制御する
  • 圧縮出力(gzip)を行うときはバッファリングと逐次処理を組み合わせる

6) 軽量バリデーション戦略と pydantic の使い分け

dataclass と小さな検証ロジックで十分なケースと、堅牢なランタイム検証が必要なケースは分けて考えます。

目的 推奨
軽量な正規化・開発効率重視 dataclass + 小さな parse/validate 関数
外部入力が不安定で安全性重視 pydantic によるランタイム検証、または attrs に validators/converters を明示的に実装する
静的解析と型互換チェック mypy と型注釈、テストで補完

attrs は型注釈を付けるだけで実行時の型検証を行うものではないため、必要な検証や変換は validators や converters などで明示します。pydantic も便利ですが、依存と検証・シリアライズの挙動を理解した上で使うことが重要です。検証エラーの記録やエラーレート閾値は運用で役立ちます。

7) テストとCIで防ぐ典型的な運用エラー

テスト設計は次の要素を含めます。

  • 変換ユーティリティの単体テスト(正常系/欠損/誤フォーマット)
  • 境界値テスト(長すぎる文字列、大きな数値、深いネスト)
  • パラメータ化テストで複数のフォーマット(日付形式など)をカバー
  • モック CSV/JSONL を使った E2E テスト(読み取り→変換→書き出し)
  • スキーマ差分検出テスト:古いスキーマ vs 新スキーマで互換性チェック

CI に入れるべき簡易スクリプト例:スキーマ差分チェック(型名と必須フィールドの差を検出)を自動化しておくと、運用時の誤変更を防げます。

8) 運用チェックリストと移行・バージョン管理の実践

運用フローに組み込む際の最小チェックリスト:

項目 説明
スキーマバージョン 各行に schema_version をメタデータで持たせる
エラー率モニタ 変換エラー率が閾値超えなら自動アラート
サンプル検査 一定割合のサンプル出力を人が確認するルール
再処理ルール 失敗行の隔離と再実行手順をドキュメント化
移行パス フィールド追加は Optional で始め、後で必須に移行する計画

既存パイプラインへの挿入例(第148/149回との接続):

  • 埋め込み生成直前に dataclass 層で正規化→埋め込み入力が安定する
  • スキーマバージョンをメタデータとして保存し、再処理時に変換ルールを選べるようにする

現場でよくある落とし穴と対策(短めチェックリスト)

問題 対策
パフォーマンスの低下 プロファイリングでホットスポットを特定、必要なら C もしくは vectorized 処理に切り替え
日時フォーマットのばらつき parse_date_or_none に複数フォーマット順試行を実装、ログで未対応フォーマットを収集
型の過信 外部入力は常に検証。pydantic を補助的に導入
スキーマ変更で壊れる バージョニングと互換性レイヤ(Optional→必須の移行プラン)

テスト設計の具体案(簡易)

代表的なテストケース例:

  • 正常行の変換が期待通りに行われる
  • 空文字・null が Optional フィールドにマップされる
  • 不正な日付はログに残して行をスキップ(あるいは None)にする
  • 大きなリスト(1000 要素)の items を扱えるか

CI ではこれらを pytest のパラメータ化で回し、カバレッジと型チェック(mypy)を組み合わせます。

まとめ — 実務で使うための短い指針

  • dataclass + 型注釈は「読みやすさ」と「移行性」を高める。まずは小さなスキーマで試す。
  • 外部入力の変換ルールは、入力型を広く取る専用ファクトリに集約し、ユーティリティ関数で分割する。
  • 読み取りはジェネレータ、書き出しは逐次処理でメモリを抑え、エラー記録と再実行ルールを運用に組み込む。
  • ランタイム検証が必要な領域は pydantic などを使い分ける。静的検査(mypy)とテストで品質を担保する。

次に繋げるトピック案

  • pydantic と attrs(validators/converters)を用いた検証の実務比較
  • mypy を組み込んだ CI の実装例
  • スキーママイグレーション自動化と再処理オーケストレーション

今回示した考え方と雛形は、現場での小さなミスや運用コストを減らすための実務改善です。まずは一つの CSV/JSONL パイプラインに導入して、変換エラー率や可読性の改善を測定してみてください。次回は pydantic を使ったランタイム検証の深掘りを予定しています。

第149回 実務で使えるPython基礎:ジェネレータとイテレータで設計するメモリ効率の良い埋め込み&検索パイプライン

はじめに:現場でつまずきやすい点に寄り添って

大きなCSVやログを扱うと、いつの間にかメモリが増えてプロセスが落ちる、あるいは埋め込みAPIに一度に大量送信して失敗する――こうした現場あるあるに悩んでいませんか。第148回を読んで埋め込みの概念は理解したが、いざ現場で大規模データを安定運用する段になると「どう設計し、どう復旧するか」が課題になりがちです。本記事では、Pythonのジェネレータ/イテレータを中心に、メモリを一定に保ちながら埋め込み→ベクトルDB登録→簡易検索までつなぐための設計上の手順とチェックリストを提示します。特定の埋め込みAPIやベクトルDBに依存しない、概念と実装方針の解説です。

学習ゴール(この記事を読み終えたとき)

  • ストリーミング読み出しとバッチ化でメモリを一定に保つ設計がわかる
  • 埋め込みAPI呼び出しの安全なリトライやレート制御の考え方を実務で使える形で理解できる
  • 障害時の途中再開(チェックポイント)や簡易メトリクスの実装ポイントが掴める
  • 小さなデータでのテスト方法とCIに組み込む観点が分かる

ステップ1:概念とパターン(ジェネレータ/イテレータの実務的理解)

ジェネレータの基本と利点

Pythonのジェネレータは、yieldで値を逐次返す仕組みです。データ全体をメモリに展開せずに「必要な分だけ処理する」ことで、長時間実行のバッチ処理や大規模ログ処理で安定化します。実務では次の点を押さえます。

  • 遅延評価:入力を1行ずつ(あるいは小さなチャンク単位で)処理することでメモリ使用量が一定に近づく。
  • パイプライン合成:読み出し→前処理→バッチ化→API送信 の各段をジェネレータで繋ぐと、各段の責務が明確になる。
  • 単体テストしやすい:入力ストリームをモック化してジェネレータ単体を検証できる。

よく使うパターン

  • chunked(固定サイズバッチ): バッチサイズNごとにまとめる
  • sliding window(可変・重複あり): 文脈ウィンドウが必要なとき
  • filter/mapチェーン: 前処理→フィルタ→正規化を逐次的に行う

メモリ使用イメージ比較

処理方法 主な特徴 メモリ使用の見積もり
一括ロード(pandas.read_csv) 実装は簡単だが全件をメモリに保持 データサイズ+行・列のオーバーヘッド(高)
ジェネレータで逐次処理 行単位や小チャンク単位で処理。低メモリ。 バッファサイズ+バッチサイズに依存(一定)
ハイブリッド(分割して複数ファイル) ファイル分割で並列処理。運用と整合性管理が必要。 プロセスごとの最大バッファを足し合わせた値

ステップ2:実装ガイド(現場で使える設計と小さな実装ヒント)

1) ストリーミング読み出しジェネレータ(CSV/JSONL/ログ)

ポイントは「逐次読取」と「最小限のパース」。簡単な考え方は次の通りです。

  • CSV: ファイルを逐行読み、必要な列だけパースしてdictをyieldする。
  • JSONL: 各行をjson.loadsしてyield。パース失敗はログに出してスキップ可能。
  • 圧縮ファイル: gzip.open や streaming ライブラリで逐次解凍しながら処理。

(実装ヒント)関数名の例:
read_csv_rows(file_path) → yield {‘id’: …, ‘text’: …}

2) 固定サイズ/可変サイズバッチの作り方

固定サイズバッチ(chunked)は最も実用的です。イメージは次の通り。

用途 実装方針 注意点
固定サイズバッチ ジェネレータからN個ずつ取り出してyield 最終バッチはサイズ未満になることを許容する
可変バッチ(サイズ上限のみ) 文字数やトークン数で上限を設定して詰める トークン推定が必要でやや複雑

(実装ヒント)chunkedの疑似手順: イテレータをループして一時リストにappend、サイズ到達でyieldしクリア。

3) 埋め込みAPI呼び出し:同期/非同期の使い分けと安全なリトライ

実務では次の観点で選択します。

  • 同期: 単純なワンオフ処理やデバッグ時に便利。失敗時の影響範囲が分かりやすい。
  • 非同期(async/await): 高スループットが必要な場合。イベントループ設計とエラーハンドリングが必要。

リトライ戦略(実務向け):

  • タイムアウト設定を必ず設ける(APIごとに最適値)。
  • 指数バックオフ+最大試行回数。短時間の一時障害を吸収する。
  • レート制御: APIのスロットリング(固定インターバル)かトークンバケット方式を採用。
  • 失敗時の扱い: 一時的失敗はバッチ単位で再試行、恒久的エラーはログと失敗キューに記録。

(実装ヒント)簡易流れ: chunkedジェネレータ → API送信(タイムアウト+リトライ)→ 成功した埋め込みは次段へyield

ステップ3:結合と耐障害設計

埋め込み失敗時の戦略

  • 短時間の失敗: バッチ単位で再試行(指数バックオフ)。再試行回数はビジネス要件で決める。
  • 恒久エラー(フォーマット不備等): そのレコードをスキップして失敗ログに保管。人手確認用のCSV/JSONとして残す。
  • 失敗キュー: 再処理用ファイルや小さなDBに失敗レコードを貯め、別ジョブで再投入する。

途中再開(チェックポイント)の実装案

長時間ジョブでは途中から再開できることが重要です。シンプルな手法:

  • 進捗マーカー(例: 最終処理行番号、最終ID)を定期的にファイルに書く。
  • 各バッチの登録完了後にチェックポイントを更新。チェックポイントは冪等に保つ(同じIDで再登録しない工夫)。
  • ベクトルDB側にアップサート(upsert)やidempotentなID設計をすることで重複登録を避ける。

ロギングと簡易メトリクス

指標 何を見るか 取得方法
処理レート records/s や batches/min 処理した件数を時間で割る(ジョブ内でカウント)
メモリ使用 プロセスのRSSやコンテナメモリ psutilやコンテナ監視で定期取得
失敗率 総バッチ数に対する失敗バッチ数 ログから集計、アラート閾値を設定

ステップ4:ベクトルDB登録と簡易検索

少量インサート vs バルクインサート

  • 少量インサート(逐次登録): レイテンシは低めだが総オーバーヘッドが大きくスループットが落ちる。
  • バルクインサート: スループット重視。失敗時のロールバック/再試行戦略が重要。

ID管理: アプリ側で一意のIDを発行しておくと、途中再開や重複検出が容易になります(例: ハッシュ+元ファイルの行番号)。

簡易検索とスコア調整

  • 検索テスト: 少量データでレイテンシと順位を確認(kとスコア閾値の調整)。
  • スコアリングの観察点: 同一ドメインの類似度分布を確認し、閾値を決める。

テストとCIのポイント

  • ジェネレータ単体テスト: ファイルストリームをモックして逐次出力を検証する。
  • 埋め込みAPIはモック化して遅延・エラーをシミュレートしたテストを用意する。
  • メモリ回帰テスト: 小データと大データでメモリ使用を比較し、閾値超過の検出をCIに組み込む。
  • E2Eテスト: 小さなサンプル(例: 100件)で読み出し→埋め込み→登録→検索までの流れを自動化する。

運用の現実的注意点と伸ばしどころ

  • Dockerでのメモリ制限: コンテナ側のメモリ上限を設定し、OOM発生時の挙動を確認する。
  • 長時間ジョブのタイムアウト: ジョブ管理(AirflowやKubernetes CronJobなど)で最大実行時間を設定。
  • 監視指標: 処理レート、メモリ、失敗率、最終正常実行時刻を継続監視する。
  • 将来的な拡張: RAGやキャッシュ戦略との連携、並列処理によるスケールアウトの設計を検討する。

実務チェックリスト

項目 確認ポイント
入力ストリーム 逐次読み出しでメモリ固定化されているか。パースエラーはログ化されるか。
バッチ設計 バッチサイズとレート制御が現場負荷に合っているか。
リトライ/バックオフ タイムアウト・最大試行回数・バックオフが設定されているか。
チェックポイント 定期的に進捗を保存し、再開時に重複登録を防げるか。
ベクトルDB整合性 ID管理とupsert戦略があるか。
テスト ジェネレータ単体・APIモック・E2EのテストがCIに含まれているか。
監視 処理レート、メモリ、失敗率をダッシュボードで確認できるか。

まとめ(実務でまず手を付ける優先順)

  • 1: まずはストリーミング読み出し(ジェネレータ)に置き換え、メモリ問題を排除する。
  • 2: 固定サイズバッチで埋め込みを試し、同期→非同期の切替はスループット要件に応じて行う。
  • 3: チェックポイントとID管理を実装して途中再開と整合性を確保する。
  • 4: ロギングと簡易メトリクスを整備して運用監視を始める。

次回候補として『RAGの品質評価とログ駆動の改善ループ』を想定しています。本稿で挙げたチェックポイントやログは、RAG運用に移行する際の良い出発点になります。ここで示したパターンは過度に特殊化せず、現場の多様なデータ形態に適用できる実務的な基準を目指しました。小さなサンプルでまず動かし、観測→改善を繰り返すことをおすすめします。

第148回 CSV逐次読み込み→断片化→モック埋め込み→SQLite保存の最小構成

CSV逐次読み込み→断片化→モック埋め込み→SQLite保存の最小構成

注意(必読)

この記事内で説明する埋め込みクライアントは配線確認用のモック実装です。実API接続、認証、レート制御、堅牢な再試行・監視などの実務要件は実装していません。本文は「CSVを逐次読み込み→断片化→埋め込み(モック)→SQLiteへ保存→小規模検索へ渡す」ための最小構成を分かりやすく説明することを目的とします。示す設計説明や検証フローは実装ガイドとして有用ですが、示した実装例がそのまま本番で動くとは受け取らないでください。

目的と範囲

大きなCSVを一括でメモリに載せず逐次処理し、各行(あるいは行から生成した断片)単位で埋め込みを作ってSQLiteに保存する。保存したベクトルはベクトル配列とメタ情報(row_id, fragment_id, model)を同じ順序で読み出し、小規模検索(Faiss等)へ渡せるようにする、という最小構成を示します。主キーや保存形式、テストでの検証ポイントを明確にします。

前提・依存

  • 想定するライブラリ(説明用): Python標準の csv/json/sqlite3、numpy(float32変換)、および任意で faiss。
  • 保存ファイルは1モジュール(例: pipeline.py)としてまとめる前提で説明します。
  • 埋め込みクライアントはモック。実APIに差し替える場合は認証・レート制御・再試行設計を追加してください。

処理の順序(概観)

  1. CSVをストリーミングで1行ずつ読み込む(メモリに全件を保持しない)。
  2. 各行を断片(fragment)に変換する。断片は少なくとも以下を含む: row_id, fragment_id, text, metadata。
  3. 断片をバッチ化(クライアントの batch_size に応じる)して埋め込みAPI(ここではモック)に投げる。
  4. 返却されたベクトルを float32 のバイナリに変換し、metadata(JSON)と共に SQLite に保存する。主キーは複合 (row_id, fragment_id, model) を採用する。
  5. 必要に応じて、SQLite からベクトルを ORDER BY row_id, fragment_id の順に読み出し、NumPy 配列(shape: N×D)と対応するレコード配列を得る。これをインデクシング(Faiss 等)や線形探索に渡す。

各段階で保存・保持する情報(明確化)

  • 断片変換段階: row_id(必須)、fragment_id(ユニーク化)、text(埋め込み対象)、metadata(title 等の補助情報)。
  • 保存段階(SQLite): row_id, fragment_id, model を複合主キーとして、vector(float32 のバイナリ)、metadata(JSON文字列)、タイムスタンプを格納する。
  • 読み出し段階: vector は NumPy の float32 配列に復元し、records には row_id/fragment_id/metadata を同じ順序で保持する(位置対応の保証)。

主要な関数・クラス(役割と相互関係)

ここではコード本文は掲載しませんが、実装時に最低限用意するべき関数とその役割を示します。すべて1ファイルにまとめる想定です。

  • stream_csv_rows(path): CSV を逐次的に読み、1行ずつ辞書で返すイテレータ。
  • row_to_fragment(row): CSV の行を断片に変換。id 列の存在を検証し、row_id, fragment_id, text, metadata を返す。
  • chunked(iterable, size): イテラブルを指定サイズで区切るユーティリティ。
  • EmbeddingClient クラス: モック埋め込みクライアント。embed_batch(texts) と embed_batch_with_retry(texts) を提供し、バッチ単位でベクトル(例: 128次元のランダム値)を返す。
  • Create_schema(conn): SQLite にテーブルを作成する。主キーは (row_id, fragment_id, model)。
  • save_embedding(conn, row_id, fragment_id, model, vector, metadata): ベクトルを float32 のバイナリに変換して保存。重複は上書き(INSERT OR REPLACE 相当)する運用を想定。
  • load_embeddings(conn, model): 指定モデルの行を ORDER BY row_id, fragment_id で取得し、(vectors: NumPy 配列, records: メタ情報リスト) を返す。np.frombuffer の場合は安全のためコピーする。
  • embed_csv(path, conn, client, model): 上記を組み合わせる本体。CSV をストリームし、断片化→バッチ埋め込み→保存→コミットを繰り返す。
  • search_sqlite(…): 小規模検証用の検索。クエリを埋め込み、Faiss が使える環境なら IndexFlatL2 を作成して検索、未導入なら単純な線形探索で距離を計算する。

pytest による最小検証の流れ(コードは掲載しません)

実際のテストでは以下の手順でパイプラインの基本動作を検証します。テストは pipeline.py を同一ディレクトリに保存する前提です。

  1. 一時ディレクトリに小さな CSV を作る(例: id,title,body のヘッダと 1 行)。※シェルの printf 等で作成可能だが、ここではコマンドは載せません。
  2. メモリ上の SQLite を開き、Create_schema を実行する。
  3. EmbeddingClient インスタンスの embed_batch をモックして、既知のベクトル(例: [0.1]*128)を返すようにする。
  4. embed_csv を呼び出す。呼び出し後、モックの embed_batch が想定されるテキスト(例: ‘Hello\nThis is a test’ のように title と body を結合した文字列)で呼ばれたことを assert する。
  5. SQLite の embeddings テーブルを読み、保存された行の row_id、fragment_id、model が期待値と一致すること、vector BLOB のバイト長が float32 の次元数×4 バイト(例: 128*4)であることを検証する。

主なトラブルシューティング項目(具体的に確認する点)

  • CSV に id 列は存在するか。欠落すると row_to_fragment で失敗する。
  • embed_csv 実行後、バッチごとのコミットが行われているか(途中失敗時の部分保存の有無を設計に合わせる)。
  • 埋め込みクライアントが返すベクトル数がバッチ内の断片数と一致するか。ミスマッチは ValueError 相当で検出する。
  • vector を float32 で保存しているか。テストでは BLOB 長=次元×4 バイトを使って検証するのが簡便。
  • load_embeddings の戻りベクトル配列と records 配列が同じ順序か。ORDER BY を使って読み出すこと、np.frombuffer の場合は copy を取ることを確認する。
  • 検索時に次元不一致がないか。クエリ埋め込みの次元と保存済みベクトルの次元は一致させる必要がある。

Faiss は任意の次段階

小規模であれば SQLite に読み出したベクトルを NumPy で線形探索するだけで十分です。高速化したい場合は Faiss を導入して IndexFlatL2(あるいは他のインデックス)を作り、NumPy 配列を一括で追加して検索する。Faiss は環境依存の導入手順が必要なため、ここでは任意の次段階として簡潔に扱います。

実装時の注意(まとめ)

  • 本文ではプレーンな設計とテスト方針を示しましたが、コード本文は掲載していません。実装例をそのままコピペして動くと保証するものではありません。
  • 本番移行時には埋め込み API 固有のエラー処理、レート制御、認証管理、監視、テレメトリを必ず追加してください。
  • 小規模検証→段階的スケールアップの手順を作り、各段階で性能・一貫性・復旧を確認してください。

以上が「CSVを逐次読み込み→断片化→埋め込み(モック)→SQLiteへ保存→小規模検索へ渡す」ための最小構成の説明です。実装時には上に挙げた関数群を1つのモジュールにまとめ、pytest による保存結果の検証(モック差し替えと BLOB 長の確認)を行うと、配線確認と基本品質の担保に有効です。

第147回 実務で使えるPython基礎:collectionsとitertoolsで書く効率的な表データ処理パターン

日々のCSVや表データ処理で「遅い」「メモリを食う」「コードが読みにくい」と感じたことはありませんか。この記事ではcollectionsとitertoolsを用いて、現場でよく出るパターン(グルーピング・集計・チャンク処理・スライディングウィンドウ・重複除去)を、貼り付けてすぐ使える関数例とともに整理します。第135回(CSVクリーニング)や第146回(スキーマ検証)の後処理として自然に使える手順を意識しています。

目的と対象読者

小規模〜中規模のCSV/表データを日常的に処理する実務担当者向け。読み終わったらそのまま業務に貼れるコード例を重視し、可読性とメモリ効率の両立を目指します。

導入:なぜcollectionsとitertoolsを使うか

標準ライブラリのcollectionsとitertoolsは追加依存が不要で、次の利点があります。

  • 可読性:意図がはっきりするデータ構造(Counter, defaultdict, deque)
  • 速度:C実装されたイテレータや最適化パターン
  • メモリ効率:ジェネレータやストリーミング処理で大きな中間リストを避けられる

実務フローでは第135回のCSVクリーニング→第146回のスキーマ検証→本記事の集計やバッチ処理、という流れでつなぐと作業が安定します。

主要ツールの早見表

モジュール/クラス 用途 注意点
collections.defaultdict 存在しないキーに対するデフォルト値の自動生成(集計が簡潔) defaultdict(list)では欠損キーごとに新しいリストが生成される。共有済みオブジェクトを返すファクトリを使う場合は共有に注意
collections.Counter 要素の頻度集計・最頻値取得 大きなユニーク数はメモリ増加の要因
collections.deque 固定長のスライディングウィンドウ、両端操作 maxlenを設定してメモリ上限を作ると安全
collections.namedtuple 軽量なレコード型(可読性向上) 不変なので安全に使える
itertools.groupby 連続する同キーのグルーピング(ソート済み前提) 事前にソートしないと期待したグループにならない
itertools.islice / tee / chain.from_iterable チャンク化、複製ストリーム、フラット化 teeはメモリを使う点に注意(複数消費時)
functools.lru_cache 関数結果のキャッシュ(計算コストの高い変換に有用) キャッシュサイズとメモリを設計する

実務パターンと短い解説

以下は現場で繰り返し出るパターンと、注意点・貼り付け可能な関数例です。

1) 安定した順序での重複除去(先頭を残す)

dictはPython3.7以降で挿入順を保ちます。行の順序を保ちながら重複を除去したいときに有用です。

def dedupe_preserve_first(rows, key_func):
    """rows: iterable of rows
    key_func: row -> key used for dedupe
    yields unique rows, keeping the first occurrence"""
    seen = {}
    for r in rows:
        k = key_func(r)
        if k not in seen:
            seen[k] = True
            yield r

注意点:ユニークキーの総数が極端に多いとseenがメモリを占有します。

2) ソート済みデータのgroupbyによる集計

groupbyは連続する同キーをまとめるため、事前にソートが必要です。

from itertools import groupby
from operator import itemgetter

def aggregate_sorted(rows, key_index, value_index, agg_func=sum):
    """rows: iterable of sequences already sorted by key_index"""
    for key, group in groupby(rows, key=itemgetter(key_index)):
        vals = (float(r[value_index]) for r in group)
        yield key, agg_func(vals)

注意点:groupbyの前に並べ替えが必要な場合、その処理で入力全体をメモリに保持することがあります。行数だけで判断せず、実データで処理時間と最大メモリ使用量を測り、収まらない場合は外部ソートや別の集計方法を検討します。

3) defaultdict/Counterによる高速集計・マージ

逐次読み込みで累積集計する基本パターン。

from collections import defaultdict, Counter

def accumulate_counts(rows, key_func):
    cnt = Counter()
    for r in rows:
        cnt[key_func(r)] += 1
    return cnt

def merge_counters(*counters):
    total = Counter()
    for c in counters:
        total.update(c)
    return total

注意点:Counterはキー数が増えるとメモリを使うため、必要に応じて中間集計をファイル化する。

4) isliceで作るチャンク処理(APIバッチ送信など)

イテレータをn件ずつ切って処理するテンプレート。nには正の整数を指定します。

from itertools import islice

def chunked(iterable, n):
    if not isinstance(n, int) or isinstance(n, bool) or n <= 0:
        raise ValueError("n must be a positive integer")
    it = iter(iterable)
    while True:
        chunk = list(islice(it, n))
        if not chunk:
            break
        yield chunk

ユースケース:APIのバッチ送信、分割書き出し。チャンクサイズはAPI制限やメモリと相談して決定。

5) dequeを使ったスライディングウィンドウ

固定長ウィンドウで近傍集計や異常検知をする場合に有用です。nには正の整数を指定します。

from collections import deque

def sliding_window(iterable, n):
    if not isinstance(n, int) or isinstance(n, bool) or n <= 0:
        raise ValueError("n must be a positive integer")
    it = iter(iterable)
    dq = deque(maxlen=n)
    for x in it:
        dq.append(x)
        if len(dq) == n:
            yield tuple(dq)

例:移動平均や最近n件の最大値を計算して閾値判定に使います。

6) chain.from_iterableでのフラット化/多段処理

入れ子になったイテラブルをフラット化して連続処理する時に便利です。

from itertools import chain

def flatten(list_of_lists):
    return chain.from_iterable(list_of_lists)

コード例とコピー可能なユースケース

ここでは「1関数=単一責務」でテストしやすい形にした実務向けサンプルを示します。標準ライブラリのみを使います。

① 行単位でメモリ節約しながら日次集計を作る(csv.reader+itertools+defaultdict)

import csv
from collections import defaultdict
from datetime import datetime

def daily_aggregate_csv(path, date_index, value_index, date_fmt='%Y-%m-%d'):
    """CSVを逐次読みして日次合計を返すジェネレータ"""
    totals = defaultdict(float)
    with open(path, newline='', encoding='utf-8') as f:
        reader = csv.reader(f)
        for row in reader:
            try:
                d = datetime.strptime(row[date_index], date_fmt).date()
                v = float(row[value_index])
            except Exception:
                continue  # 実務ではログを残す
            totals[str(d)] += v
    for day, total in totals.items():
        yield day, total

② API呼び出しをチャンク化してPrefect/スケジューラに渡す最小実装

def send_in_batches(iterable, batch_size, send_func):
    """send_funcはリストを受け取る関数(同期)"""
    for batch in chunked(iterable, batch_size):
        send_func(batch)

③ スライディングウィンドウで移動平均・閾値判定

def moving_average_threshold(iterable, window_size, threshold):
    for window in sliding_window(iterable, window_size):
        avg = sum(window) / window_size
        if avg > threshold:
            yield True, avg
        else:
            yield False, avg

④ 重複除去+最初の有効行を残す例

def dedupe_keep_first_by_key(rows, key_index):
    def key_func(r):
        return r[key_index]
    yield from dedupe_preserve_first(rows, key_func)

性能・メモリの実務チェックリスト

確認項目 測定方法 判断のしかた
入力全体を保持していないか 代表的な実データで最大メモリ使用量を測る 実行環境の利用可能メモリに余裕がなければ、逐次処理や分割処理へ変更する
チャンクサイズ 複数のサイズで処理時間と最大メモリ使用量を比較する APIの件数制限や書き込み先の制約を守り、実測値から決める
groupby前のソート ソートと集計を分けて処理時間・メモリ使用量を測る メモリや許容時間を超える場合は、外部ソートや別の集計方法を検討する
dequeのサイズ 業務上必要な期間・件数を確認する 目的からウィンドウ幅を決め、不要な履歴を保持しない
Counterのキー数 ユニークキー数と最大メモリ使用量を記録する メモリに収まらない場合は、分割集計やデータベースの利用を検討する

測定にはtimeittracemallocなどを利用できます。ただし、固定した行数や割合をすべての環境に当てはめることはできません。本番に近いデータ量、実行環境、許容時間をそろえて比較します。

落とし穴と運用上の注意

  • groupbyの事前ソート忘れ:ソートしないと意図しないグループに分割される。
  • defaultdictのファクトリの罠:defaultdict(list)では欠損キーごとに新しいリストが生成される。一方、同じ可変オブジェクトを返すファクトリを指定すると複数のキーで共有され、副作用が起きることがある。
  • 巨大ファイル処理でのバッファ戦略:読み込み時のバッファサイズやopenのencodingを運用基準として明記する。
  • デバッグ用サンプルデータ作成:代表的なエッジケース(空行、欠損、型エラー)を含めた小サンプルを最低1つ用意する。
  • ログ出力ポイント:中間集計時のキー数や最長処理時間をログに残すと再現可能な失敗ケースを作りやすい。

テスト・CI向けの小節

関数は純粋関数(副作用を少なく)にしてpytestでテストします。例:

import pytest

def test_dedupe_preserve_first():
    rows = [("a",1),("b",2),("a",3)]
    out = list(dedupe_preserve_first(rows, lambda r: r[0]))
    assert out == [("a",1),("b",2)]

def test_chunked_empty():
    assert list(chunked([], 10)) == []

def test_chunked_rejects_non_positive_size():
    with pytest.raises(ValueError):
        list(chunked([1, 2], 0))

def test_sliding_window_rejects_non_positive_size():
    with pytest.raises(ValueError):
        list(sliding_window([1, 2], 0))

API呼び出しはモック(unittest.mock)でsend_funcを代替して呼び出し回数や最後に送ったデータを検証します。

実装テンプレート配布案

標準ライブラリのみのミニテンプレートをGitHub Gistで配布すると実務導入が早まります。リポジトリには:

  • functions.py(今回の関数群)
  • tests/test_functions.py(pytestテスト)
  • examples/(小さなCSVサンプル、README)

依存が少ないため社内に取り込みやすく、CIへの組み込みも容易です。

まとめ

collectionsとitertoolsは標準ライブラリでありながら、実務で直接役立つ強力なツール群です。この記事で示したパターン(重複除去、groupby集計、defaultdict/Counter、チャンク処理、スライディングウィンドウ、フラット化)はそのまま業務に貼って使える形にしています。ポイントは「メモリを意識した設計」「単一責務の小さな関数」「ログとテストを忘れないこと」です。

次の一歩(第148回以降の案)

次回以降では、並列処理(concurrent.futuresやmultiprocessing)とcollections/itertoolsの組み合わせ、またストリーミングETLでのbackpressure対策や外部ソートの実践を扱う予定です。今回の関数群を土台に、並列化やスケーリングを段階的に導入していくことをおすすめします。

Manage AI シリーズ「AIとPythonの実務」 — 本記事は第147回目です。業務に直結する小さな改善を積み重ねていきましょう。