結論:maxsize付きQueueで受け渡し、全件にtask_doneを対応させる
入力を作る速度が処理速度より速いと、未処理データをリストへため続けるだけではメモリが増える。スレッド間の受け渡しにはqueue.Queue(maxsize=...)を使うと、待ち行列が満杯の間はput側を待たせられる。消費側ではgetした各項目にtask_doneを一度だけ対応させ、処理済みであることを通知する。
上限があるのは待ち行列内の項目数であり、項目そのもののバイト数ではない。大きな画像2枚と小さな整数2個では必要なメモリが違うし、取り出して処理中の項目は行列の外へ出ている。ワーカー数や入力の大きさも含めて、許容する保持量を考えたい。maxsizeだけでプロセス全体の使用量を厳密に決めることはできない。
そのまま動かせる例
例は容量2のキューへ整数5個を投入し、1つの消費スレッドが処理する。入力3は意図的なValueErrorとし、失敗を記録したうえで残りの処理を続ける。最後に専用のSTOPオブジェクトを入れ、それを受け取ったら終了する。標準ライブラリだけで動き、外部通信やファイル操作はしない。
putとgetには2秒、スレッドの終了待ちには4秒の期限を設定している。正常なSTOPを確認し、スレッドが止まった後にqueue.joinで未処理数が0になったことを待ち合わせる構成である。queue.join自体にはtimeout引数がないため、失敗したワーカーを放置したまま無条件に呼ぶ例にはしていない。仕事は5件の短い処理に限定した。
from queue import Queue, Empty
from threading import Thread
queue = Queue(maxsize=2)
STOP = object()
results, errors = {}, {}
normal_stop = []
def consumer():
while True:
try:
item = queue.get(timeout=2)
except Empty:
errors["worker"] = "idle timeout"
return
try:
if item is STOP:
normal_stop.append(True)
return
if item == 3:
raise ValueError("invalid item")
results[item] = item * 2
except ValueError as exc:
errors[item] = str(exc)
finally:
queue.task_done()
thread = Thread(target=consumer, name="consumer")
thread.start()
try:
for value in range(5):
queue.put(value, timeout=2)
queue.put(STOP, timeout=2)
finally:
thread.join(timeout=4)
assert not thread.is_alive()
assert normal_stop == [True]
queue.join() # The worker handled every item including STOP.
assert results == {0: 0, 1: 2, 2: 4, 4: 8}
assert errors == {3: "invalid item"}
print("queue capacity:", queue.maxsize)
print("results:", sorted(results.items()))
print("errors:", sorted(errors.items()))
print("worker stopped and queue drained: True")
実行結果
queue capacity: 2
results: [(0, 0), (1, 2), (2, 4), (4, 8)]
errors: [(3, 'invalid item')]
worker stopped and queue drained: True
task_doneは取得した回数と対にする
getでキューから項目がなくなっても、その処理が完了したとは限らない。Queueは未完了数を別に管理し、task_doneが呼ばれるとその数を減らす。joinは未完了数が0になるまで待つ。例では処理中にValueErrorが起きてもfinallyでtask_doneを呼ぶので、失敗項目だけ未完了として残り続けることがない。
STOPもputされた一項目なのでtask_doneが必要である。一方、getがEmptyで失敗したときには取得していないため、task_doneを呼ばない。回数が不足するとjoinが戻らず、多すぎればエラーとなる。「処理が成功した回数」ではなく「取得した項目の処理を完了として扱った回数」と対応させると整理しやすい。
終了用の印と失敗の方針
通常のデータと区別できる専用オブジェクトを使い、isで判定する。Noneを有効なデータとして扱う処理なら、Noneを終了印にしない方が混同しにくい。複数の消費スレッドを起動する場合は、それぞれが終了印を受け取れるように必要な数を投入するなど、全ワーカーの停止手順を設計する。この例は一つの消費側だけに限定している。
失敗を記録して継続するか、全体を中断するかは用途による。ここではValueErrorだけを想定した入力エラーとして保存している。どんな例外も無視して進めると、処理が壊れているのに全件完了に見える。中断方式なら、生産側へ失敗を伝える方法、待ち項目をどう扱うか、終了待ちの期限なども必要になる。task_doneは仕事の成功を保証する通知ではない。
待ち時間と観測値の注意
qsizeやemptyは、その瞬間の目安であり、次のgetやputが必ず待たずに成功する保証にはならない。別のスレッドが直後に状態を変え得るため、空かどうかを先に見てから判断するより、getのtimeoutや例外で扱う方が安全である。この例もEmptyを直接処理し、入力が来ないまま無期限に待ち続けないようにしている。
Pythonの新しい版ではQueue.shutdownなど追加された停止機能もあるが、この例はCPython 3.12で使える終了印の方式である。公式資料の最新版にある機能を手元でも使えると判断せず、対応版を確かめよう。キューを小さくして例外を一件混ぜ、全結果と失敗、未処理数、スレッドの終了を一緒に確認すると、実データへ広げる前に停止漏れを見つけやすい。
確認環境と参考資料
例はLinux・CPython 3.12.14で実行した。掲載した出力はこの環境での結果である。公式資料のstable版や最新版は更新されるため、手元のバージョンと対応する仕様も確認してほしい。
- Python公式:queue(2026年10月2日参照)
関連項目:ThreadPoolExecutorで結果と失敗を取りこぼさず回収する / 共有変数の更新がずれる原因は?Lockで競合を防ぐ
