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