第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回目です。業務に直結する小さな改善を積み重ねていきましょう。

第146回 実務で使える表データのスキーマ検証と自動修復ワークフロー — Pythonで作る現場向けバリデーションの手順

現場で受け取るCSVやExcelがそのままではAIモデルや集計に使えず、作業が止まってしまう──そんな経験は多いはずです。本記事では「受け取り側の契約(スキーマ)」を明確にし、検証→自動修復→記録→隔離という実務ワークフローをPythonで実装します。CSVとExcel(.xlsx)の読み込み、検証失敗時の例外処理、修復ログと隔離ファイルの出力までを一つの最小構成にまとめます。第134回・第135回のファイルI/O/クリーニング記事の知識がある前提です。

狙いと前提

目的は、現場で受け取る表データを「そのまま処理できる状態」に自動で近づけることです。完全自動で正解にすることを約束するのではなく、許容できる自動修復ポリシーを定義して運用する点を重視します。

この記事の実装例では、入力ファイルを変更せず、処理結果を次の3種類に分けて出力します。

  • accepted.csv:検証を通過したレコード
  • rejected.csv:修復または検証に失敗したレコード
  • audit.jsonl:行番号、処理結果、修復前後、適用ルールを記録した監査ログ

ステップ1:スキーマ設計の実務手順

まずは業務要件から逆算して、スキーマ(契約)を設計します。下表は実務で使いやすいテンプレートです。

カラム名 必須 許容値/正規表現 欠損ポリシー ユニーク 外部参照
customer_id string 必須 ^[A-Z0-9\-]{8,}$ 欠損は隔離 顧客マスタ参照
email string 条件付き 業務で合意した形式 空文字を許容するか明示 業務次第 なし
order_date date 必須 出力はISO日付 パース不可は隔離 false なし
amount number 必須 >=0 欠損は隔離、または業務ポリシーを別途定義 false 通貨変換ルール

設計時のチェックリスト:

  • 「必須」か「許容する欠損」かを明確にする
  • 許容されるデータ型とフォーマットを具体的に例示する
  • 自動修復ルールを優先順位付きで決める
  • ユニーク制約や外部参照を行単位で確認するか、バッチ後に確認するかを決める
  • スキーマのバージョン管理と互換性方針を定める

ステップ2:環境構築と入力ファイル

以下の例はPython 3.10以降を想定します。CSVは標準ライブラリのcsv、Excel(.xlsx)はpandasopenpyxlで読み込みます。

mkdir table-validator
cd table-validator
python -m venv .venv

# macOS / Linux
source .venv/bin/activate

# Windows PowerShellの場合
# .venv\Scripts\Activate.ps1

python -m pip install --upgrade pip
python -m pip install pandas openpyxl jsonschema "pydantic>=2" pandera pytest

動作確認用のinput.csvを作成します。

customer_id,email,order_date,amount
ABCD-1234, MAILTO:USER@EXAMPLE.COM ,31/01/2025,"1,200"
BAD,user@example.com,not-a-date,-10
EFGH-5678,,2025-02-01,300

1行目は修復後に検証を通過し、2行目は日付のパースに失敗して隔離されます。3行目はそのまま、または軽微な正規化後に通過します。

ステップ3:dataclassからJSON Schemaを生成する軽量実装

次のコードをworkflow.pyとして保存してください。dataclassの型情報とメタデータからJSON Schemaの辞書を生成し、jsonschema.validateで行単位に検証します。

from __future__ import annotations

import argparse
import csv
import json
import unicodedata
from dataclasses import MISSING, asdict, dataclass, field, fields
from datetime import datetime
from pathlib import Path
from typing import Any, get_type_hints

import pandas as pd
from jsonschema import FormatChecker, ValidationError, validate


INPUT_FIELDS = ["customer_id", "email", "order_date", "amount"]


@dataclass
class OrderRow:
    customer_id: str = field(
        metadata={"pattern": r"^[A-Z0-9\-]{8,}$"}
    )
    order_date: str = field(
        metadata={"format": "date"}
    )
    amount: float = field(
        metadata={"minimum": 0}
    )
    email: str = ""


def dataclass_to_json_schema(model_class: type) -> dict[str, Any]:
    """この例で使うstrとfloatをJSON Schemaへ変換する。"""
    type_map = {
        str: "string",
        int: "integer",
        float: "number",
        bool: "boolean",
    }
    hints = get_type_hints(model_class)
    properties: dict[str, Any] = {}
    required: list[str] = []

    for item in fields(model_class):
        python_type = hints[item.name]
        if python_type not in type_map:
            raise TypeError(f"未対応の型です: {item.name}={python_type}")

        property_schema = {"type": type_map[python_type]}
        for key in ("pattern", "format", "minimum", "maximum"):
            if key in item.metadata:
                property_schema[key] = item.metadata[key]

        properties[item.name] = property_schema

        if item.default is MISSING and item.default_factory is MISSING:
            required.append(item.name)

    return {
        "$schema": "https://json-schema.org/draft/2020-12/schema",
        "type": "object",
        "properties": properties,
        "required": required,
        "additionalProperties": False,
    }


ORDER_SCHEMA = dataclass_to_json_schema(OrderRow)


def normalize_text(value: Any) -> str:
    text = "" if value is None else str(value)
    return unicodedata.normalize("NFKC", text).strip()


def repair_row(row: dict[str, Any]) -> tuple[dict[str, Any], list[dict[str, Any]]]:
    """許可した修復だけを適用し、修復内容も返す。"""
    repaired = dict(row)
    changes: list[dict[str, Any]] = []

    def set_value(field_name: str, new_value: Any, rule: str) -> None:
        old_value = repaired.get(field_name)
        repaired[field_name] = new_value
        if old_value != new_value:
            changes.append(
                {
                    "field": field_name,
                    "rule": rule,
                    "before": old_value,
                    "after": new_value,
                }
            )

    customer_id = normalize_text(repaired.get("customer_id"))
    set_value("customer_id", customer_id, "NFKC変換と前後空白削除")

    email = normalize_text(repaired.get("email")).lower()
    if email.startswith("mailto:"):
        email = email.removeprefix("mailto:").strip()
    set_value("email", email, "小文字化、前後空白削除、mailto:削除")

    amount_text = normalize_text(repaired.get("amount")).replace(",", "")
    try:
        amount = float(amount_text)
    except ValueError as exc:
        raise ValueError(f"amountを数値に変換できません: {amount_text!r}") from exc
    set_value("amount", amount, "NFKC変換、カンマ削除、float変換")

    date_text = normalize_text(repaired.get("order_date"))
    parsed_date = None
    selected_format = None
    for date_format in ("%Y-%m-%d", "%d/%m/%Y"):
        try:
            parsed_date = datetime.strptime(date_text, date_format).date()
            selected_format = date_format
            break
        except ValueError:
            continue

    if parsed_date is None:
        raise ValueError(f"order_dateを日付に変換できません: {date_text!r}")

    set_value(
        "order_date",
        parsed_date.isoformat(),
        f"{selected_format}で解析しISO日付へ変換",
    )

    return repaired, changes


def read_rows(path: Path) -> list[dict[str, Any]]:
    """CSVまたは.xlsxを辞書のリストとして読み込む。"""
    suffix = path.suffix.lower()

    if suffix == ".csv":
        with path.open("r", encoding="utf-8-sig", newline="") as handle:
            return [dict(row) for row in csv.DictReader(handle)]

    if suffix == ".xlsx":
        dataframe = pd.read_excel(
            path,
            engine="openpyxl",
            dtype=str,
            keep_default_na=False,
        )
        return dataframe.to_dict(orient="records")

    raise ValueError("対応形式は.csvまたは.xlsxです")


def write_csv(
    path: Path,
    rows: list[dict[str, Any]],
    fieldnames: list[str],
) -> None:
    path.parent.mkdir(parents=True, exist_ok=True)
    with path.open("w", encoding="utf-8-sig", newline="") as handle:
        writer = csv.DictWriter(handle, fieldnames=fieldnames)
        writer.writeheader()
        writer.writerows(rows)


def process_file(input_path: Path, output_dir: Path) -> None:
    accepted: list[dict[str, Any]] = []
    rejected: list[dict[str, Any]] = []
    audit_records: list[dict[str, Any]] = []

    for line_number, original in enumerate(read_rows(input_path), start=2):
        changes: list[dict[str, Any]] = []

        try:
            repaired, changes = repair_row(original)
            validate(
                instance=repaired,
                schema=ORDER_SCHEMA,
                format_checker=FormatChecker(),
            )

            # 検証後のキーと型を固定する。
            normalized = asdict(
                OrderRow(
                    customer_id=repaired["customer_id"],
                    email=repaired.get("email", ""),
                    order_date=repaired["order_date"],
                    amount=repaired["amount"],
                )
            )
            accepted.append(normalized)
            status = "repaired" if changes else "accepted"

            audit_records.append(
                {
                    "line": line_number,
                    "status": status,
                    "original": original,
                    "result": normalized,
                    "changes": changes,
                }
            )

        except (ValueError, ValidationError) as exc:
            rejected_row = dict(original)
            rejected_row["_line"] = line_number
            rejected_row["_error"] = str(exc)
            rejected.append(rejected_row)

            audit_records.append(
                {
                    "line": line_number,
                    "status": "rejected",
                    "original": original,
                    "changes": changes,
                    "error_type": type(exc).__name__,
                    "error": str(exc),
                }
            )

    output_dir.mkdir(parents=True, exist_ok=True)

    write_csv(
        output_dir / "accepted.csv",
        accepted,
        ["customer_id", "order_date", "amount", "email"],
    )

    extra_reject_fields = sorted(
        {
            key
            for row in rejected
            for key in row
            if key not in INPUT_FIELDS and key not in {"_line", "_error"}
        }
    )
    reject_fields = INPUT_FIELDS + extra_reject_fields + ["_line", "_error"]
    write_csv(output_dir / "rejected.csv", rejected, reject_fields)

    with (output_dir / "audit.jsonl").open("w", encoding="utf-8") as handle:
        for record in audit_records:
            handle.write(json.dumps(record, ensure_ascii=False) + "\n")

    print(f"accepted={len(accepted)} rejected={len(rejected)}")
    print(f"output={output_dir.resolve()}")


def main() -> None:
    parser = argparse.ArgumentParser()
    parser.add_argument("input", type=Path)
    parser.add_argument("--output-dir", type=Path, default=Path("output"))
    args = parser.parse_args()
    process_file(args.input, args.output_dir)


if __name__ == "__main__":
    main()

保存後、次の順番で実行します。

python workflow.py input.csv --output-dir output

Excelを入力する場合は、同じ列を持つinput.xlsxを用意して次のように実行します。出力形式はCSVとJSON Linesです。

python workflow.py input.xlsx --output-dir output

repair_row内の変換に失敗した場合はValueError、JSON Schemaに違反した場合はjsonschema.ValidationErrorを捕捉します。いずれも処理全体を停止させず、その行をrejected.csvへ隔離します。入力に想定外の列がある場合も、additionalProperties: falseにより検証エラーになります。

ステップ4:実装パターンの比較

パターン 概要 メリット デメリット
軽量(dataclass + jsonschema) dataclassからJSON Schemaを生成し、行単位で検証する 処理順と修復ルールを明示しやすい 型変換や複雑な修復は手実装になる
Pydantic モデル単位で型変換と検証を行う 詳細な検証エラーを取得しやすい 独自の修復ルールにはvalidatorの実装が必要
Pandera pandas DataFrame全体の列型や制約を検証する ユニーク制約などの表単位の検査を記述しやすい 行単位の修復・隔離は別途設計が必要

Pydanticによる行単位の検証例

Pydantic v2を使う場合は、model_validateで入力辞書を検証し、ValidationErrorを捕捉します。次は日付、金額、メールをモデル内で正規化する最小例です。

import unicodedata
from datetime import date, datetime
from typing import Any

from pydantic import (
    BaseModel,
    ConfigDict,
    Field,
    ValidationError as PydanticValidationError,
    field_validator,
)


class OrderModel(BaseModel):
    model_config = ConfigDict(extra="forbid")

    customer_id: str = Field(pattern=r"^[A-Z0-9\-]{8,}$")
    order_date: date
    amount: float = Field(ge=0)
    email: str = ""

    @field_validator("customer_id", mode="before")
    @classmethod
    def normalize_customer_id(cls, value: Any) -> str:
        return unicodedata.normalize("NFKC", str(value)).strip()

    @field_validator("amount", mode="before")
    @classmethod
    def normalize_amount(cls, value: Any) -> float:
        text = unicodedata.normalize("NFKC", str(value)).strip()
        return float(text.replace(",", ""))

    @field_validator("order_date", mode="before")
    @classmethod
    def normalize_order_date(cls, value: Any) -> date:
        if isinstance(value, date):
            return value

        text = str(value).strip()
        for date_format in ("%Y-%m-%d", "%d/%m/%Y"):
            try:
                return datetime.strptime(text, date_format).date()
            except ValueError:
                continue
        raise ValueError(f"日付に変換できません: {text!r}")

    @field_validator("email", mode="before")
    @classmethod
    def normalize_email(cls, value: Any) -> str:
        text = "" if value is None else str(value)
        text = unicodedata.normalize("NFKC", text).strip().lower()
        return text.removeprefix("mailto:").strip()


raw = {
    "customer_id": "ABCD-1234",
    "email": " MAILTO:USER@EXAMPLE.COM ",
    "order_date": "31/01/2025",
    "amount": "1,200",
}

try:
    order = OrderModel.model_validate(raw)
    print(order.model_dump(mode="json"))
except PydanticValidationError as exc:
    for error in exc.errors():
        print(error)

実運用では、成功したmodel_dump()の結果を正常系へ送り、PydanticValidationErrorを捕捉したレコードを隔離します。修復前の辞書は上書きせず、監査ログへ保存してください。

PanderaによるDataFrame全体の検証例

行単位の修復後に、ユニーク制約や列全体の型を確認する場合はPanderaを利用できます。次の例は、先ほど生成したoutput/accepted.csvを再検証します。

import pandas as pd
import pandera as pa


dataframe = pd.read_csv(
    "output/accepted.csv",
    keep_default_na=False,
)

order_schema = pa.DataFrameSchema(
    {
        "customer_id": pa.Column(
            str,
            checks=pa.Check.str_matches(r"^[A-Z0-9\-]{8,}$"),
            unique=True,
        ),
        "order_date": pa.Column("datetime64[ns]", nullable=False),
        "amount": pa.Column(float, checks=pa.Check.ge(0)),
        "email": pa.Column(str, nullable=False),
    },
    strict=True,
    coerce=True,
)

try:
    validated_dataframe = order_schema.validate(dataframe, lazy=True)
    print(validated_dataframe)
except pa.errors.SchemaErrors as exc:
    print(exc.failure_cases)
    exc.failure_cases.to_csv(
        "output/pandera_failure_cases.csv",
        index=False,
        encoding="utf-8-sig",
    )

lazy=Trueでは、検出した複数の違反をSchemaErrorsとしてまとめて確認できます。failure_casesは違反内容の調査に使えますが、元レコードを隔離する処理は、利用するインデックスや業務上の主キーに合わせて別途実装してください。

ステップ5:自動修復ポリシー

フィールド 修復ルール 優先度 ログ/差分
文字列 NFKC正規化、前後空白の削除 修復前後と適用ルールを保存
amount 全角数字の正規化、カンマ削除、float変換 変換前後を保存
order_date %Y-%m-%d、%d/%m/%Yの順で試し、ISO日付へ統一 採用した日付形式を保存
email 小文字化、前後空白削除、mailto:削除 低〜中 修復前後を保存

金額のクリッピングや欠損値の0埋めは、正しい値を推測する変更になり得ます。そのため上記の実装には含めていません。実施する場合は、業務上の閾値、承認条件、HITL(人の確認)を明文化してください。

  • 修復前の入力ファイルを上書きしない
  • 修復理由を構造化ログとして残す
  • 重大な修復は自動確定せず、担当者の確認対象にする
  • 日付のように解釈が分かれる形式は、送信元と優先順位を合意する

ステップ6:ワークフロー統合と運用フロー

バリデーションはETLの初段に組み込みます。定期実行にはcron、ワークフロー管理ツール、CIなど、運用環境に合う方法を選びます。重要なのは、検証失敗時に処理全体を停止するのか、行単位で隔離して継続するのかを事前に決めることです。

ケース 自動対応 運用アクション
軽微な修復で処理可能 修復後に再検証して正常系へ送る 修復ログを記録
自動修復では確定できない 隔離ファイルへ出力 担当者が原本とエラーを確認
重大な整合性違反 本番投入を止めるか対象行を隔離 送信元へのフィードバックとスキーマ再合意

SLOの例:

  • 受信ファイルの90%を自動修復で本番投入可能にする(週次評価)
  • 重大な隔離レコードは24時間以内に原因調査を開始する

テストとCI

次のコードをtest_workflow.pyとして保存します。

import pytest
from jsonschema import FormatChecker, ValidationError, validate

from workflow import ORDER_SCHEMA, repair_row


def test_repair_and_validate():
    raw = {
        "customer_id": " ABCD-1234 ",
        "email": " MAILTO:USER@EXAMPLE.COM ",
        "order_date": "31/01/2025",
        "amount": "1,200",
    }

    repaired, changes = repair_row(raw)

    validate(
        instance=repaired,
        schema=ORDER_SCHEMA,
        format_checker=FormatChecker(),
    )

    assert repaired["order_date"] == "2025-01-31"
    assert repaired["amount"] == 1200.0
    assert repaired["email"] == "user@example.com"
    assert changes


def test_negative_amount_is_rejected():
    raw = {
        "customer_id": "ABCD-1234",
        "email": "",
        "order_date": "2025-01-31",
        "amount": "-10",
    }

    repaired, _ = repair_row(raw)

    with pytest.raises(ValidationError):
        validate(
            instance=repaired,
            schema=ORDER_SCHEMA,
            format_checker=FormatChecker(),
        )


def test_invalid_date_is_rejected():
    raw = {
        "customer_id": "ABCD-1234",
        "email": "",
        "order_date": "not-a-date",
        "amount": "100",
    }

    with pytest.raises(ValueError):
        repair_row(raw)

テストを実行します。

pytest -q
CI段階 チェック内容
PR時 pytestで検証と自動修復のユニットテストを実行
マージ後 匿名化したサンプルデータで統合テストを実行

スキーマ設計用YAMLテンプレート

次のYAMLは送信元との合意やレビューに使う設計テンプレートです。上記のPythonコードがこのYAMLを直接読み込むわけではありません。

fields:
  - name: customer_id
    type: string
    required: true
    pattern: '^[A-Z0-9\-]{8,}$'
  - name: order_date
    type: date
    required: true
    formats:
      - '%Y-%m-%d'
      - '%d/%m/%Y'
  - name: amount
    type: number
    required: true
    minimum: 0
  - name: email
    type: string
    required: false

運用上の注意点と落とし穴

  • 過度な厳格化は業務停止を招くため、最低限の必須ルールから始める
  • 入力側とのスキーマ合意とサンプル交換を定期的に行う
  • スキーマをバージョン管理し、破壊的変更は事前に通知する
  • ログや監査情報を保存し、どの入力にどの修復を行ったか追跡できるようにする
  • CSVの文字コード、区切り文字、Excelの対象シートも受け渡し契約に含める

次の一歩

  • 主要テーブル1つから導入し、隔離理由の多い項目を確認する
  • ユニーク制約や外部参照をバッチ検証へ追加する
  • 生成したスキーマを使ったデータ生成テストで破壊的変更を検知する

まとめ

  • スキーマは受け取り側と渡し手の「契約」であり、型、必須、形式、修復範囲を合意することが重要です。
  • dataclassとjsonschemaを使うと、行単位の修復、検証、例外処理、隔離を小さく実装できます。
  • Pydanticはモデル単位の型変換と詳細なエラー、PanderaはDataFrame全体の制約確認に利用できます。
  • 自動修復では入力を上書きせず、修復前後と適用ルールを監査ログへ残してください。
  • pytestとCIで正常系、逸脱例、境界値を継続的に確認してください。

第145回 実務で使えるPython基礎:クラスとオブジェクト指向で作る再利用可能なデータ処理コンポーネント

実務でPythonを使っていると、「似た処理が何度も現れる」「一度作った処理のテストや運用が難しい」と感じることはありませんか。この記事では、クラス設計のパターンを、CSVクリーナー、推論ラッパー、Pipelineのコードで示します。標準ライブラリを中心に実装し、pytestによるテストまで確認します。

なぜクラスを使うか(現場メリット)

関数の集まりでも十分なケースは多いですが、状態やリソースの管理、初期化コストの節約、モックしやすいインターフェースを作る点でクラスは有用です。以下の表は代表的な利点と注意点です。

利点 現場で役立つ理由 注意点
状態管理 重いリソースを初回だけロードして再利用できる(モデル、セッション) メモリやマルチプロセスを考慮する必要あり
リソース管理 ファイルやセッションを確実に開放しやすい 誤ったライフサイクル設計でリークする
テスト性 モックしやすいインスタンスメソッドで外部依存を切れる 複雑な状態はテストを書く負担を増やす

状態を持つコンポーネント vs ステートレス関数

どちらを選ぶかは次の観点で判断します。

  • 初期化コストが高い処理(モデルロード、外部接続)には、状態を持つクラスを検討する。
  • 受け取ったデータだけを変換する処理には、ステートレス関数を使う。
  • テストやモックのしやすさを優先する場合は、外部依存をコンストラクターから注入できるようにする。

リソース管理:context managerを使う

ファイルやセッションは明示的に開閉する必要があります。context manager(withステートメント)を使うと、処理中に例外が発生した場合もファイルを閉じられます。

例1:ファイルハンドルを安全に扱うCSVクリーナークラス

以下のコードをcomponents.pyとして保存します。入力CSVの各文字列から前後の空白を取り除き、別のCSVへ保存する例です。

# components.py
import csv
from threading import Lock


class CSVCleaner:
    def __init__(self, input_path, output_path, encoding="utf-8"):
        self.input_path = input_path
        self.output_path = output_path
        self.encoding = encoding
        self._input_file = None
        self._output_file = None

    def __enter__(self):
        self._input_file = open(
            self.input_path,
            "r",
            encoding=self.encoding,
            newline="",
        )
        try:
            self._output_file = open(
                self.output_path,
                "w",
                encoding=self.encoding,
                newline="",
            )
        except Exception:
            self._input_file.close()
            self._input_file = None
            raise
        return self

    def __exit__(self, exc_type, exc, traceback):
        try:
            if self._output_file is not None:
                self._output_file.close()
        finally:
            if self._input_file is not None:
                self._input_file.close()
        return False

    def clean(self):
        if self._input_file is None or self._output_file is None:
            raise RuntimeError("CSVCleaner must be used with a with statement")

        reader = csv.DictReader(self._input_file)
        if reader.fieldnames is None:
            raise ValueError("CSV header is required")

        writer = csv.DictWriter(
            self._output_file,
            fieldnames=reader.fieldnames,
        )
        writer.writeheader()

        count = 0
        for row in reader:
            cleaned_row = {
                key: value.strip() if isinstance(value, str) else value
                for key, value in row.items()
            }
            writer.writerow(cleaned_row)
            count += 1

        return count

次のようにwith構文で使用します。clean()の途中で例外が発生した場合も、__exit__が入出力ファイルを閉じます。

from components import CSVCleaner

with CSVCleaner("raw.csv", "clean.csv") as cleaner:
    processed_count = cleaner.clean()

print(processed_count)

合成(composition)と継承(inheritance)の使い分け

一般の実務では合成を優先します。合成は役割ごとにコンポーネントを分け、再利用とテストを容易にします。継承は明確に「is-a」の関係があるとき、またはフレームワークの拡張時に限定して使います。

選び方 合成(composition) 継承(inheritance)
使う場面 異なる振る舞いを組み合わせるとき 共通の振る舞いを拡張するとき
メリット 変更に強く、テストしやすい コードの重複を減らせるが複雑になりがち

実践例:推論ラッパーとPipeline

次は、モデルを必要になるまでロードしない推論ラッパーと、CSV読み込み、前処理、推論、保存を合成するPipelineクラスです。以下のコードは、先ほど作成したcomponents.pyの末尾へ追記します。

例2:推論ラッパークラス(モデルロード/predict)

class ModelWrapper:
    def __init__(self, loader):
        self.loader = loader
        self._model = None
        self._load_lock = Lock()

    def _get_model(self):
        if self._model is None:
            with self._load_lock:
                if self._model is None:
                    self._model = self.loader()
        return self._model

    def predict(self, records):
        model = self._get_model()
        return model.predict(records)

loaderには、引数なしでモデルを返す関数を渡します。モデルは最初のpredict()呼び出し時にロードされ、その後は同じインスタンスが再利用されます。ロックが保護するのはロード処理であり、各モデルのpredict()自体がスレッドセーフであることを保証するものではありません。

例3:Pipelineクラス(CSV読み込み→前処理→推論→保存)

class Pipeline:
    def __init__(self, model, preprocessor):
        self.model = model
        self.preprocessor = preprocessor

    def run(self, input_path, output_path, encoding="utf-8"):
        with open(input_path, "r", encoding=encoding, newline="") as input_file:
            reader = csv.DictReader(input_file)
            if reader.fieldnames is None:
                raise ValueError("CSV header is required")

            fieldnames = list(reader.fieldnames)
            if "prediction" in fieldnames:
                raise ValueError("Input CSV already has a prediction column")

            original_rows = list(reader)

        processed_rows = [
            self.preprocessor(row) for row in original_rows
        ]
        predictions = list(self.model.predict(processed_rows))

        if len(predictions) != len(original_rows):
            raise ValueError(
                "The number of predictions must match the number of input rows"
            )

        output_fieldnames = fieldnames + ["prediction"]
        with open(output_path, "w", encoding=encoding, newline="") as output_file:
            writer = csv.DictWriter(
                output_file,
                fieldnames=output_fieldnames,
            )
            writer.writeheader()

            for row, prediction in zip(original_rows, predictions):
                output_row = dict(row)
                output_row["prediction"] = prediction
                writer.writerow(output_row)

        return len(original_rows)

Pipelineは、predict(records)を持つオブジェクトと、CSVの1行を前処理する関数を受け取ります。特定の機械学習ライブラリには依存していません。

Pipelineをローカルで実行する

動作確認用として、次の内容をinput.csvへ保存します。

id,amount
1,80
2,120

続いて、以下をrun_pipeline.pyとして保存します。

from components import ModelWrapper, Pipeline


class ThresholdModel:
    def __init__(self, threshold):
        self.threshold = threshold

    def predict(self, records):
        return [
            "high" if record["amount"] >= self.threshold else "low"
            for record in records
        ]


def load_model():
    return ThresholdModel(threshold=100.0)


def preprocess(row):
    return {"amount": float(row["amount"].strip())}


model = ModelWrapper(loader=load_model)
pipeline = Pipeline(model=model, preprocessor=preprocess)
processed_count = pipeline.run("input.csv", "output.csv")
print(f"processed: {processed_count}")

components.pyinput.csvrun_pipeline.pyを同じディレクトリへ置き、次の順序で実行します。

python run_pipeline.py

実行後のoutput.csvは次の内容になります。

id,amount,prediction
1,80,low
2,120,high

CSVが存在しない場合のFileNotFoundError、数値へ変換できない場合のValueError、モデル内部の例外は呼び出し元へ伝わります。ここでは広い例外を握りつぶさず、バッチ処理やAPIなどの呼び出し側でログ記録、再試行、処理中断を選べるようにしています。

テストとモックの方針

外部依存(ファイル、HTTP、モデル)はモックで切り離し、入出力の最小契約をテストします。ここではpytestのtmp_pathと、標準ライブラリのunittest.mock.Mockを使います。

次のコードをtests/test_components.pyとして保存します。

import csv
from unittest.mock import Mock

from components import CSVCleaner, ModelWrapper, Pipeline


def test_csv_cleaner_closes_files_and_strips_values(tmp_path):
    input_path = tmp_path / "raw.csv"
    output_path = tmp_path / "clean.csv"
    input_path.write_text(
        "id,name\n1, Alice \n",
        encoding="utf-8",
    )

    with CSVCleaner(input_path, output_path) as cleaner:
        assert cleaner.clean() == 1

    with output_path.open("r", encoding="utf-8", newline="") as file:
        rows = list(csv.DictReader(file))

    assert rows == [{"id": "1", "name": "Alice"}]


def test_model_wrapper_loads_model_only_once():
    loaded_model = Mock()
    loaded_model.predict.side_effect = [["first"], ["second"]]
    loader = Mock(return_value=loaded_model)
    wrapper = ModelWrapper(loader=loader)

    assert wrapper.predict([{"amount": 1}]) == ["first"]
    assert wrapper.predict([{"amount": 2}]) == ["second"]

    loader.assert_called_once_with()
    assert loaded_model.predict.call_count == 2


def test_pipeline_writes_predictions(tmp_path):
    input_path = tmp_path / "input.csv"
    output_path = tmp_path / "output.csv"
    input_path.write_text(
        "id,amount\n1,120\n",
        encoding="utf-8",
    )

    model = Mock()
    model.predict.return_value = ["high"]

    pipeline = Pipeline(
        model=model,
        preprocessor=lambda row: {"amount": float(row["amount"])},
    )

    assert pipeline.run(input_path, output_path) == 1
    model.predict.assert_called_once_with([{"amount": 120.0}])

    with output_path.open("r", encoding="utf-8", newline="") as file:
        rows = list(csv.DictReader(file))

    assert rows == [
        {"id": "1", "amount": "120", "prediction": "high"}
    ]

pytestが利用できる環境で、プロジェクトのルートディレクトリから次のコマンドを実行します。

python -m pytest

テストのポイントは次の通りです。

  • ModelWrapperのloaderはMockへ差し替えられ、ロード回数を検査できる。
  • Pipelineに渡すモデルもMockにできるため、実際のモデルをロードせず入出力を確認できる。
  • ファイル入出力にはtmp_pathを使い、テスト終了後に残るファイルを減らす。

シリアライズ・設定・バージョン対応(dataclass併用例)

設定をdataclassで型付きにし、JSONで保存する例です。長期運用では、設定データのバージョンを明示し、読み込めないバージョンを検出できるようにします。

項目 実務での推奨
設定管理 dataclassとJSONを使い、設定バージョンを明示する。環境変数で上書きする場合は優先順位も文書化する
シリアライズ 長期保存するデータでは、Python固有の形式だけを前提にせず、JSONなどでスキーマを管理する
モデルファイル ファイル名にバージョンを含め、manifestでメタ情報を管理する

次のコードは、設定の保存と読み込みを行う最小例です。

# config.py
import json
from dataclasses import asdict, dataclass
from pathlib import Path


@dataclass(frozen=True)
class AppConfig:
    model_version: str
    threshold: float
    schema_version: int = 1


def save_config(config, path):
    Path(path).write_text(
        json.dumps(asdict(config), ensure_ascii=False, indent=2),
        encoding="utf-8",
    )


def load_config(path):
    data = json.loads(Path(path).read_text(encoding="utf-8"))

    if data.get("schema_version") != 1:
        raise ValueError("Unsupported configuration schema version")

    return AppConfig(**data)


config = AppConfig(model_version="v1", threshold=100.0)
save_config(config, "config.json")
loaded_config = load_config("config.json")
print(loaded_config)

load_config()は未対応のschema_versionを検出するとValueErrorを送出します。実運用で環境変数による上書きを追加する場合は、JSONを読み込んだ後に適用し、その優先順位を明示してください。

運用チェックリストとデプロイ時の注意点

実務で詰まりやすい点をチェックリスト形式でまとめます。

カテゴリ チェック項目
初期化コスト 重いモデルは遅延ロードやwarmupを行っているか
メモリ/プロセス マルチプロセスやコンテナでモデルを使うときのメモリ増加を確認したか
シリアライズ 特定バージョンのPythonやライブラリだけで読める形式を長期互換の前提にしていないか
設定管理 環境変数とファイルの優先順位を文書化したか
ログ・メトリクス 処理時間、エラー率、入力サイズを計測しているか
障害時 再起動手順、再実行時の重複対策、ロールバック手順があるか

まとめと次の一歩

この記事の要点は次の通りです。

  • 状態を持つ処理(モデル、セッション)はクラスにして、初回ロードやライフサイクルを管理する。
  • 合成を優先してコンポーネントを小さく保ち、テストとモックで外部依存を切り離す。
  • context managerでファイルや接続を安全に閉じ、設定はdataclassとJSONで明示的に管理する。
  • 例外を無条件に握りつぶさず、ログや再試行を担当する呼び出し側へ伝える。

読了後のアクション提案:

  1. components.pyrun_pipeline.pyを保存し、サンプルPipelineをローカルで実行する。
  2. tests/test_components.pyを追加し、python -m pytestでテストする。
  3. 業務の小さな処理を1つ選び、外部依存を注入できるクラスへ分割する。
  4. デプロイ前に、ログ、設定の優先順位、再実行手順、バージョン互換性を確認する。

Manage AI(https://manageai.online)では、次回以降で運用・点検軸をさらに掘り下げます。まずは再利用可能な小さなコンポーネントを作り、実装からテストまでの流れを確認してみてください。

第144回 実務で使えるPython基礎:リストと辞書で作る表データ処理の基本パターン

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

CSVやJSONLを扱うと、欠損や型、行の重複、結合ルールといった「地味に面倒な問題」に時間を取られがちです。本記事では、現場で最も出番の多い「リスト(行の列挙)と辞書(行を表す)」を使った実務的な処理パターンを、すぐ使える関数設計と運用チェックリスト付きで整理します。134/135回のファイル入出力・クリーニング回とつなげて読めるよう設計しています。

基本方針(入出力分離・小さな再利用可能コンポーネント)

関数はできるだけ「入力(rows)を受け取って新しい出力を返す」スタイルにします。ミュータブルな引数を書き換えるとバグになりやすいため、状態変更は明示的に行い、副作用は最小限にします。型ヒントと短いdocstringをつけると再利用性が上がります。

よく使う基本パターン(一覧と例)

パターン 関数署名(例) 説明・使い方(例スニペット)
フィルタ def filter_rows(rows: list, pred) -> list rows = [r for r in rows if pred(r)] 例: pred = lambda r: int(r[‘qty’]) > 0
変換(map) def transform_rows(rows: list, fn) -> list rows = [fn(r.copy()) for r in rows] 例: fnで日付正規化や型変換
グルーピング def group_by(rows: list, key_fn) -> dict collections.defaultdict(list)で集約。agg関数で合計や平均を作る
マージ(アップサート) def merge_rows(base, updates, key: str, prefer=’updates’) -> list 主キーでインデックス化し優先度ルールで列の上書きと衝突ログを出す
チャンク処理(大規模データ) def chunked_iter(it, size: int) yieldでメモリを節約。itertools.isliceで実装

集約とグルーピングの使い分け

手法 適用例 メモリ ソート要否
collections.defaultdict(list) 任意のキーでまとめて集約・集計 キー数×行参照分のメモリ 不要
collections.Counter 出現頻度や簡単な頻度集計 キー数分のメモリ 不要
itertools.groupby ソート済みデータの連続ブロック処理(ストリーミング向け) ほぼ定数(ソート前提) 必要(事前に同キーでソート)

軽い比較ベンチマークのポイント: defaultdictはランダム順のまま使える反面、キー数が多いとメモリが増える。groupbyは事前ソートが必要でソートコストが主要なボトルネックになります。

マージ(アップサート)パターンの実務設計

運用的に重要なのは「冪等性」と「衝突の可視化」です。典型的な実装手順は次のとおりです。

  • baseをkey→rowの辞書に変換する(index化)
  • updatesを走査して既存行があれば優先度ルールで列を更新、なければ挿入
  • 変更点はログ(衝突レコード)として出力する
  • 最後に辞書をリスト化して返す(順序が必要なら元順序を保持)
ポイント 注意点
優先度設計 列ごとに“どちらを優先するか”を明示するルール表を運用ドキュメント化する
冪等性 同じ更新を何度適用しても結果が変わらないように設計する(タイムスタンプ比較など)
衝突ログ 変更前後の値をCSV/JSONで出力して監査できるようにする

大規模データ対策(チャンクとストリーミング)

技法 利点 実務向け注意点
ジェネレータ + itertools.islice メモリ使用量を一定に抑えられる ステートフル処理(グルーピングなど)はチャンク境界の扱いに注意
外部ソート(重複除去) メモリに乗らないキーでのユニーク化が可能 ディスクI/Oのコストを見積もること
ハッシュリング(簡易) キーをハッシュで分割し、分割ごとに処理・マージする 分割数設計とファイル管理が運用コストになる

小さな再利用コンポーネント設計の例

関数 署名 戻り値 docstringのポイント
filter_rows filter_rows(rows: list, pred: Callable[[dict], bool]) -> list フィルタ後の行リスト predの期待するキーと型、欠損時の動作を明記
merge_rows merge_rows(base, updates, key: str, prefer: str=’updates’) -> list マージ後の行リストとオプションで衝突ログ 優先ルールと冪等性の挙動を説明

pytestでのテスト例(境界ケース・欠損値)

テスト名 簡単な内容(擬似アサーション)
test_filter_missing_key assert filter_rows([{‘a’:’1′},{ }], lambda r: r.get(‘a’)==’1′) == [{‘a’:’1′}]
test_merge_prefer_updates assert merge_rows([{‘id’:’1′,’v’:’old’}],[{‘id’:’1′,’v’:’new’}],’id’,prefer=’updates’)[0][‘v’]==’new’
test_chunked_iter_boundary assert len(list(chunked_iter(range(5), 2))) == 3

実務運用チェックリスト

チェック項目 目的・方法
エンコーディング確認 ファイル読み込み時にencodingを明示(UTF-8/BOMなど)
欠損値ルール 欠損の扱い(除外/補完/デフォルト値)を自動チェックで実行
型変換の検証 サンプル変換後の型チェックと失敗行の報告
行数差分とサンプルハッシュ 入出力の行数、先頭N行のハッシュで大きなズレを検出
監査ログ出力 マージや上書きの前後をログ化(143回の監査ログ方針に準拠)

落とし穴と改善策(実務でよくある問題)

  • ミュータブルな行をそのまま編集して意図せず元データを書き換える: コピーして操作する。
  • キーの正規化ミス(空白・大文字小文字): 正規化ルールを統一して統合テストに入れる。
  • パフォーマンスの盲点: プロファイラでボトルネックを特定し、必要ならpandas導入を判断(全体の読みやすさと運用コストを比較)。

次の一歩(ワークフロー化への橋渡し)

作成したコンポーネント群をPrefectやcronに組み込み、リトライ、監査ログ、テストを回すことで運用に耐えるフローになります。まずは1つの取り込み->正規化->マージの小さなパイプラインを作り、SLO(処理時間・成功率)を設定してください。

まとめ

リストと辞書は、表データ処理の基本であり続けます。重要なのは「小さな、再利用可能で副作用の少ない関数」を設計すること、そして運用チェックリストで品質を保つことです。本記事のパターンをテンプレートにして、まずは小さなパイプライン一つを作ってみてください。次回以降はワークフロー化と監査自動化の具体例へつなげます。

第143回 実務で使える運用と点検:監査ログ・SLO・運用チェックリストをPythonで自動化する手順

運用を始めると「何を記録すれば良いか」「どこまで自動化できるか」で迷うことが多いです。本記事は、監査ログ(予測の起点・入力・モデルバージョン・出力)を確実に残し、SLO/SLIを定義して日次・週次チェックをPythonで自動化する実務手順をまとめます。現場でそのまま試せるスキーマ例、チェックリスト、簡易ランブック、サンプルコードと運用上の注意点を提供します。

目次

  • 概要と要件定義
  • 監査ログ設計:JSONLスキーマ例と必須フィールド
  • SLO/SLIの定義と閾値設計
  • Pythonで作る日次/週次自動チェック(サンプルスクリプト)
  • ランブックと対応フロー
  • CI/スケジューラへの組み込みとレポート配信
  • 運用でよくある落とし穴と改善ループ
  • まとめと次につなげる企画案

① 概要と要件定義

まず最初に決めるのは「どのイベントを残すか」と「保持期間・マスキング方針」です。実務では記録不足で原因追跡ができないことがよくあるため、後から必要になりそうな情報は最小限でも残す方針が安全です。以下のポイントを決めてから実装を始めてください。

  • 記録するイベント:リクエスト受信、前処理、モデル呼び出し、応答、エラー、モデルデプロイ/ロールバック
  • 保持期間:法規や契約に従う。一般的には30日~1年を検討(個人情報含む場合は短縮)
  • マスキング:個人情報(PII)は保存しないかハッシュ化/部分マスクする
  • アクセス制御:監査ログは読み取り権限を限定
  • 改ざん防止:ハッシュや署名で改ざん検知できる仕組みを検討

② 監査ログ設計:JSONLスキーマ例と必須フィールド

実務では行指向(JSONL)が扱いやすく、ストリーム処理や検索にも適しています。まずは必須フィールドを定義し、運用に必要なメタデータを含めます。

JSONLスキーマ(例)

{"timestamp": "2026-08-31T09:12:34.123Z", "request_id": "req_0001", "user_id_hash": "sha256:...", "model": "predictor", "model_version": "v1.2.3", "input_summary": "token_count=45, features_ok=true", "input_hash": "sha256:...", "prediction": "approve", "confidence": 0.87, "latency_ms": 210, "error": null, "provenance": {"pipeline_version": "p1.0"}, "signature": "hmac:..."}

上記の1行が1レコード(JSONL)です。運用ではgzipで圧縮して保管することもあります。

フィールド 説明 必須
timestamp UTCタイムスタンプ(ISO8601) 必須
request_id リクエストを一意に識別するID 必須
user_id_hash 個人を特定しないハッシュ/マスク済みID 必須(PIIは非保存が原則)
model / model_version 使用したモデルとバージョン 必須
input_hash / input_summary 入力のハッシュと簡易要約(全文保存しない場合) 必須
prediction / confidence 出力と確信度 必須
latency_ms 推論所要時間(ミリ秒) 推奨
error エラー情報(発生時) 推奨
provenance 処理パイプラインや前処理のバージョン 推奨
signature ログ改ざん検知用の署名(省略可) オプション

③ SLO/SLIの定義と閾値設計

SLO(Service Level Objective)は運用目標、SLI(Service Level Indicator)はそれを測る指標です。現実的な閾値を決めるためには、基線となるメトリクスを数週間観察することが重要です。

指標(SLI) 式(例) SLO目標例
応答時間(P95 latency) 95パーセンタイルの遅延(ms) P95 < 500ms(可用性:99%)
エラー率 エラー数 / 総リクエスト数 < 1%(月間)
データ品質(欠損率) 欠損あるいは不整合レコードの割合 < 0.5%(週次)

SLOは技術的に達成可能で、かつビジネスに意味のある水準で設定します。最初は緩めに設定して運用で引き締めるのが現実的です。

SLOを計算するサンプル関数(Python)

def calc_error_rate(records):
    total = len(records)
    errors = sum(1 for r in records if r.get("error"))
    return errors / total if total else 0.0

def calc_p95_latency(latencies):
    if not latencies:
        return None
    latencies = sorted(latencies)
    import math
    idx = math.ceil(0.95 * len(latencies)) - 1
    return latencies[idx]

④ Pythonで作る日次/週次自動チェック(サンプルスクリプト付き)

ここでは、JSONLログを読み込み、SLOを計算してCSVと簡易HTMLレポートを出力し、Slackへ送る軽量な流れを示します。標準ライブラリ中心で書き、必要に応じてpandas導入を検討します。

ポイントで使う標準要素

  • logging と RotatingFileHandler(ログローテーション)
  • json / csv モジュール
  • datetime, timezone(UTC管理)
  • pathlib(ファイル操作)
  • 例外処理と簡易リトライ

サンプル:日次チェックの簡易スクリプト

import logging
from logging.handlers import RotatingFileHandler
import json
from pathlib import Path
from datetime import datetime, timezone
import csv
import requests

# ログ設定
logger = logging.getLogger("audit_checker")
handler = RotatingFileHandler("logs/audit_checker.log", maxBytes=5*1024*1024, backupCount=3)
formatter = logging.Formatter("%(asctime)s %(levelname)s %(message)s")
handler.setFormatter(formatter)
logger.addHandler(handler)
logger.setLevel(logging.INFO)

DATA_DIR = Path("/var/logs/audit")
REPORT_DIR = Path("/var/reports/audit")
REPORT_DIR.mkdir(parents=True, exist_ok=True)

# JSONLを読み込むユーティリティ
def read_jsonl(path):
    with path.open("r", encoding="utf-8") as f:
        for line in f:
            try:
                yield json.loads(line)
            except json.JSONDecodeError:
                logger.warning(f"invalid json: {path} : {line[:100]}")

# 集計処理
def aggregate_for_period(paths):
    records = []
    latencies = []
    for p in paths:
        for r in read_jsonl(p):
            records.append(r)
            if r.get("latency_ms") is not None:
                latencies.append(r["latency_ms"])
    return records, latencies

# SLO計算(先ほどの関数を活用)
def calc_error_rate(records):
    total = len(records)
    errors = sum(1 for r in records if r.get("error"))
    return errors / total if total else 0.0

def calc_p95_latency(latencies):
    if not latencies:
        return None
    latencies = sorted(latencies)
    import math
    idx = math.ceil(0.95 * len(latencies)) - 1
    return latencies[idx]

# レポート出力
def write_csv_report(path, summary):
    with path.open("w", newline="", encoding="utf-8") as f:
        writer = csv.writer(f)
        writer.writerow(["metric", "value"])
        for k, v in summary.items():
            writer.writerow([k, v])

# Slackへ送る(軽量な例)
def post_slack(webhook_url, text):
    try:
        requests.post(webhook_url, json={"text": text}, timeout=5)
    except Exception as e:
        logger.warning(f"slack post failed: {e}")

# メイン処理(例:前日分)
def run_daily_check(date_str, slack_webhook=None):
    # 日付文字列 -> ファイルを選ぶルールは運用で決める
    day_file = DATA_DIR / f"audit_{date_str}.jsonl"
    records, latencies = aggregate_for_period([day_file])
    error_rate = calc_error_rate(records)
    p95 = calc_p95_latency(latencies)
    summary = {"date": date_str, "total_requests": len(records), "error_rate": error_rate, "p95_latency_ms": p95}
    csv_path = REPORT_DIR / f"summary_{date_str}.csv"
    write_csv_report(csv_path, summary)
    logger.info(f"daily check done: {summary}")
    if slack_webhook:
        post_slack(slack_webhook, f"Daily audit report {date_str}: {summary}")

# 実行例
if __name__ == '__main__':
    today = datetime.now(timezone.utc).strftime('%Y-%m-%d')
    run_daily_check(today, slack_webhook=None)

上記は最小限の例です。実運用ではログのローテーション、古いログの削除、圧縮、取得失敗時のリトライなどを追加してください。

⑤ ランブックと対応フロー(トリアージ、エスカレーション、復旧手順)

平時から復旧までの手順を文章化しておくと、担当が不在でも対応が回ります。以下は簡易ランブックの例です。

アラート種別 優先度 初動対応(担当) エスカレーション先
SLO違反(エラー率急上昇) ログ確認、直近のデプロイ確認(オンコール) 開発リード/インフラ担当
モデルパフォーマンス劣化(急低下) サンプルの入力と出力を抽出して比較(運用担当) MLエンジニア
ログ書き込み失敗 ディスク容量・権限を確認(運用担当) インフラ担当

トリアージのポイント

  • まずは影響範囲(ユーザー数、バッチ処理かリアルタイムか)を評価する
  • 短時間で復旧可能かを判断し、ロールバックや一時遮断を検討する
  • 根本原因調査は別タスクにして、まずサービス維持を優先する

⑥ CI/スケジューラへの組み込みとレポート配信

定期チェックはcronやGCP Cloud Scheduler、GitHub Actionsなどに組み込みます。成果物はCSV/HTMLレポートとして保存し、Slackやメールで配信します。レポートは短く要点を示す形式が受け入れられやすいです。

軽量なSlack送信例(requests使用)は前述の post_slack を参照してください。メール送信はSMTPライブラリ、あるいは外部サービスのAPI利用を検討します。

⑦ 運用でよくある落とし穴と改善ループ

  • 大量ログのコスト増:すべて保存せずサンプリングや条件付き保存(エラー時のみ詳細を保存)を設計する
  • プライバシー:PIIは生データで保存しない。ハッシュ化や部分マスキングの具体例をポリシー化する
  • 改ざん対策:ログに単純なHMACを付与して改ざん検知の第一歩とする(秘密鍵は別管理)
  • SLOの現実的妥協:初期は広めに設定し、運用で引き締める。過度なSLOは誤検知や過剰対応を招く
  • 担当1人でも回せる仕組み:手順を簡潔に、復旧手順を優先度付きで定義する

運用チェックリスト(そのまま使える表)

項目 チェック内容 推奨頻度 責任者 備考
ログ保持 保存ポリシーが適用されているか(古いログは削除/アーカイブ) 月次 運用担当 法規遵守を確認
アクセス制御 ログストアのアクセス権を確認 月次 セキュリティ担当 最小権限の原則
マスキング PIIがマスク/ハッシュされているか 週次 運用担当 抜き打ち検査を推奨
異常検知閾値 SLO/SLIの閾値が現状に合っているか 四半期 SRE/運用 基線を見て調整
通知先 Slack/メールの通知先が最新か 週次 運用担当 オンコール表と同期
抜き打ちテスト テストリクエストでログ整合性を確認 月次 運用担当 結果は記録して改善へ

まとめ

監査ログ設計、SLO/SLI定義、日次/週次の自動チェック、そしてランブック化は運用の基本です。本記事では現場で使えるJSONLスキーマ例、SLO計算、Pythonを使った自動チェックの骨子、実務的な注意点を示しました。まずは最小限のログを確実に取り、SLOを設定して自動化することを優先してください。運用は作って終わりではなく、定期的な見直し(改善ループ)が重要です。

次につなげる企画案

  • ポストモーテムと改善ループの作り方(テンプレート付き)
  • 監査ログを検索・分析する軽量な検索基盤構築(ELK代替)
  • 監査ログからの自動改善サイクル(フィードバックループ)実装例

この記事は「AIとPythonの実務」シリーズの一部です。実務での導入や具体的なスクリプト改良の相談があれば、次回以降でより深掘りします。

第142回 実務で使えるPython基礎:設定とシークレット管理で安全にAPIキーを運用する手順

設定やシークレット管理で困っていませんか。ローカルでは動くのに本番で失敗した、誰かのAPIキーが誤って公開リポジトリに混入した、などの経験は多くの実務担当者にとって身近な問題です。本記事は「実務で確実に使える」手順とチェックリストを、具体例(.env/pydantic、クラウドシークレット、CIルール、ローテーション)を交えて提供します。

なぜ設定とシークレット管理が重要か(現実的リスク)

AIワークフローではAPIキーやモデル用トークンが多数発生します。これらが漏れると、予期せぬコスト発生、不正利用、顧客データ流出などの重大インシデントに直結します。よくある現実的リスク:

  • 誤ってコミットされた.envファイル(公開リポジトリに流出)
  • テスト環境のキーが本番で使用されていたため権限過剰になっていた
  • ローテーションが未実施で長期間同一キーが使われていた

基本原則と設計方針

  • 環境の分離(local / staging / production)を明示する
  • 最小権限:キーには必要最小限のスコープだけ付与する
  • 構成(config)とシークレット(secret)を明確に分ける
  • 運用チェックリストを作り、CIで自動検査する

要件化の簡単な例

リスク 要件
公開リポジトリにキーが入る git-precommitで秘密値検査、CIで再スキャン
テスト用と本番用の混同 環境ごとに設定を分離し、環境変数で明示
漏洩検出後の対処が遅い インシデント対応フローとローテーション手順の整備

ローカル開発の典型パターン(.env + python-dotenv + pydantic)

ローカルではシークレットを直接ファイルで扱いがちですが、.envは開発専用に限定し、チームでの共有は避けます。安全な基本パターン:.env(ローカル・例外的)+環境変数(CI/本番)+型検証(pydantic)です。

サンプル:.env と python-dotenv の読み込み(最小例)

開発時のみ .env を使い、常に環境変数で上書きできるようにします。

# .env (例、絶対にコミットしない)
OPENAI_API_KEY=sk-xxxxx
ENV=local

pydantic の BaseSettings を使った型・バリデーションの例:

from pydantic import BaseSettings, Field

class Settings(BaseSettings):
    openai_api_key: str = Field(..., env='OPENAI_API_KEY')
    env: str = Field('production', env='ENV')

    class Config:
        env_file = '.env'
        env_file_encoding = 'utf-8'

settings = Settings()

ポイント:

  • env_file を指定してローカルでの利便性を保ちつつ、本番では環境変数で上書きする運用にする
  • 必須値に対しては Field(… ) で明示し、起動時に不足があれば例外を出す
  • 型・制約を入れることで誤った値の混入を防ぐ

本番での安全な配置パターン(比較表)

本番環境ではファイルベースよりも環境変数やマネージドシークレットが望ましいです。以下は現場でよく選ばれる選択肢の比較です。

方式 長所 短所 推奨規模
環境変数(直接) 導入容易、ランタイムで読み取り簡単 プロセス内でアクセス可能な全員が見える 小〜中規模、単一ホスト
コンテナのシークレットマウント ファイルとしてマウントし管理しやすい シークレットをファイルとして扱うためアクセス管理が必要 中規模、コンテナ運用
Kubernetes Secrets クラスタ管理、RBACや監査と連携可能 Base64エンコードで平文化のまま管理されるケースがあり運用が必要 中〜大規模、K8s利用
クラウド Secret Manager 自動ローテーション・監査・強力なアクセス制御 コストがかかる、導入運用の学習コスト 推奨(中〜大)

推奨パターン:小規模チームはまず「環境変数 + CIでの検査」から始め、中〜大規模は「クラウドSecret Manager + IAM/RBAC + ローテーション自動化」を目指すとよいでしょう。

クラウドシークレット管理の実務例(概要と最小実装)

ここでは AWS Secrets Manager の呼び出しの最小フロー例を示します。設計では「ローカルフォールバック」を用意しておくと便利です(例:本番はシークレットマネージャ、ローカルは環境変数)。

# boto3 を使った最小取得例(実際は IAM ロールでアクセス)
import boto3
import base64
from botocore.exceptions import ClientError

def get_secret(secret_name, region_name='us-east-1'):
    client = boto3.client('secretsmanager', region_name=region_name)
    try:
        resp = client.get_secret_value(SecretId=secret_name)
        if 'SecretString' in resp:
            return resp['SecretString']
        else:
            return base64.b64decode(resp['SecretBinary'])
    except ClientError:
        raise

実務ポイント:

  • 呼び出しは最小限にし、アプリ起動時にキャッシュする(頻繁な API 呼び出しは避ける)
  • IAM ポリシーは最小権限で Secret の読み取りのみを許可する
  • ローテーションは可能なら自動化(Secrets Manager は Lambda 連携で自動ローテーション可)

CI とコード側での漏洩対策・テスト

コード側・CI パイプラインで自動検査を入れることで漏洩リスクを大幅に下げられます。代表的な対策:

  • pre-commit フック(git-secrets など)によるコミット時スキャン
  • CI ビルドでの再スキャン(push 時や PR 時)
  • テスト環境は固定のモックキーを用意し、本物のキーは一切配備しない
  • 証跡(audit)を残す:誰がいつキーを作成・参照したかをログに残す

pre-commit の簡単な設定例(概念):

# .pre-commit-config.yaml の一部(例)
-   repo: https://github.com/awslabs/git-secrets
    rev: v1.3.0
    hooks:
    -   id: git-secrets

インシデント準備とキーのローテーション手順

鍵が漏れた場合の手順を事前に決めておくと対応が早くなります。基本フロー:

  • 検出:自動検出(CI・監査ログ)または手動報告
  • 無効化:該当キーを即時無効化(可能なら読み取り専用から取り消す)
  • ローテーション:新しいキーを発行し、アプリへ段階的に展開
  • 確認:動作確認とアクセスログの確認
  • 対策:なぜ漏れたかを特定しルールを改訂

ロールプレイ用の簡易チェックリスト(自動化可能)

ステップ 実施内容
検出 CI スキャンアラート、監査ログ、セキュリティチーム報告
無効化 管理コンソール或いは API で即時無効化
新キー発行 Secret Manager などで新キーを発行し、必要権限を付与
展開 段階的に環境変数/デプロイで差し替え、ステージングで検証
レビュー 原因究明と再発防止策の実施

現場での導入手順(段階的チェックリスト)

小規模チームが今日から始められる短期〜長期ステップを示します。

期間 作業
短期(今週〜1ヶ月) 1) .gitignore に .env を追加 2) pre-commit で git-secrets 導入 3) pydantic で起動時バリデーション
中期(1〜3ヶ月) 1) CI に漏洩スキャンを追加 2) 本番は環境変数を利用 3) 簡易ローテーション手順を文書化
長期(3〜12ヶ月) 1) Secret Manager の導入 2) 自動ローテーションと監査ログの整備 3) RBAC の見直し

関連記事との接続(シリーズ案内)

本記事は「AIとPythonの実務」シリーズの一部です。関連回:

  • 第126回:外部API連携(設定の読み替えポイント)
  • 第137回:トークン管理(設計の一貫性)
  • 第139回:推論API運用(運用面の補完)

まとめ

設定とシークレット管理は、AIワークフローの信頼性と安全性の基礎です。まずは小さな改善(.env の扱い、pydantic によるバリデーション、pre-commit/CI の導入)から始め、組織の成長に合わせて Secret Manager や自動ローテーションを導入するのが現実的な道筋です。最後に、必ず「検出→無効化→ローテーション→再発防止」のフローを文書化して運用に落とし込んでください。

次回は、実際に AWS Secret Manager を使って Lambda と連携する具体的なハンズオンを予定しています。シリーズを通じて、現場で使える手順を積み上げていきましょう。

第141回 実務で使えるPython基礎:例外設計とエラーハンドリングで作る説明しやすいAIワークフロー

はじめに — うまく説明できないエラーで時間を取られていませんか?

AIを業務に組み込むと、予期せぬエラーが発生しがちです。外部APIの一時障害、入力データの不整合、モデルの推論失敗、ファイルI/Oの問題……どれも現場では頻出です。しかし重要なのは“何が起きたかを説明できて、対応できること”です。本記事では、実務で使える例外設計とハンドリング手順を、コード例・ログ設計・テスト・運用チェックリストとともに示します。読了後には実装に移れるレベルを目指します。

問題提起:現場でよくある失敗パターン

カテゴリ 症状 原因(例)
外部API タイムアウト、502/503 ネットワーク/レート制限/一時的障害
データ不整合 パース失敗、スキーマ違反 入力CSVの欠損・型違い
モデル推論 推論例外、メモリ不足 モデルの入力前提違反、リソース不足
ファイルI/O 読み書き失敗 権限、ディスク容量、同時アクセス

例外設計の基本

実務向けには、標準例外を乱用せず、業務ドメインに応じたカスタム例外を追加することを推奨します。命名規約は「場所+原因+Action」で分かりやすくします(例:CsvParseError、ApiTimeoutError、ModelInferenceError)。
下記は基本の階層例です。

基底 目的
AppError(Exception) アプリケーション全体で共通の基底 ログやアラートの一括ハンドリング用
TransientError(AppError) 再試行で回復が期待できる一時的エラー ApiTimeoutError、ServiceUnavailableError
PermanentError(AppError) 再試行しても意味がない恒久的エラー InvalidInputError、DataSchemaError
ResourceError(AppError) 外部リソースに関わる問題 FileIOError、DatabaseError

振る舞いマップ(例外ごとの実務アクション)

次の表は例外種別と現場で望ましい振る舞いのマッピングです。実運用ではこれをルール化してワークフローに落とし込みます。

例外種別 現場アクション ログ/メトリクス
TransientError(外部APIタイムアウト等) 自動再試行(指数バックオフ)、失敗後はアラート retry_count、last_error_type
PermanentError(入力データ不備等) 処理をスキップしてユーザー通知/エラーレポート bad_record_count、validation_errors
ResourceError(ディスク/DB権限) 即時停止+オンコールへ通知、想定外ならロールバック resource_status、failure_rate
ModelError(推論失敗) 再試行+軽量なフォールバック(簡易モデル)または手動介入 model_failures、latency

実装ガイド

カスタム例外(dataclassベース)

from dataclasses import dataclass

class AppError(Exception):
    """アプリ全体の基底例外"""
    pass

@dataclass
class TransientError(AppError):
    message: str
    retry_after: float = 0.0

@dataclass
class PermanentError(AppError):
    message: str
    field: str = ""

例外にコンテキストを持たせる(traceback保持)

raise from を使い、元の例外を保持しておくと原因追跡が容易になります。

try:
    result = external_api.call(payload)
except TimeoutError as e:
    raise TransientError("API timeout", retry_after=2.0) from e

例外をラップしてドメインに翻訳するパターン

外部ライブラリの例外はそのまま流すのではなく、ドメイン例外へ翻訳します。

def fetch_with_wrap():
    try:
        return third_party.fetch()
    except third_party.HttpError as e:
        # 外部のHTTPエラーを自ドメインのTransientErrorに変換
        raise TransientError("third_party http error", retry_after=5.0) from e

ロギング&可観測性連携

標準loggingで構造化ログを出し、SentryやPrometheusに必要な情報を渡します。ポイントは「何が」「どこで」「どの程度」の情報を一貫して出すことです。

構造化ログの例

import logging
logger = logging.getLogger(__name__)

def handle_error(exc, context):
    logger.error("processing_error", extra={
        "error_type": type(exc).__name__,
        "message": getattr(exc, 'message', str(exc)),
        "context": context,
    })

メトリクスの置き場所

Prometheus等には次のようなメトリクスを送ります:

  • error_count{type=…}
  • retry_count
  • processing_latency_seconds
ログ/メトリクス 用途 閾値の目安
error_rate (5m) 全体の健全性 >1%で要調査
transient_error_rate 外部依存の劣化 急増は外部障害の兆候

テストとCI

例外処理はユニットテストとCIでの回帰チェックが必須です。外部依存はモックで一貫して疑似障害を作るとよいです。

ユニットテストの例(pytest)

def test_api_timeout(monkeypatch):
    def fake_call(_):
        raise TimeoutError("timeout")
    monkeypatch.setattr('external_api.call', fake_call)

    with pytest.raises(TransientError) as excinfo:
        fetch_with_wrap()
    assert excinfo.value.retry_after == 5.0

CIでの検証項目(例)

  • 例外階層のドキュメントが最新か
  • 主要な例外でのユニットテスト通過
  • ログに必要なフィールドが出力されているかのスナップショット

デプロイ/運用チェックリスト

項目 確認内容
例外リスト 主要例外と期待される振る舞い(retry/skip/alert)がドキュメント化されている
アラートテスト オンコールに通知されることをテスト済み
ログ/メトリクス 構造化ログと主要メトリクスが収集されている
ロールバック手順 障害時のロールバック手順が用意されている

小さなハンズオン:CSV一括推論パイプライン

以下は簡易サンプルです。CSVを読み、1行ずつ推論して結果を書き出す。APIタイムアウトは再試行、データ不備はスキップしてレポートします。

import csv
import time
import logging

logger = logging.getLogger(__name__)

class CsvParseError(PermanentError):
    pass

class ApiTimeoutError(TransientError):
    pass

def predict_with_retry(payload, max_retry=3):
    for i in range(max_retry):
        try:
            return external_api.predict(payload)
        except TimeoutError as e:
            if i == max_retry - 1:
                raise ApiTimeoutError("api timeout", retry_after=2.0) from e
            time.sleep(2 ** i)

def process_csv(in_path, out_path):
    bad_rows = []
    with open(in_path) as inf, open(out_path, 'w', newline='') as outf:
        reader = csv.DictReader(inf)
        writer = csv.DictWriter(outf, fieldnames=reader.fieldnames + ['prediction'])
        writer.writeheader()
        for row in reader:
            try:
                if not row.get('text'):
                    raise CsvParseError('missing text', field='text')
                pred = predict_with_retry(row['text'])
                row['prediction'] = pred
                writer.writerow(row)
            except PermanentError as e:
                logger.warning('skip_row', extra={'reason': e.message, 'row': row})
                bad_rows.append(row)
            except TransientError as e:
                logger.error('transient_failure', extra={'reason': e.message})
                raise
    return bad_rows

読後アクション(まず取り組む3つ)

  • 自分のワークフローで発生しうる主要例外を一覧化する。
  • 各例外に対して期待される振る舞い(retry/skip/alert)を書き出す。
  • 1種類の例外(例:APIタイムアウト)を実装してユニットテスト&デプロイする。

まとめ

実務で使えるエラーハンドリングは、例外を分類し「期待される振る舞い」を明文化することが出発点です。カスタム例外の設計、例外のラップ、構造化ログとメトリクス連携、そしてユニットテストとCIでの回帰検証をセットにすることで、何が起きたか説明でき、対応できる運用が作れます。本記事の表やチェックリストを自分のワークフローに当てはめ、まず一つの例外を実装・検証してみてください。

関連:第124回(再試行)、第126回(API連携)、第123回(可観測性)も合わせて参照すると、今回の設計を運用に落とし込みやすくなります。

第140回 実務で使えるPython基礎:標準ライブラリで作る生産性向上ツール

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

標準ライブラリだけで仕事に使えるツールを作ろうとすると、「依存を減らしたい」「でも本当に標準で足りるのか?」と迷うことが多いはずです。本記事ではその迷いに寄り添い、まずは標準ライブラリだけで安全に、保守しやすく、実用的な小さなツールを作る手順とテンプレートを示します。

導入:なぜまず標準ライブラリを使うか

標準ライブラリを優先する理由は主に次の3点です。

  • 依存管理が簡単(追加のパッケージ管理が不要)
  • 保守性が高い(長期的互換性やドキュメントが豊富)
  • セキュリティと信頼性(広くレビューされたコード群)

標準ライブラリで十分かを判断するチェックリスト

項目 標準で十分か? コメント
HTTPクライアント 限定的 簡単なGET/POSTはurllibで可能。認証やセッション管理が複雑ならrequestsの検討を。
CSV/JSON処理 十分 標準のcsv/jsonでほとんどの業務処理はこなせる。
データ解析(大規模) 不十分 大量データや高度な集計はpandasの方が現実的。
非同期IOや高スループット 場合による asyncioは強力だが、使いこなしコストを考慮。

まずは小さなPoCを標準ライブラリで作り、性能・可読性・開発コストを見てから外部依存を検討するのが実務的です。

ハンズオン:主要モジュールと短い実例コード

ここでは実務でよく使うモジュールを例とともに紹介します。各コードはそのまま貼って試せる最小例です。

モジュール 用途
pathlib ファイル/パス操作(プラットフォーム非依存)
logging + RotatingFileHandler ログ設計とローテーション
configparser / os.environ 設定管理(ファイル/環境変数)
csv / json 軽量データ入出力
itertools ストリーム処理(メモリ節約の変換)
collections defaultdict, Counter, deque などの便利構造
functools.lru_cache 簡易キャッシュ
subprocess 外部コマンド実行
tempfile / shutil 一時ファイルとアーカイブ

pathlib(パス操作)

from pathlib import Path
p = Path('data')
for f in p.glob('*.csv'):
    print(f.resolve())

logging + RotatingFileHandler(ログ設計)

import logging
from logging.handlers import RotatingFileHandler

logger = logging.getLogger('myapp')
handler = RotatingFileHandler('logs/app.log', maxBytes=10_000_00, backupCount=5)
formatter = logging.Formatter('%(asctime)s %(levelname)s %(message)s')
handler.setFormatter(formatter)
logger.addHandler(handler)
logger.setLevel(logging.INFO)

logger.info('start')

configparser / os.environ(設定管理)

import configparser
import os

config = configparser.ConfigParser()
config.read('config.ini')
api_key = os.environ.get('API_KEY') or config.get('default', 'api_key', fallback=None)

csv / json(軽量データ入出力)

import csv
import json

# CSV読み取り
with open('input.csv', newline='', encoding='utf-8') as f:
    reader = csv.DictReader(f)
    for row in reader:
        print(row['id'], row['value'])

# JSON書き出し
with open('out.json', 'w', encoding='utf-8') as f:
    json.dump({'ok': True}, f, ensure_ascii=False, indent=2)

itertools(ストリーミング変換)

import itertools

def chunks(iterable, n):
    it = iter(iterable)
    while True:
        chunk = list(itertools.islice(it, n))
        if not chunk:
            break
        yield chunk

for c in chunks(range(10), 3):
    print(c)

collections(defaultdict / Counter / deque)

from collections import defaultdict, Counter, deque

cnt = Counter(['a','b','a'])
print(cnt.most_common())

d = defaultdict(list)
d['k'].append(1)

q = deque(maxlen=100)
q.append(1)

functools.lru_cache(キャッシュ)

from functools import lru_cache

@lru_cache(maxsize=128)
def expensive(x):
    # 重い計算や外部呼び出しの結果をキャッシュ
    return x * x

print(expensive(3))

subprocess(外部コマンド実行)

import subprocess

# 引数をリストで渡してシェルインジェクションを避ける
res = subprocess.run(['ls', '-la'], capture_output=True, text=True)
print(res.stdout)

tempfile / shutil(一時ファイルとアーカイブ)

import tempfile
import shutil
from pathlib import Path

with tempfile.TemporaryDirectory() as td:
    tmp = Path(td) / 'data.txt'
    tmp.write_text('hello')
    shutil.make_archive('archive', 'gztar', root_dir=td)

実務パターン:すぐ使えるテンプレート3本

ここではテンプレートごとに必要モジュール、設定例、ログ設定、最小テスト例を示します。

1) CLIバッチ用ユーティリティ(ファイル変換等)

  • 主なモジュール:argparse / pathlib / csv / logging / configparser
  • 設定ファイル例(config.ini):
[default]
input_dir = ./data
output_dir = ./out
log_file = ./logs/app.log
  • 簡易ログ設定:前述のRotatingFileHandlerを使う
  • 最小テスト(pytest)例:
def test_cli_process(tmp_path):
    # 入出力のダミーファイルを作って処理が期待通り動くかを確認
    in_file = tmp_path / 'input.csv'
    in_file.write_text('id,value\n1,100\n')
    # process_cliは実装関数
    process_cli(str(in_file), str(tmp_path / 'out'))
    assert (tmp_path / 'out' / 'input.csv').exists()

2) ファイル監視→処理の小型ワーカー

  • 主なモジュール:pathlib / time / logging / collections.deque / tempfile
  • 構成:監視ループで新着ファイルをキュー(deque)に入れ、ワーカーが順次処理
from collections import deque
from pathlib import Path
import time

queue = deque()
watch = Path('inbox')
while True:
    for f in watch.glob('*.csv'):
        queue.append(f)
        f.rename(f.with_suffix('.processing'))
    if queue:
        process(queue.popleft())
    time.sleep(5)
  • 設定は環境変数優先(os.environ)にして、config.iniはフォールバックとする
  • テスト:ファイルの到着と処理開始をモックまたは一時ディレクトリで確認

3) 日次アーカイブ+アップロードスクリプト

  • 主なモジュール:shutil / tempfile / subprocess / logging / datetime
  • 処理の流れ:古いファイルを集めて一時ディレクトリで圧縮→外部コマンド(curl/s3cli等)でアップロード→クリーンアップ
import shutil
from datetime import datetime
from pathlib import Path

src = Path('logs')
archive_name = f'logs-{datetime.today().date().isoformat()}'
shutil.make_archive(archive_name, 'gztar', root_dir=src)
# subprocessでアップロードコマンドを実行

既存ワークフローとの接続例

過去記事との連携例をテーブルで示します。

ステップ 該当記事 標準ライブラリでの実装例
CSV入出力 第134回 csv.DictReader → itertoolsでストリーム変換
CSV前処理 第135回 collections.Counterで簡易集計後、jsonで保存
推論API登録 第139回 subprocessやurllibでバッチ登録リクエスト

例:CSVクリーニング→itertoolsでストリーム処理→subprocessで外部ETL呼び出し→FastAPIにバッチ登録

import csv
import itertools
import subprocess
from urllib import request, parse

# 1. ストリームで読み、変換
with open('in.csv', newline='') as f:
    reader = csv.DictReader(f)
    cleaned = ({'id': r['id'], 'v': float(r['value'])} for r in reader if r['value'])
    # 2. バッチ処理(外部ETLを呼ぶ)
    for chunk in itertools.islice(cleaned, 1000):
        subprocess.run(['etl_tool', '--input', '-'], input=str(chunk).encode())
# 3. FastAPI に登録
data = parse.urlencode({'job': 'daily'}).encode()
req = request.Request('https://example/api/register', data=data)
request.urlopen(req)

運用チェックリスト

項目 確認ポイント
ログローテーション 適切なmaxBytes/backupCountが設定されているか
ログレベル運用 本番はINFO以上、デバッグ時にDEBUGを有効化する仕組み
例外ハンドリング リトライポリシーとアラート(最低限stdout/exit codeで監視可能に)
一時ファイルのクリーンアップ TemporaryDirectoryやtry/finallyで確実に削除
メトリクス出力 標準出力へ簡易メトリクス(処理数、エラー数)を出す

実装上の注意と落とし穴

  • subprocessの呼び出しは引数をリストで渡し、ユーザー入力を直接渡さない(シェルインジェクション対策)。
  • itertoolsでメモリ節約をしても、中間でlist化すると効果が消える。isliceやジェネレータを意識する。
  • loggingはモジュール間でグローバルに影響するため、テストではログ設定のリセットを行う。
  • 標準ライブラリで性能限界に達したら、代替ライブラリ導入の判断基準(処理時間、メモリ、開発コスト)を明確にする。

成果物と次の一歩

本記事で紹介したテンプレートはGitHub向けに最小リポジトリを用意すると実務導入が早まります。例:

ファイル/ディレクトリ 説明
README.md 導入手順と設定例
config.ini サンプル設定ファイル(ENV優先)
src/ スクリプト本体(__main__.py とモジュール)
tests/ pytest用の最小テスト
.github/workflows/ci.yml pytestを回すCI設定(例)

WordPressに貼りやすいコードブロック例(すでに本文中で複数提示)をそのままコピーして使ってください。

次回候補:itertools+collectionsの深掘り、標準ライブラリだけで作る小さなCLIアプリ公開フロー。

まとめ

標準ライブラリは、依存管理や保守性、セキュリティ面で強い利点があります。まず標準でプロトタイプを作り、実際の負荷や運用を見てから外部ライブラリ導入を判断するのが実務的です。本稿のテンプレートやチェックリストを使えば、すぐに試せるツールが用意できます。まずは小さく始めて、運用で得た知見を元に改善していきましょう。

第139回 実務で使える軽量推論API入門:FastAPIで作るバッチ&オンデマンド推論と運用チェックリスト

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

「ローカルでモデルは動くが、実際の業務にAPIとして出すと想定外の問題が出る」——こうした悩みはよく聞きます。本稿では、短時間で動く軽量な推論APIをFastAPIで設計・実装し、オンデマンドとバッチの両方を同一サービスで扱うパターンと、運用で必要なチェック項目を手順で示します。実務でそのまま使える観点を優先し、落とし穴と対処法も整理します。

1) 要件定義と設計(オンデマンド vs バッチ)

まず設計方針を明確にしましょう。用途によってSLA・レイテンシ・コスト感が変わります。

用途 要求 実装上の着眼点
オンデマンド推論(同期) 低レイテンシ(数百ms〜数秒) 非同期I/O、リクエストタイムアウト、スケール方針
バッチ推論(大量一括) スループット重視、コスト効率 バッチング、バックグラウンドキュー、バッチサイズ制御
混在ケース 同一サービスで両立 優先度制御、リソース隔離、レート制限

設計時に決めるべき最小項目:

  • SLA(最大許容遅延、成功率)
  • 一回の推論当たりコスト想定
  • 同時実行数とスケール戦略(垂直 vs 水平)
  • 冪等性と再試行ポリシー

2) FastAPIの最小セットアップとエンドポイント設計

入力検証にはpydanticを使い、落ちる可能性を早期に検出します。ここではエンドポイント設計をテーブルで示します(コード例はポイントのみ記述)。

エンドポイント メソッド 目的 備考(入力検証)
/predict POST オンデマンド推論(同期応答) pydanticモデルで必須フィールド・型チェック
/predict_async POST 非同期キュー登録(バッチング対象) リクエストIDとコールバックURLを受け取る設計が実務向け
/batch/status/{job_id} GET バッチジョブの状態取得 ステータスとエラー情報を返す
/health, /ready GET ヘルスとレディネスチェック 監視用に簡潔な応答を返す

pydantic例(説明のみ):入力モデルはサイズ制限、文字列長、必須キーを明示する。大きなペイロードは事前チェックで弾く。

3) バッチング戦略 — バックグラウンドキュー+集約タイムウィンドウ

バッチ処理は「一定時間で集めて一度に処理する」か「サイズでトリガする」方式が基本です。ここでは簡易ワーカーと非同期キューの考え方を示します。

方式 特徴 向き不向き
時間ウィンドウ(例:100msごと) レイテンシばらつきが小さくバッチ効率が良い 短めの応答要件を満たしつつスループットを稼ぎたい場合
サイズトリガ(例:batch_size=32) バッチ効率が最大化されるが待ちが発生しうる バッチ性能を最優先にする業務
ハイブリッド どちらのしきい値でもトリガ可能 実運用でよく使われる

実装ヒント(要点のみ、簡潔に):

  • APIは受信時にリクエストを軽くバリデートし、非同期キューに格納する。
  • バックグラウンドタスクがタイマーやサイズでバッチを切り出し、推論ワーカーに渡す。
  • ワーカーは外部モデルAPI呼び出し/ローカルモデル推論を行い、結果を保存またはコールバック。
  • 非同期実装はasyncio.QueueやRedis/RQを利用すると堅牢。単純構成ならuvicornのバックグラウンドタスクで十分。

4) エラー処理と再試行方針

実運用では予期せぬ例外、タイムアウト、外部API障害が起きます。方針を明確にしておきましょう。

  • 例外ハンドラでHTTP 500を返す前に、詳細はログに残す。クライアントには説明的なエラーコードとメッセージを返す。
  • 再試行は冪等性が担保できる場合のみ行う(重複結果の影響を考慮)。指数バックオフ+最大試行回数を設定する。
  • タイムアウトはAPIゲートウェイとアプリ双方で設定する(例:APIは短め、バッチジョブでは長めに)。
  • 長時間処理のバッチはジョブ状態を保存し、失敗時に部分成功を扱えるようにする。

5) ロギング・メトリクス・ヘルスチェック

監視のための最低限のエンドポイントとログ設計を示します。

名称 用途 実装例(返す情報)
/health プロセスが稼働しているか(簡易) {“status”: “ok”}
/ready 依存(モデルロード、外部サービス接続)が整っているか {“ready”: true, “model_loaded”: true}
/metrics Prometheus互換のメトリクス公開 リクエスト数、エラーレート、レイテンシヒストグラム

ログ設計のポイント:

  • 構造化ログ(JSON)でリクエストID、ジョブID、処理時間、エラー詳細を出力する。
  • ログレベルは運用と開発で分ける。デバッグはファイル/外部でのみ残す。
  • 例外はスタックトレースを追えるように保存しつつ、機密情報はマスクする。

6) Docker化と簡易CI

軽量APIはコンテナ化してデプロイするのが運用しやすいです。CIは次の流れを推奨します。

  • pytestでFastAPI TestClientを使ったエンドポイントテストを用意する(/predictの成功系・異常系、バッチ登録の統合テスト)。
  • Dockerfileは小さなベースイメージ(python:3.x-slim)を使い、依存は最小化する。
  • CI(例:GitHub Actions)でテスト→イメージビルド→(必要なら)イメージのスキャン→registryへpushの流れを作る。

CIに含める最小ステップ(例):

ステップ 目的
pytest実行 機能回帰チェック
lint(optional) コード品質チェック
docker build & push イメージ化とレジストリ反映

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

以下は現場ですぐチェックできる実用的な項目です。最後に表で整理します。

  • リソース制限(CPU/メモリ)の設定とOOM対策
  • レートリミットと優先度制御(オンデマンド優先/バッチはスロット確保)
  • セキュリティ:TLS、認証、シークレット管理
  • 監視:メトリクス収集、アラートルール(エラーレート、レイテンシ、スループット)
レベル 項目 チェック内容(そのまま使える)
必須 ヘルス/レディネス /health と /ready を監視に登録し、復旧不能時は自動再起動
必須 リクエストタイムアウト APIゲートウェイとアプリに最大タイムアウトを設定(例:30s/300s)
必須 ログレベル・フォーマット JSONログでリクエストIDを必ず出力
推奨 メトリクス Prometheus互換でリクエスト数/エラー/レイテンシを収集
推奨 リソース制限 コンテナごとにCPU/メモリ上限を設ける
発展 認証・シークレット管理 VaultやクラウドKMSでキー管理、APIはBearerトークンで保護

8) よくある落とし穴とトラブルシューティング

実務で遭遇しやすい問題と対処法を短くまとめます。

  • 問題:バッチの遅延が増える。対処:ウィンドウ幅を短縮、優先度でオンデマンド優先にする。
  • 問題:メモリ不足でOOM。対処:バッチサイズ制限、ワーカー数削減、モデルを軽量化。
  • 問題:外部APIの不安定。対処:タイムアウト、回路遮断(circuit breaker)、キャッシュ導入。
  • 問題:再試行で重複データ。対処:リクエストIDによる冪等性チェック、または重複排除ロジック。

まとめ

本稿では、FastAPIを使った軽量推論APIの設計と実装の考え方を、オンデマンドとバッチを一つのサービスで扱う観点から整理しました。重要なポイントは次の通りです:

  • 設計段階でSLAとコスト感を明確にすること
  • 入力検証(pydantic)と構造化ログで問題の早期把握を行うこと
  • バッチは時間ウィンドウとサイズトリガのハイブリッドがおすすめで、非同期キューを使って安全に処理すること
  • ヘルス/レディネス/メトリクスを整備し、CIでテスト・コンテナ化を自動化すること

最後に、運用チェックリストを実運用で回しながら改善を続けることが、安定稼働への近道です。次回は第138回で触れたテスト方針を踏まえ、具体的なpytestによるエンドツーエンドテストの例を紹介します。

第138回 実務で使えるテストとCI:pytest・モック・テストデータでAIワークフローの品質を担保する手順

AIワークフローを業務に組み込もうとすると、「モデルは動くけれど、変更すると突然動かなくなる」「外部APIの遅延で処理が止まる」「実運用データでの想定外ケースに気づけない」といった悩みに直面しがちです。本記事では、Pythonで実装した日常的なAIワークフロー(CSV前処理、API/LLM呼び出し、推論バッチ)を対象に、実務で使えるテスト戦略とCIの構成を手順とコード例で解説します。現場で再現しやすい具体例に沿って進めますので、まずは小さく始め、運用に耐える品質を徐々に高めていきましょう。

本記事の対象と目的

想定読者:AIを仕事で活用したい実務担当者、個人事業主、中小企業の担当者。目的は、既存のPythonワークフローにテストとCIを導入して運用リスクを下げ、変更時の回帰を防ぐことです。

章立て(概要)

  • テスト戦略の決め方(ユニット/統合/エンドツーエンド)
  • pytestを使った単体テストの書き方(fixtures, parametrize)
  • 外部API/LLM呼び出しのモック手法
  • テストデータの生成と管理(faker, factory_boy, CSVサンプル)
  • 統合テストと小型データセットでの実行方針
  • GitHub ActionsでのCIパイプライン例
  • 運用時の注意点と失敗パターン
  • 導入チェックリストとテンプレート/30分クイックスタート

1. テスト戦略(どのテストをどこまで書くか)

まずは役割を明確にします。小さなチームや一人運用では、全てのテストを同時に充実させるのは非現実的です。優先度を決めて段階的に増やすと良いです。

テスト種別 対象 速度 導入優先度(実務)
ユニットテスト 関数単位、前処理ロジック 速い CSVの列変換、正規化関数
統合テスト モジュール間連携、データパイプライン CSV読み込み→前処理→推論ラッパー
エンドツーエンド(E2E) 本番に近いフロー(API/外部サービス含む) 遅い 小さなデータセットでのバッチ実行 低→定期実行推奨

実務的な優先度の目安

  • まずユニットテスト:データ前処理、ビジネスロジック、スコア計算。
  • 次に統合テスト:モジュールの連携、不整合を早期検出。
  • E2Eは定期実行(夜間/週次)で本番寄せの検証を行う。

2. pytestを使った単体テストの書き方(基本とコツ)

pytestはfixtureやparametrizeで可読性の高いテストが書けます。ここではCSV前処理関数の例を示します。

CSV前処理(例)

def clean_row(row):
    # 入力: dict(CSVの1行)
    # 戻り値: 正規化された dict
    name = row.get("name", "").strip()
    age = row.get("age")
    try:
        age = int(age)
    except (TypeError, ValueError):
        age = None
    return {"name": name, "age": age}

pytestの例

import pytest
from myproject.csv_utils import clean_row

@pytest.mark.parametrize(
    "input_row,expected",
    [
        ({"name": " Alice ", "age": "30"}, {"name": "Alice", "age": 30}),
        ({"name": "", "age": "x"}, {"name": "", "age": None}),
    ],
)
def test_clean_row_param(input_row, expected):
    assert clean_row(input_row) == expected

@pytest.fixture
def sample_row():
    return {"name": " Bob ", "age": "45"}

def test_clean_row_fixture(sample_row):
    out = clean_row(sample_row)
    assert out["name"] == "Bob"
    assert isinstance(out["age"], int)

assertion のコツ

  • 等価性は == を使い、浮動小数点は pytest.approx を使う。
  • 複雑なオブジェクトは必要最小限に比較(キーのみや主要値のみ)。
  • 失敗時に読みやすいメッセージが出るよう、pytestのassertをご活用ください。

3. 外部API/LLM呼び出しのモック

外部依存はテストを不安定にするため、ユニット/統合レベルで適切にモックします。代表的な手法と使い分けを下表に示します。

ライブラリ 用途 向く場面
unittest.mock 関数/メソッドの置換(軽量) 内部ラッパーの戻り値制御、エラー再現
responses requestsによるHTTPレスポンスのモック 外部REST APIのユニットテスト
requests-mock requestsのSession単位でのモック 細かいHTTP振る舞いを制御したい場合
VCR.py 実際のHTTP交流を録画して再生 既知のAPI応答をキャッシュして再現性を得たいとき

外部APIラッパーの例とテスト

# myproject/api_client.py
import requests

class APIClient:
    def __init__(self, base_url):
        self.base_url = base_url

    def get_score(self, payload):
        r = requests.post(f"{self.base_url}/score", json=payload, timeout=5)
        r.raise_for_status()
        return r.json()
# tests/test_api_client.py
from unittest.mock import patch
from myproject.api_client import APIClient

@patch("myproject.api_client.requests.post")
def test_get_score_mock(mock_post):
    mock_resp = mock_post.return_value
    mock_resp.json.return_value = {"score": 0.9}
    mock_resp.raise_for_status.return_value = None

    client = APIClient("https://api.example")
    out = client.get_score({"x": 1})
    assert out["score"] == 0.9
    mock_post.assert_called_once()

よりHTTP層で試したい場合は responses を使います:

import responses
from myproject.api_client import APIClient

@responses.activate
def test_get_score_responses():
    responses.add(
        responses.POST,
        "https://api.example/score",
        json={"score": 0.8},
        status=200,
    )
    client = APIClient("https://api.example")
    assert client.get_score({"x": 1})["score"] == 0.8

4. テストデータの生成と管理

良いテストは良いデータから生まれます。現実の代表ケースと境界値、エッジケースを意図的に含めましょう。個人情報(PII)を含む実データは必ず合成化するか、マスクして使用してください。

ルール 説明
代表ケース抽出 現場データから頻度の高いパターンを抽出してサンプル化
異常値/境界値 欠損、極端値、異常文字列(例: 半角/全角混在)を含める
PIIの扱い 本番データを使う場合は必ず匿名化または合成データで代替
小さなCSVサンプル 数十行の代表的セットを用意して統合テストに利用

faker と factory_boy の簡単例

from faker import Faker
fake = Faker()

def make_user_record():
    return {"name": fake.first_name(), "email": fake.email(), "age": fake.random_int(18, 80)}

生成したレコードはCSVに書き出して統合テスト用の小データとして利用できます。

小さなCSVサンプル設計(例)

行種別 説明
正常 代表的な1〜2行(標準的な列と値)
欠損 必須列が欠けている行(ageが空など)
境界 年齢=0、最大値、長い文字列など
異常 数値が文字列、特殊文字を含む列

5. 統合テストの実行方針

統合テストは重くなりがちなので、次の方針を推奨します。

  • ローカルではユニットテストを高速に回す(pre-commit / pre-pushで実行)。
  • 統合テストは小さなデータセットで速く回るように設計し、CIでは長いE2Eは別ジョブ・夜間に実行。
  • 外部APIはスタブやVCRで再現性を担保。定期的に実際のAPIで罹患テストを走らせる。

統合テスト(ワークフロー実行)の簡単な例

# myproject/workflow.py
import csv
from myproject.csv_utils import clean_row
from myproject.api_client import APIClient

def run_batch(csv_path, api_client):
    results = []
    with open(csv_path, newline="", encoding="utf-8") as f:
        reader = csv.DictReader(f)
        for r in reader:
            r2 = clean_row(r)
            res = api_client.get_score(r2)
            results.append(res)
    return results
# tests/test_workflow_integration.py
from myproject.workflow import run_batch
from types import SimpleNamespace

def test_run_batch(tmp_path):
    csv_file = tmp_path / "sample.csv"
    csv_file.write_text("name,age\nAlice,30\nBob,\n")

    class DummyClient:
        def get_score(self, payload):
            return {"score": 0.5 if payload.get("age") else 0.0}

    out = run_batch(str(csv_file), DummyClient())
    assert len(out) == 2
    assert out[0]["score"] == 0.5
    assert out[1]["score"] == 0.0

6. GitHub ActionsでのCIパイプライン(テンプレート)

ここではPRで自動的にpytestを回し、カバレッジを計測して報告する最小テンプレートを示します。ファイル名は .github/workflows/ci.yml を想定してください。

name: CI

on:
  pull_request:
  push:
    branches: [ main ]

jobs:
  test:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - name: Set up Python
        uses: actions/setup-python@v4
        with:
          python-version: '3.10'
      - name: Install dependencies
        run: |
          python -m pip install --upgrade pip
          pip install -r requirements-dev.txt
      - name: Run tests with coverage
        run: |
          pytest --junitxml=reports/junit.xml --cov=myproject --cov-report=xml
      - name: Upload coverage report
        uses: actions/upload-artifact@v4
        with:
          name: coverage-report
          path: coverage.xml

実務では coverage の結果を Codecov や Coveralls と連携したり、PR上で失敗を通知する仕組み(Reviewdog や GitHub Checks)を追加します。

7. 運用時の注意点と失敗パターン

問題 原因 対策
フレークテスト(不安定なテスト) 外部依存、時間依存、並列性の問題 依存のモック化、retry回避、テストの分離
テストが遅い 大きなデータセット、E2Eを頻繁に実行 ユニットと統合を分離、夜間にE2Eを実行
カバレッジ偏重 量だけ増えて中身が薄いテスト 重要ロジックに対する深いアサーション、コードレビューでの品質担保

品質指標とゲート条件(目安)

指標 目安
全体テストカバレッジ 70〜85%(業務重要度に応じて高める)
重要モジュール 90%前後を目指す
PRゲート ユニットテスト全通+カバレッジの変化が大きい場合は要確認

CIで失敗した時の簡単なエスカレーションフロー

  • 1) PRに修正コミットを追加して再実行(まず試す)
  • 2) それで直らない場合は失敗ログを貼って担当者に@メンション
  • 3) 夜間のE2Eで発生した場合は運用チームにチケット登録し、翌稼働日までブロッキング回避策を適用

8. 導入チェックリストとテンプレート

項目 完了 注記
requirements-dev.txt に pytest, pytest-cov を追加 [ ] 最低限のテスト実行に必要
少なくとも1つのユニットテストを追加 [ ] CSV前処理などビジネスロジック
CIワークフローを追加(.github/workflows/ci.yml) [ ] PRでテストが回るように
テストデータ(小さなCSV)をリポジトリに保存 [ ] 例: tests/data/sample.csv(PIIに注意)
外部APIをモックするガイドを作成 [ ] responses や unittest.mock の例を含める

9. 30分でできるクイックスタート(実践手順)

  1. requirements-dev.txt に次を追加:pytest, pytest-cov, responses、保存(5分)。
  2. リポジトリに tests/test_csv_utils.py を1つ作成し、上記の clean_row テストを貼る(10分)。
  3. .github/workflows/ci.yml を追加し、上記テンプレートを貼る(10分)。
  4. PRを作って動作を確認。失敗したらログを見て修正して再コミット(5分)。

補助リソースと推奨ライブラリ(導入順)

ライブラリ 目的
1 pytest, pytest-cov テスト実行とカバレッジ測定
2 unittest.mock 軽量なモック・パッチ
3 responses / requests-mock HTTP APIのモック
4 faker, factory_boy テスト用データ生成
5 VCR.py HTTP録画再生(必要時)
6 GitHub Actions CIの実行環境

まとめ

テストとCIは、運用リスクを低減し、変更時の安心感を高めます。現場では小さく始めることが重要です。まずは主要な前処理やスコア計算に対するユニットテストを1つ書き、CIで自動実行する流れを作ってください。次に統合テストを小さなデータセットで追加し、外部依存はモックで切り離します。最終的には定期的に実運用に近いE2Eを走らせることで、運用上の齟齬を早期に発見できます。

Manage AI(https://manageai.online)では、今回のような実務寄りの手順を今後も扱っていきます。まずはこの記事のクイックスタートに従い、既存リポジトリに1つのユニットテストとCIを追加してみてください。

次の一歩(推奨):ローカルで pytest が通ることを確認したら、PR を作成して GitHub Actions で自動化される流れを体験してください。困ったときはログ全文を保存してチームで共有すると原因特定が速くなります。