ThreadPoolによる非同期タスク処理
Pythonの`multiprocessing.pool`モジュールを利用することで、スレッド管理を簡素化し、I/Oバウンドな処理の効率を向上させることができます。以下のコードは、スレッドプールを生成し、複数のタスクを非同期に実行する例です。
import time
from multiprocessing.pool import ThreadPool
# タスク処理用の関数定義
def async_task(task_id):
print(f"タスク {task_id} を開始します...")
time.sleep(2) # I/O待機のシミュレーション
print(f"タスク {task_id} が完了しました。")
return task_id * 100
start_time = time.time()
# スレッドプールのサイズを3に設定して生成
pool = ThreadPool(3)
tasks = []
# 5つのタスクを非同期に投入
for i in range(5):
# apply_asyncを使用すると、タスクの完了を待たずに次の処理へ進む
result = pool.apply_async(async_task, args=(i,))
tasks.append(result)
# これ以上タスクを受け付けないようにプールを閉じる
pool.close()
# 全てのタスクが完了するのを待つ
pool.join()
# 結果の確認(get()で戻り値を取得)
for t in tasks:
print(f"戻り値: {t.get()}")
print(f"総経過時間: {time.time() - start_time:.2f}秒")
この実行結果において、スレッドプールのサイズが3であるため、3つのタスクが並列に処理されます。それぞれのタスクに2秒の待機時間があるため、5つのタスクを完了させるには理論上約4秒(2秒 + 2秒)かかります。Pythonのグローバルインタプリタロック(GIL)によりCPUバウンドな計算処理には限界がありますが、ファイルアクセスやネットワーク通信のようなI/O待機が発生する処理では、スレッドによる並列化が有効です。
ProcessPoolを使った並列処理
次に、CPUバウンドな処理やメモリ空間を分離したい場合に適したプロセスプール(`multiprocessing.Pool`)の実装例を示します。ここでは、ブロッキング動作を行う`map`メソッドと、非ブロッキング動作を行う`map_async`メソッドの違いを確認します。
import time
from multiprocessing import Pool
def heavy_process(data):
print(f"プロセスがデータ {data} を処理中...")
time.sleep(1)
return data ** 2
if __name__ == '__main__':
proc_start = time.time()
# プロセス数を2に設定
with Pool(2) as process_pool:
# 1. mapメソッド(同期的・ブロッキング)
print("--- mapによる実行(ブロッキング) ---")
# 全ての処理が終わるまでメインプロセスは待機する
results_map = process_pool.map(heavy_process, [1, 2, 3, 4])
print(f"mapの結果: {results_map}")
# 2. map_asyncメソッド(非同期的・ノンブロッキング)
print("\n--- map_asyncによる実行(非ブロッキング) ---")
# 結果オブジェクトが即座に返される
result_obj = process_pool.map_async(heavy_process, [5, 6, 7, 8])
print("コードはここでブロックされずに進行します")
# 必要に応じて後で結果を待つ
result_obj.wait()
print(f"map_asyncの結果: {result_obj.get()}")
print(f"プロセス処理の総経過時間: {time.time() - proc_start:.2f}秒")
`map`および`apply`メソッドはタスクが完了するまで呼び出し元をブロックしますが、末尾に`_async`がつくメソッドは即座に戻り値を返し、バックグラウンドで処理を継続します。これにより、メインプロセスのフローを止めずに他の処理を行うことが可能になります。
スレッドプールを用いた並行サーバーの実装
最後に、上記のスレッドプールを応用して、複数のクライアントからの接続を同時に処理するTCPサーバーを構築します。`ThreadPool`の代わりに`Pool`(プロセスプール)を使用することも可能であり、API仕様が共通であるため、モジュールのインポートとクラス名を変更するだけで切り替えが可能です。
import socket
from multiprocessing.pool import ThreadPool
# クライアントとの通信を処理するワーカー関数
def client_handler(client_sock, client_addr):
print(f"新しい接続: {client_addr}")
try:
while True:
# データの受信
raw_data = client_sock.recv(1024)
if not raw_data:
break
decoded_msg = raw_data.decode('utf-8')
print(f"受信メッセージ [{client_addr}]: {decoded_msg}")
# エコーバック(受信したデータをそのまま送信)
client_sock.sendall(raw_data)
except Exception as e:
print(f"エラー発生: {e}")
finally:
client_sock.close()
print(f"接続終了: {client_addr}")
def start_server():
# ソケットの設定
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server_socket.bind(('0.0.0.0', 8888))
server_socket.listen(5)
print("サーバー起動中... ポート8888で待機")
# スレッドプールの生成(最大4つの同時接続を処理)
pool = ThreadPool(4)
try:
while True:
# クライアントの接続待ち
conn, addr = server_socket.accept()
# 接続処理をスレッドプールに非同期で委譲
pool.apply_async(client_handler, args=(conn, addr))
except KeyboardInterrupt:
print("サーバーを停止します")
finally:
server_socket.close()
if __name__ == '__main__':
start_server()
このサーバー実装では、メインスレッドがループ内で`accept()`を行い、新しい接続が確立するたびに`pool.apply_async`を通じてワーカー関数に処理を委譲します。これにより、特定のクライアントとの通信でブロックが発生しても、他のクライアントの接続受け付けに影響を与えない、堅牢な並行処理サーバーが実現できます。