Pythonのログ集計を並列化しても速くならない原因と対処|ThreadPoolExecutorとProcessPoolExecutorの使い分け・GIL・max_workers・chunksize・Windowsのif name

Pythonのログ集計を並列化しても速くならない原因と対処

こんにちは、おうどんです🍜

うどん屋の閉店後。 今日の注文ログを集計します。

1日ぶんで約12MB、それが8日ぶん。合わせて100MBほど。 Python で1ファイルずつ読んで、ステータスコードを数えるだけ。……なのに、ちょっと待たされます。

店長(おうどん)「遅いなあ。店員を増やそう。8人でいっせいに数えれば8倍速いでしょ」

スレッドさん「8人集まりました!」

店長「よし、何秒?」

スレッドさん「3.91秒です」

店長「1人のときは?」

スレッドさん「3.84秒です」

……増やしたら、ちょっと遅くなった。

(会話は説明用の架空の場面です。数字は今回の検証で実際に計測したものです)

これ、Python ではよくある話なんです。 しかも仲間がいます。

  • スレッドを増やしても、集計が速くならない
  • プロセスにしたら速くなったけど、4並列で1.5倍くらいしか速くならない
  • 8並列にしたら、4並列より遅くなった
  • Windows で RuntimeError: An attempt has been made to start a new process... が出る

今回は、店員を増やしても速くならない理由と、Python の並列処理の使い分けを、実際に計測しながら見ていきます。

最初に、初心者向けのチェックリストです。

  • 計算や文字列処理が中心(CPU を使う仕事)なら ProcessPoolExecutor、待ち時間が中心(ダウンロードなど)なら ThreadPoolExecutor
  • 並列数は、まず os.process_cpu_count()(使える論理CPUの数)までにする
  • プロセスを使うスクリプトは、if __name__ == "__main__": の中で動かす
目次

スレッドを8人にしても速くならず、プロセスにすると速くなった

結論から見せます。 同じ集計でも、スレッドではほぼ変わらず、プロセスでは速くなりました。ただし、増やせば増やすほど速くなるわけではありません。

検証用に、架空のアクセスログを8ファイル作りました。1ファイル20万行、合計約100MBです。1行はこんな形です。

2026-10-01T00:00:00.123+09:00 192.0.2.10 GET /menu/kitsune 200 123ms

(ログの中身は乱数で作った架空のデータです。IP アドレスは説明用の予約範囲 192.0.2.0/24 を使っています)

これを、1ファイル=1つの仕事として、直列(1つずつ)・スレッド・プロセスで集計して比べます。

import glob
import os
import re
import time
from collections import Counter
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor

LINE = re.compile(r"^\S+ \S+ (\w+) (\S+) (\d{3}) (\d+)ms$")


def analyze(path):
    """Count status codes and sum response times for one log file."""
    status = Counter()
    total_ms = 0
    with open(path, encoding="utf-8") as f:
        for line in f:
            m = LINE.match(line.rstrip("\n"))
            if m:
                status[m.group(3)] += 1
                total_ms += int(m.group(4))
    return status, total_ms


def merge(results):
    status, total_ms = Counter(), 0
    for s, ms in results:
        status.update(s)
        total_ms += ms
    return status, total_ms


def run(label, make_pool):
    """Time 3 rounds. Each round includes starting and closing the pool."""
    secs = []
    for _ in range(3):
        start = time.perf_counter()
        if make_pool is None:
            status, total_ms = merge(map(analyze, files))
        else:
            with make_pool() as ex:
                status, total_ms = merge(ex.map(analyze, files))
        secs.append(time.perf_counter() - start)
    times = " ".join(f"{s:5.2f}" for s in secs)
    print(f"{label:<12} {times}  (500={status['500']} total_ms={total_ms})")


if __name__ == "__main__":
    files = sorted(glob.glob("logs/access-*.log"))
    print(f"files={len(files)} cpu_count={os.cpu_count()} process_cpu_count={os.process_cpu_count()}")

    run("serial", None)
    for n in (2, 4, 8):
        run(f"thread x{n}", lambda: ThreadPoolExecutor(max_workers=n))
    for n in (2, 4, 8):
        run(f"process x{n}", lambda: ProcessPoolExecutor(max_workers=n))
files=8 cpu_count=4 process_cpu_count=4
serial        3.84  3.60  3.58  (500=48274 total_ms=722558839)
thread x2     4.00  4.25  4.05  (500=48274 total_ms=722558839)
thread x4     3.89  3.93  4.01  (500=48274 total_ms=722558839)
thread x8     3.91  3.96  5.15  (500=48274 total_ms=722558839)
process x2    3.90  3.73  3.94  (500=48274 total_ms=722558839)
process x4    2.35  2.66  2.53  (500=48274 total_ms=722558839)
process x8    2.90  2.85  2.88  (500=48274 total_ms=722558839)

(今回の検証環境の Windows 10・Python 3.13.1、2コア4スレッドのノートPC向けCPU で実行した結果です。数字は秒で、各3回。プールの起動と終了の時間も含みます)

読み方はこうです。

  • 右側の 500=48274 と total_ms=... は、どのやり方でも同じ。集計結果は正しく、速さだけが違う
  • スレッドは2人でも8人でも、直列(約3.6〜3.8秒)とほぼ同じか、少し遅い
  • プロセスは4並列で約2.4〜2.7秒。速くはなったけれど、4倍ではない
  • プロセス8並列は約2.9秒で、4並列より遅い

店長「8人雇って、1人のときより遅いの?」

スレッドさん「全員で1本の包丁を回してたので……」

……厨房に包丁が1本しかない店で、店員だけ増やしていた。

この包丁が、次の章の主役です。

じゃあ、ぜんぶプロセスにしておけば安心?

そうでもないんです。 待ち時間が中心の仕事では、逆にスレッドのほうが圧勝します。順番に見ていきましょう。

スレッドで速くならないのは、GIL という包丁が1本だから

スレッドさんの言い訳は、半分ほんとうです。 ふつうの Python(CPython)では、Python のコードを実行できるスレッドは、同時に1つだけなんです。

CPython は、python.org から入れるふつうの Python のことです。 Python 用語集の global interpreter lock(GIL)の項には、GIL は「一度に1つのスレッドだけが Python のバイトコードを実行する」ことを保証するしくみで、インタープリター全体をロックするぶん、マルチプロセッサーのマシンで得られる並列性の多くを犠牲にしている、と書かれています。

バイトコードは、Python が実行する直前に変換した、中間の命令のことです。 つまり、スレッドを何本作っても、Python の行を進められるのは常に1本ずつ。今回の集計は、1行ずつ正規表現にかけて数を足す、Python のコードがずっと動く仕事なので、スレッドを増やしても順番待ちが増えるだけでした。

スレッドさん「包丁、わたしの番まだですか?」

別のスレッドさん「いま3人前です」

店長「……順番待ちの列、レジより長くない?」

……人件費だけが並列化されている。

同じ用語集には、こうも書かれています。

  • 圧縮やハッシュ計算のような重い処理では、GIL を手放すように作られた拡張モジュールもある
  • 入出力(I/O)の最中は、GIL は必ず手放される

I/O は、ファイルやネットワークの読み書きのことです。 この2つ目が、スレッドさんの見せ場につながります。

じゃあスレッドって、何のためにあるの?

待つためです。 ……と書くと、すごく情けない店員に見えますが、待つのがうまい店員はとても役に立ちます。

待ち時間が主役の仕事なら、スレッドが圧勝する

ダウンロードや API の応答待ちのように、CPU を使わずに待っている時間が長い仕事なら、スレッドで並列にするだけで速くなります。

ログをサーバーから取ってくる仕事を、0.5秒待つだけの関数でまねして比べます。

import time
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor


def fetch(name):
    """Pretend to download one log from a remote server (waits 0.5 s)."""
    time.sleep(0.5)
    return name


def run(label, func):
    start = time.perf_counter()
    got = list(func())
    print(f"{label:<12} {time.perf_counter() - start:5.2f}s  files={len(got)}")


if __name__ == "__main__":
    names = [f"access-{i:02d}.log" for i in range(1, 9)]
    run("serial", lambda: map(fetch, names))
    with ThreadPoolExecutor(max_workers=8) as ex:
        run("thread x8", lambda: ex.map(fetch, names))
    with ProcessPoolExecutor(max_workers=4) as ex:
        run("process x4", lambda: ex.map(fetch, names))
serial        4.03s  files=8
thread x8     0.52s  files=8
process x4    1.33s  files=8

(今回の検証環境の Python 3.13.1 で実行した結果です。time.sleep で待ち時間をまねたもので、実際のダウンロードはしていません)

直列は 0.5秒 × 8回で約4秒。 スレッド8本なら、8つの待ちが同時に進むので約0.5秒。 プロセス4つは、4つずつ2回待つので1秒、それに起動の手間が乗って1.33秒でした。

スレッドさん「待つのは得意です。8人で同時に、出前の電話を待てます」

店長「さっきと別人みたいに頼もしい」

……包丁を使わない仕事なら、全員が同時に働ける。

まとめるとこうです。

仕事の種類例向いているもの
CPU が中心(CPU バウンド)ログの解析、正規表現、集計、圧縮していない計算ProcessPoolExecutor
待ちが中心(I/O バウンド)ダウンロード、API 呼び出し、DB の応答待ちThreadPoolExecutor

CPU バウンドは「CPU の計算がボトルネック(いちばん遅いところ)」、I/O バウンドは「待ちがボトルネック」という意味です。

店長「じゃあ、電話番を100人にしたら、100本の出前を0.5秒で受けられる?」

……受けられても、茹でる人が4人なので、厨房の前に100杯ぶんの行列ができます。待ちが速くなっても、そのあとの仕事が詰まるだけ、ということもあるんです。

multiprocessing のドキュメントの冒頭にも、multiprocessing はスレッドの代わりにサブプロセス(子プロセス)を使うことで GIL を事実上回避し、マシンの複数のプロセッサーを十分に活用できる、と書かれています。ProcessPoolExecutor は、このしくみの上に作られた使いやすい窓口です。

ログの集計は、ファイルを読むから I/O バウンドじゃないの?

いい質問です。 今回のログは、1回目の実行でファイルの中身が OS のキャッシュ(メモリ上の一時置き場)に乗るので、読み込みの待ちはほとんどありません。重いのは、1行ずつの正規表現と足し算のほうでした。スレッドで速くならなかったのが、その証拠です。ネットワーク越しの共有フォルダーなど、読み込みが遅い場所のログなら、話は変わってきます。

並列数は、まず「使える論理CPUの数」まで

じゃあ何並列にするか。 CPU が中心の仕事なら、まずは os.process_cpu_count() の数まで。それより増やしても、今回は遅くなりました。

os モジュールのドキュメントでは、2つの関数がこう区別されています。

  • os.cpu_count():システム全体の論理CPUの数
  • os.process_cpu_count():いまのプロセスの呼び出し元スレッドが使える論理CPUの数(Python 3.13 で追加)。CPU アフィニティ(このプロセスはこの CPU だけ、という割り当て)によっては cpu_count() より少なくなる

論理CPUは、OS から見た CPU の数です。1つのコアが2つぶんの仕事を受け付ける機能(ハイパースレッディングなど)があると、コアの数より多く見えます。今回の検証環境は2コア4スレッドで、どちらの関数も 4 を返しました。

そして、引数を省略したときの既定値です。concurrent.futures のドキュメントによると、

  • ProcessPoolExecutor の max_workers を省略すると os.process_cpu_count()(3.13 から。それより前は os.cpu_count())
  • ThreadPoolExecutor は min(32, (os.process_cpu_count() or 1) + 4)(3.13 から)。今回の環境なら 4 + 4 = 8

つまり、ProcessPoolExecutor は、何も書かなければ使える論理CPUの数で動きます。今回いちばん速かった4並列と同じ数です。

店長「コンロが4口なら、茹で係は4人まで。5人目からは、横で見てるだけ」

茹で係の5人目「見てるだけじゃないです。ときどき交代してます」

……交代するたびに、鍋の前で場所を譲り合う時間が増える。

店長「じゃあ、いっそコンロを買い足せば?」

経理「ノートPCに、コンロは増設できません」

……厨房ごと買い替える話になってしまう。

雇いすぎた店員、どこで待ってるの?

OS が、4つの論理CPUを8つのプロセスで交代に使わせています。上級者向けの章で、交代でどれくらい1ファイルあたりの時間が伸びたかを見ます。

🔰 ここまで読めば今日から困らない

ここまでで、最低限の判断はできます。

  • 正規表現・集計・計算など、CPU を使う仕事 → ProcessPoolExecutor
  • ダウンロード・API・DB の応答待ちなど、待つ仕事 → ThreadPoolExecutor
  • 並列数は、まず省略(=プロセスなら os.process_cpu_count())。増やすときは計測してから
  • 速くなったかは、勘ではなく time.perf_counter() で測る。結果が直列と同じかも必ず確かめる

迷ったら、この形をコピーしてください。

from concurrent.futures import ProcessPoolExecutor

def analyze(path):
    ...  # 1ファイルぶんの集計をして、小さな結果を返す

if __name__ == "__main__":
    files = [...]
    with ProcessPoolExecutor() as ex:
        results = list(ex.map(analyze, files))

(形を示すためのひな形で、このままでは動きません。動く完成コードは中級者向けの章にあります)

店長「明日から、包丁を使う仕事は支店に、電話番は本店のスレッドさんに」

スレッドさん「電話番、がんばります」

……適材適所が、やっと始まった。

ProcessPoolExecutor を使えば、いつでも4倍速い?

なりません。今回は4並列で約1.5倍でした。 しかも、使い方を間違えると、直列より何倍も遅くなります。そこから先が、次の章からの話です。

ここから先は中級者向け。 Windows で起きるエラー、遅くなる使い方、完成コードを見ていきます。

Windows では if __name__ == "__main__": を書かないと止まる

Windows では、プロセスを使うスクリプトの本体を if __name__ == "__main__": の中に入れないと、RuntimeError で止まります。

わざと書かずに動かしてみます。

from concurrent.futures import ProcessPoolExecutor


def double(x):
    return x * 2


print("start")
with ProcessPoolExecutor(max_workers=2) as ex:
    print(list(ex.map(double, [1, 2, 3])))
start
start
start
RuntimeError: 
        An attempt has been made to start a new process before the
        current process has finished its bootstrapping phase.
(中略)
concurrent.futures.process.BrokenProcessPool: A process in the process pool was terminated abruptly while the future was running or pending.

(今回の検証環境の Windows 10・Python 3.13.1 で実行した出力から、標準出力と標準エラーの主要な行を抜き出したものです。実際はトレースバックが3つ出て、合計で150行ほどありました。終了コードは 1 でした)

start が 3回出ています。 1回目は本人(親プロセス)、2回目と3回目は、作ろうとした2つの子プロセスです。

multiprocessing のドキュメントによると、Windows と macOS の既定の起動方法は spawn です。spawn は、新しい Python を起動して、子プロセスの仕事に必要なものだけを引き継ぐ方式です。

そのとき子プロセスは、関数 double を見つけるために、スクリプトをもう一度頭から読み込みます。 ガードがないと、子プロセスもスクリプトの最後の with ProcessPoolExecutor(...) まで実行して、さらに子プロセスを作ろうとする。それを止めるために、Python が RuntimeError を出しているわけです。

ドキュメントのプログラミングガイドラインの「Safe importing of main module」の節にも、spawn や forkserver では、メインモジュールが新しいインタープリターから安全にインポートできるようにすること、if __name__ == '__main__': でプログラムの入り口を守ること、と書かれています。

支店「開店マニュアルを頭から読みます。……最後に『支店を2つ開くこと』と書いてあります」

本店「それは本店向けの指示!」

……支店が支店を開こうとして、チェーン展開が止まらない。

直し方は、スクリプトの本体をガードの中に入れるだけです。 関数の定義(def)は、ガードの外の上のほうに置いたままで大丈夫です。子プロセスが関数を見つけられるように、むしろ外に置きます。

もう1つ、同じ理由で起きるエラーがあります。

from concurrent.futures import ProcessPoolExecutor

if __name__ == "__main__":
    with ProcessPoolExecutor(max_workers=2) as ex:
        try:
            print(list(ex.map(lambda x: x * 2, [1, 2, 3])))
        except Exception as e:
            print(type(e).__name__, ":", e)
PicklingError : Can't pickle <function <lambda> at 0x000002403E803920>: attribute lookup <lambda> on __main__ failed

(今回の検証環境の Python 3.13.1 で実行した結果です)

プロセスに仕事を渡すとき、関数と引数は pickle(Python のオブジェクトをバイト列に詰める標準のしくみ)で送られます。 concurrent.futures のドキュメントにも、submit() と map() の関数と引数は pickle できる必要があり、REPL(対話モード)で定義した関数や lambda は動くと思わないように、と書かれています。同じページには、ProcessPoolExecutor は __main__ モジュールをワーカーがインポートできる必要があり、対話モードでは動かない、ともあります。

lambda「わたし、名前がないので、支店に住所を書いて送れないんです」

……匿名すぎて、宅配便が受け付けてくれない。

関数は def で名前を付けて、モジュールのいちばん外側に置きましょう。

支店で書いたメモは、本店に届かない

プロセスで動いた関数がグローバル変数を書き換えても、親プロセスの変数は変わりません。スレッドとの大きな違いです。

from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor

errors = 0


def count(n):
    global errors
    errors += n
    return errors


if __name__ == "__main__":
    with ThreadPoolExecutor(max_workers=4) as ex:
        list(ex.map(count, [1, 2, 3, 4]))
    print("thread : errors =", errors)

    errors = 0
    with ProcessPoolExecutor(max_workers=4) as ex:
        returned = list(ex.map(count, [1, 2, 3, 4]))
    print("process: errors =", errors, " returned =", returned)
thread : errors = 10
process: errors = 0  returned = [1, 3, 6, 10]

(今回の検証環境の Python 3.13.1 で実行した結果です)

スレッドは、同じプロセスの中で同じ変数を共有しているので、1 + 2 + 3 + 4 = 10 になりました。 プロセスは、子プロセスがそれぞれ自分の errors を持っているので、親の errors は 0 のままです。

ただし、ここで1つ「へえ」があります。 戻り値が [1, 3, 6, 10] と、足し算が積み上がっていますよね。今回の実行では、4つの仕事を1つの子プロセスが続けて受け持ったので、その子の中の errors がたまっていったんです。どの子プロセスがどの仕事を受け持つかは実行ごとに変わりうるので、子プロセスのグローバル変数に頼った集計は、親にも届かないうえに、結果も安定しません。

支店「本日のエラー、10件とメモしました!」

本店「こっちの帳簿は0件だけど」

支店「メモは支店の冷蔵庫に貼ってあります」

……報告は、紙を持って帰ってきてから。

正しい方法は、関数の戻り値で結果を返して、親で合計することです。1章の merge() がそれをしています。

スレッドなら共有できるから、スレッドのほうが便利では?

共有できるぶん、2つのスレッドが同時に同じ変数を書き換える心配が出てきます。スレッドでも、戻り値を返して親で合計する形にしておくと、あとでプロセスに切り替えるときも書き直さずにすみます。

1行ずつ渡すと、直列より何倍も遅くなる(chunksize)

プロセスに仕事を細かく渡しすぎると、仕事そのものより、渡す手間のほうが重くなります。

1ファイルの先頭5万行を、1行=1つの仕事としてプロセス4つに渡してみます。

import re
import time
from concurrent.futures import ProcessPoolExecutor

LINE = re.compile(r"^\S+ \S+ (\w+) (\S+) (\d{3}) (\d+)ms$")


def parse(line):
    m = LINE.match(line.rstrip("\n"))
    return m.group(3) if m else None


if __name__ == "__main__":
    with open("logs/access-01.log", encoding="utf-8") as f:
        lines = f.readlines()[:50_000]

    start = time.perf_counter()
    serial = [parse(x) for x in lines]
    print(f"serial          {time.perf_counter() - start:6.2f}s")

    for size in (1, 100, 5000):
        start = time.perf_counter()
        with ProcessPoolExecutor(max_workers=4) as ex:
            got = list(ex.map(parse, lines, chunksize=size))
        print(f"chunksize={size:<5} {time.perf_counter() - start:6.2f}s  same={got == serial}")
serial            0.12s
chunksize=1      28.44s  same=True
chunksize=100     0.60s  same=True
chunksize=5000    0.39s  same=True

(今回の検証環境の Python 3.13.1 で実行した結果です)

直列は 0.12秒。 プロセス4つで1行ずつ(chunksize=1)渡すと、28.44秒。直列の200倍以上かかりました。

店長「ネギ1本ずつ、支店に宅配便で送ってたの?」

本店「5万回、送り状を書きました」

……ネギより送り状のほうが高い。

concurrent.futures のドキュメントによると、ProcessPoolExecutor の map() は、入力をいくつかのかたまり(チャンク)に分けて、別々の仕事としてプールに送ります。その大きさを chunksize で指定でき、既定は 1。とても長い入力では、大きな chunksize にすると既定の 1 より性能が大きく上がることがある、と書かれています。ThreadPoolExecutor では chunksize は効きません。

chunksize=5000 にすると 0.39秒。それでも直列の 0.12秒には負けました。 この仕事は小さすぎて、プロセスの起動と受け渡しの手間を取り返せないんです。

だから、ログの集計では、1行ずつではなく1ファイルずつ渡すのがおすすめです。1章の bench.py がそうしています。ファイルが1つしかない巨大ログなら、行の範囲で分ける方法もありますが、今回は試していません。

じゃあ chunksize を100万にすれば最強?

チャンクの数がプロセスの数より少ないと、仕事をもらえない子プロセスが出ます。5万行に100万なら、チャンクは1つだけ。4人雇って1人で働く店になります。

結果を全部返すと遅い:支店では集計まで済ませる

子プロセスからの戻り値も pickle で送られるので、大きな結果を返すと、そのぶん遅くなります。

同じ8ファイルで、「パースした行を全部返す」と「数えた結果だけ返す」を比べました。

import glob
import re
import time
from collections import Counter
from concurrent.futures import ProcessPoolExecutor

LINE = re.compile(r"^\S+ \S+ (\w+) (\S+) (\d{3}) (\d+)ms$")


def rows(path):
    """Return every parsed row (big result)."""
    out = []
    with open(path, encoding="utf-8") as f:
        for line in f:
            m = LINE.match(line.rstrip("\n"))
            if m:
                out.append((m.group(2), m.group(3), int(m.group(4))))
    return out


def summary(path):
    """Return only the counts (small result)."""
    c = Counter()
    with open(path, encoding="utf-8") as f:
        for line in f:
            m = LINE.match(line.rstrip("\n"))
            if m:
                c[m.group(3)] += 1
    return c


if __name__ == "__main__":
    files = sorted(glob.glob("logs/access-*.log"))
    for func in (rows, summary):
        start = time.perf_counter()
        with ProcessPoolExecutor(max_workers=4) as ex:
            results = list(ex.map(func, files))
        if func is rows:
            n500 = sum(1 for r in results for _, s, _ in r if s == "500")
        else:
            n500 = sum(c["500"] for c in results)
        print(f"return {func.__name__:<8} {time.perf_counter() - start:5.2f}s  500={n500}")
return rows      3.92s  500=48274
return summary   2.31s  500=48274

(今回の検証環境の Python 3.13.1 で実行した結果です。rows の時間には、親で500を数える時間も含みます)

答えは同じ 48274件。 でも、160万行ぶんのタプルを親に送り返す rows は 3.92秒で、直列とほぼ同じになってしまいました。

支店「本日の伝票、160万枚お送りします!」

本店「件数だけでいいって言ったよね」

……段ボールが届くたびに、本店の床が見えなくなる。

店長「じゃあ伝票を送るのをやめて、支店に電話で読み上げてもらおう」

……160万枚を読み上げる電話、たぶん明日の開店に間に合いません。送り方を変えても、量が同じなら同じことです。

使用メモリは今回は測っていませんが、全部の行を親に集めれば、そのぶん親のメモリも使います。 子プロセスでは集計まで済ませて、小さな結果だけを返す。これが並列化で速くするための基本の形です。

完成コード:ファイルごとに並列で集計するスクリプト

ここまでの注意を全部入れた、コピーして使える形です。

  • 1ファイル=1つの仕事。子プロセスでは集計まで済ませて、小さな辞書だけ返す
  • 並列数の既定は、ファイル数・os.process_cpu_count()・61 のうち一番小さい数
  • 1ファイルが壊れていても、ほかのファイルの集計は続けて、最後に失敗を報告する
  • .gz の圧縮ログもそのまま読める
"""parallel_log_count.py - aggregate access logs file by file in worker processes.

usage: python parallel_log_count.py [--workers N] LOG [LOG ...]
"""
import argparse
import gzip
import os
import re
import sys
import time
from collections import Counter
from concurrent.futures import ProcessPoolExecutor, as_completed

LINE = re.compile(r"^\S+ \S+ (\w+) (\S+) (\d{3}) (\d+)ms$")


def analyze(path):
    """Runs in a worker. Returns only small summaries, never the raw lines."""
    status, slow, bad = Counter(), Counter(), 0
    total_ms = 0
    opener = gzip.open if path.endswith(".gz") else open
    with opener(path, "rt", encoding="utf-8", errors="replace") as f:
        for line in f:
            m = LINE.match(line.rstrip("\n"))
            if not m:
                bad += 1
                continue
            ms = int(m.group(4))
            status[m.group(3)] += 1
            total_ms += ms
            if ms >= 800:
                slow[m.group(2)] += 1
    return {"status": status, "slow": slow, "bad": bad, "total_ms": total_ms}


def default_workers(n_files):
    cpus = os.process_cpu_count() or 1  # Python 3.13+
    return max(1, min(n_files, cpus, 61))  # 61: Windows limit of ProcessPoolExecutor


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--workers", type=int, default=None)
    ap.add_argument("logs", nargs="+")
    args = ap.parse_args()
    workers = args.workers or default_workers(len(args.logs))

    start = time.perf_counter()
    total = {"status": Counter(), "slow": Counter(), "bad": 0, "total_ms": 0}
    failed = []
    with ProcessPoolExecutor(max_workers=workers) as ex:
        futures = {ex.submit(analyze, p): p for p in args.logs}
        for fut in as_completed(futures):
            path = futures[fut]
            try:
                r = fut.result()
            except Exception as e:  # one broken file must not hide the others
                failed.append(f"{path}: {type(e).__name__}: {e}")
                continue
            total["status"].update(r["status"])
            total["slow"].update(r["slow"])
            total["bad"] += r["bad"]
            total["total_ms"] += r["total_ms"]

    n = sum(total["status"].values())
    print(f"workers={workers} files={len(args.logs)} failed={len(failed)} "
          f"elapsed={time.perf_counter() - start:.2f}s")
    print(f"lines={n} bad_lines={total['bad']} avg_ms={total['total_ms'] / max(n, 1):.1f}")
    print("status:", dict(sorted(total["status"].items())))
    print("slow top3:", total["slow"].most_common(3))
    for msg in failed:
        print("FAILED", msg, file=sys.stderr)
    return 1 if failed else 0


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

8つのログと、わざと存在しないファイル logs/missing.log を渡して実行しました。

workers=4 files=9 failed=1 elapsed=2.51s
lines=1600000 bad_lines=0 avg_ms=451.6
status: {'200': 1439228, '302': 64190, '404': 48308, '500': 48274}
slow top3: [('/menu/tempura', 25851), ('/menu/kake', 25789), ('/menu/kitsune', 25699)]
FAILED logs/missing.log: FileNotFoundError: [Errno 2] No such file or directory: 'logs/missing.log'

(今回の検証環境の Python 3.13.1 で、python parallel_log_count.py logs/access-01.log …(8ファイルを並べて指定)… logs/missing.log として実行した結果です。終了コードは 1 でした)

比較のため、--workers 1(1プロセス)で8ファイルだけ渡した結果です。

workers=1 files=8 failed=0 elapsed=4.12s
lines=1600000 bad_lines=0 avg_ms=451.6
status: {'200': 1439228, '302': 64190, '404': 48308, '500': 48274}
slow top3: [('/menu/tempura', 25851), ('/menu/kake', 25789), ('/menu/kitsune', 25699)]

(同じ環境で実行した結果です。終了コードは 0 でした)

ポイントは3つです。

  • 集計結果が、並列数を変えても同じ。速さを比べる前に、ここを必ず確かめる
  • 存在しないファイルは FAILED として最後に報告し、終了コード 1。ほかの8ファイルの集計は止まらない
  • 引数は glob で展開していません。Windows の cmd や PowerShell は *.log を自動では展開しないので、必要ならスクリプト側で glob.glob() を使う形に変えてください(この変更版は今回は実行していません)

店長「1人で4.12秒、4人で2.51秒。うちの厨房なら、これが実力だね」

茹で係たち「4人とも、ちゃんと数え終わりました!」

……今度は、全員の伝票が本店に届いている。

as_completed って何?

終わった順に結果を受け取る関数です。 map() は渡した順に結果が返るので、最初のファイルが遅いと、ほかが終わっていても待たされます。どの順でもいい集計なら、as_completed で終わったものから合計していけます。

店長「早く茹で上がった順に出すってことね。うちの店と同じだ」

お客さん「先に頼んだ天ぷらうどんは?」

……集計なら順番は気にしなくていいけれど、お店でやるとクレームになります。

切り分けチェックリスト

症状から逆引きできるようにまとめました。

症状まず疑うこと確かめ方直し方
スレッドを増やしても速くならないCPU が中心の仕事で GIL の順番待ちスレッド1本と比べるProcessPoolExecutor にする
プロセスにしたら逆に遅い仕事が小さすぎる、1件ずつ渡している直列の時間と、1件あたりの時間を見る1ファイル単位にする、chunksize を大きく
プロセスでも思ったほど速くない戻り値が大きい、並列数が CPU 数より多い戻り値の大きさ、os.process_cpu_count()子で集計まで済ませる、並列数を減らす
RuntimeError: An attempt has been made to start a new process...if __name__ == "__main__": がないスクリプトの最後を見る本体をガードの中へ
PicklingErrorlambda や入れ子の関数を渡している渡している関数を見るdef でモジュールの外側に定義
集計結果が0、または毎回違う子プロセスでグローバル変数に足している戻り値で返しているか見る戻り値を親で合計
ValueError: max_workers must be <= 61Windows の上限max_workers の値61 以下にする

スレッドさん「わたしが原因の行、1つだけですね」

店長「1行目だけどね」

……容疑者リストの先頭に立つ店員。

表、長い。結局どれから見ればいい?

1行目と2行目です。 「仕事の種類に合った道具か」と「仕事の大きさは十分か」。この2つで、ほとんどのハマりは説明できます。

ここから先は上級者向け。読み飛ばしてもOKです。 4並列で4倍にならない理由、Windows とそれ以外の起動方法の違い、バージョンごとの変更を見ていきます。

4並列でも4倍にならない:1ファイルあたりの時間が伸びていた

1章で、プロセス4並列は約1.5倍でした。 今回の計測では、並列にすると、1ファイルあたりの処理時間そのものが伸びていました。

子プロセスの中で、1ファイルの処理時間を測ってみます。

import glob
import os
import time
from concurrent.futures import ProcessPoolExecutor

from bench import analyze


def timed(path):
    start = time.perf_counter()
    analyze(path)
    return os.getpid(), time.perf_counter() - start


if __name__ == "__main__":
    files = sorted(glob.glob("logs/access-*.log"))
    print("serial  per file:", " ".join(f"{timed(f)[1]:.2f}" for f in files))
    for n in (1, 2, 4, 8):
        start = time.perf_counter()
        with ProcessPoolExecutor(max_workers=n) as ex:
            res = list(ex.map(timed, files))
        wall = time.perf_counter() - start
        pids = len({p for p, _ in res})
        print(f"process x{n} wall={wall:.2f}s workers_used={pids} per file:",
              " ".join(f"{t:.2f}" for _, t in res))
serial  per file: 0.45 0.52 0.46 0.43 0.49 0.47 0.46 0.42
process x1 wall=3.76s workers_used=1 per file: 0.46 0.43 0.47 0.45 0.42 0.42 0.46 0.45
process x2 wall=3.84s workers_used=2 per file: 0.88 0.89 0.94 0.91 0.83 0.83 0.96 0.95
process x4 wall=2.33s workers_used=4 per file: 0.99 1.01 1.03 1.00 0.99 0.99 1.00 0.99
process x8 wall=2.87s workers_used=8 per file: 2.15 2.12 2.17 2.16 2.13 2.11 2.08 1.83

(今回の検証環境の Python 3.13.1 で実行した結果です。from bench import analyze は、1章のスクリプトを bench.py として同じフォルダーに置いて読み込んでいます)

1ファイルの処理は、1人なら約0.45秒。 2並列では約0.9秒、4並列では約1.0秒、8並列では約2.1秒に伸びています。

8並列は分かりやすくて、4つの論理CPUを8つのプロセスで交代に使うので、1ファイルにかかる時間がほぼ倍になりました。 2並列で1ファイルが倍の時間になったのは、今回の計測だけでは原因を特定できていません。今回の CPU は2コア4スレッドなので、2つのプロセスが同じコアの2つの論理CPUに割り当てられた可能性や、同時に動くコアが増えると CPU のクロックが下がる可能性などが考えられますが、確かめてはいません。

ここで言えるのは、論理CPUの数だけ並列にしても、その数の倍率で速くなるとは限らないということです。だから、並列数は計算で決めず、自分の環境で測って決めます。

店長「4口のコンロに4人並べたら、1人あたりの茹で時間が倍になった」

茹で係「火力が、4つの鍋で分け合いになってまして……」

店長「……それ、ガスの元栓の話?」

……厨房の比喩が、ガス会社の契約にまで広がってきた。

もう1つ、プロセスには起動の手間があります。

import time
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor


def noop(x):
    return x


if __name__ == "__main__":
    for label, make in [("thread x4", ThreadPoolExecutor), ("process x4", ProcessPoolExecutor)]:
        start = time.perf_counter()
        with make(max_workers=4) as ex:
            list(ex.map(noop, range(4)))
        print(f"{label:<11} start+4 empty tasks+close: {time.perf_counter() - start:.3f}s")
thread x4   start+4 empty tasks+close: 0.002s
process x4  start+4 empty tasks+close: 0.363s

(今回の検証環境の Python 3.13.1 で実行した結果です)

何もしない仕事でも、プロセス4つの起動と終了で約0.36秒。スレッドは0.002秒でした。 spawn では、子プロセスごとに新しい Python を起動してスクリプトを読み込むので、このくらいの手間がかかります。合計で数秒以下の仕事なら、並列にしないほうが速いこともあるわけです。

店長「4支店ぶんの開店と閉店で、約0.36秒か」

支店「看板を作って、のれんを掛けて、マニュアルを頭から読んでいます」

……1秒で終わる仕事に支店を4つ開くと、開店準備で一日が終わる。

起動方法とバージョンの違い:fork・spawn・forkserver、3.13 と 3.14

Windows と macOS は spawn が既定ですが、Linux などの POSIX は Python のバージョンで既定が変わっています。同じスクリプトでも、OS とバージョンで子プロセスの作られ方が違うんです。

multiprocessing のドキュメントの起動方法の説明をまとめるとこうです。

起動方法どう作るか既定になる環境
spawn新しい Python を起動して、必要なものだけ引き継ぐWindows と macOS
forkos.fork() で親プロセスを丸ごと複製する(POSIX のみ)3.14 からは、どの環境でも既定ではない
forkserver先に作ったサーバー役のプロセスから複製する3.14 から POSIX(Linux など、Unix パイプでファイル記述子を渡せる環境)の既定

中級者向けの章で見た「Safe importing of main module」の注意は、ドキュメントでは spawn と forkserver の節に書かれています。fork については、マルチスレッドのプロセスを安全に fork するのは問題がある、と書かれています。 3.14 からは、fork が必要なコードは get_context() や set_start_method() で明示的に指定する必要があります。concurrent.futures のドキュメントにも、3.14 で ProcessPoolExecutor の既定の起動方法が fork から変わったので、fork が必要なら mp_context=multiprocessing.get_context("fork") を渡すように、とあります。

これが実務で効いてくるのは、こんな場面です。

  • Linux で fork が既定だったころ(3.13 以前)に書いたスクリプトは、ガードの書き忘れなど、spawn では問題になる書き方が見過ごされている可能性がある
  • 同じスクリプトを Windows で動かす、または Python 3.14 に上げると、起動方法が spawn や forkserver に変わる。動きが変わらないか、移行の前に確かめる

今回の検証は Windows(spawn)だけで、Linux の fork・forkserver の動きは実行していません。ここは公式ドキュメントで確認した範囲の話です。

fork「親の記憶ごと、そっくりコピーで出勤します」

spawn「わたしは毎回、新人研修からです」

……どちらが楽かは、研修の資料がちゃんと書けているかで決まる。

ほかにも、バージョンや OS で違う点があります(どれも concurrent.futures のドキュメントと os のドキュメントで確認したものです)。

  • Python 3.13 から、ProcessPoolExecutor の既定の並列数は os.cpu_count() ではなく os.process_cpu_count()。CPU アフィニティで使える CPU が絞られている環境では、少なくなる
  • 3.13 から、-X cpu_count オプションや PYTHON_CPU_COUNT 環境変数で、os.cpu_count() と os.process_cpu_count() の値を上書きできる(今回は試していません)
  • Windows では ProcessPoolExecutor の max_workers は 61 以下。省略時も、CPU がもっと多くても 61 までになる

61 の上限は、今回の環境でも確かめました。

from concurrent.futures import ProcessPoolExecutor

if __name__ == "__main__":
    try:
        ProcessPoolExecutor(max_workers=62)
    except ValueError as e:
        print("ValueError:", e)
ValueError: max_workers must be <= 61

(今回の検証環境の Windows 10・Python 3.13.1 で実行した結果です)

店長「62人目の店員は?」

Windows「定員オーバーです」

……2コアのノートPCで62人雇う予定は、そもそもなかった。

最後に、GIL そのものの話です。 Python 用語集には、Python 3.13 から、--disable-gil を付けてビルドした Python(free-threaded build)では GIL を無効にでき、その場合は -X gil=0 か環境変数 PYTHON_GIL=0 を指定して実行する、と書かれています。マルチスレッドのアプリケーションの性能が上がり、マルチコアの CPU を効率よく使いやすくなる、とも説明されています。

今回の検証に使った Python 3.13.1 は、sys._is_gil_enabled() が True を返す、GIL ありのビルドでした。free-threaded build での計測はしていないので、今回のスレッドの結果がどう変わるかは書きません。

スレッドさん「包丁、1人1本もらえる日が来るんですか?」

店長「そういう厨房も、別の建物にはあるらしい」

……引っ越しの前に、新しい厨房の使い勝手を確かめてから。

まとめ:店員の数より、仕事の種類と大きさ

冒頭では、店員を8人に増やしたのに、ログの集計が速くなりませんでした。

  • スレッドは、GIL があるので、Python のコードを動かすのは同時に1本だけ。CPU が中心の集計では速くならない
  • 待ち時間が中心の仕事なら、スレッドで大きく速くなる(今回は約4秒 → 約0.5秒)
  • CPU が中心なら ProcessPoolExecutor。並列数は、まず os.process_cpu_count() まで(今回は4並列で約1.5倍、8並列では遅くなった)
  • Windows では if __name__ == "__main__": が必須。関数は def でモジュールの外側に
  • 子プロセスのグローバル変数は親に届かない。戻り値で返して親で合計する
  • 1行ずつ渡さない(chunksize)、全部の行を返さない。子で集計まで済ませる

スレッドさん「明日から、わたしは出前の電話番ですね」

店長「そう。包丁は支店のみんなに任せて」

茹で係たち「4人で、ちょうどいいです!」

……やっと、人数と仕事が噛み合った。

じゃあ、並列化はもうしなくていい?

してください。ただし、測ってから。 今回の厨房では、4人でちょうどよかった。でも、あなたの厨房のコンロの数は、あなたの環境で測るまで分かりません。

並列化は、仕事の種類で道具を選び、仕事を大きめに切り、結果を小さくして返す。それから並列数を測って決める。それだけで、店員を増やしたのに遅くなる、という悲しい夜はかなり減ります。

まずは、自分の集計スクリプトを直列のまま time.perf_counter() で測ってみてください。 その数字が、支店を開くかどうかを決める最初の材料になります🍜

参考資料

  • Python:concurrent.futures — ThreadPoolExecutor と ProcessPoolExecutor の max_workers の既定値(3.8・3.13 の変更)、Windows の max_workers は 61 以下、main モジュールがインポートできる必要があること、関数と引数は pickle できる必要があり lambda は動かないこと、map の chunksize、3.14 で既定の起動方法が fork から変わったこと
  • Python:multiprocessing — サブプロセスで GIL を事実上回避すること、spawn・fork・forkserver の説明と既定(Windows と macOS は spawn、3.14 から POSIX は forkserver)、Safe importing of main module(if name == ‘main‘ で入り口を守る)
  • Python 用語集:global interpreter lock — 一度に1つのスレッドだけがバイトコードを実行すること、I/O の間は GIL が手放されること、3.13 の free-threaded build と -X gil=0
  • Python:os.cpu_count と os.process_cpu_count — システムの論理CPU数と、プロセスが使える論理CPU数の違い(3.13 で追加)、-X cpu_count と PYTHON_CPU_COUNT による上書き

確認日:2026-10-09。

よかったらシェアしてね!
  • URLをコピーしました!
  • URLをコピーしました!

この記事を書いた人

おうどん|癒しと創作をたのしむ雑食クリエイター
写経アプリやクレイセラピー、優しい和風デザインがすき。
ZARDと刀剣と文字に癒されて、今はアプリ作ってます。

コメント

コメントする

目次