第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)では、実務に使える小さなツール群とテンプレートを継続的に紹介します。まずはこの記事のサンプルパターンを試し、貴社の運用要件に合わせてカスタマイズしてください。