はじめに — 現場でよくあるつまずきに寄り添う
大量の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_retryとstop_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)の設定を先に確認すると運用へスムーズにつなげられます。