コルーチンの基本概念
コルーチンは、マイクロスレッドやファイバーとも呼ばれ、ユーザー空間で動作する軽量なスレッドのようなものです。従来のサブルーチン(関数)は、AからBを呼び出し、BからCを呼び出すといった階層的な呼び出し構造を持ち、必ず決まった順序で戻り値を返します。一方コルーチンは、実行途中で一時停止し、別のルーチンに制御を移した後、再開時に元の位置から処理を継続することができます。
メリット
- 高い実行効率: スレッドのコンテキストスイッチによるオーバーヘッドが発生しません。プログラム自身で制御を移行するため、スレッド数が増えるほどスレッド切り替えコストの削減効果が顕著になります。
- ロックが不要: 単一スレッド内で動作するため、複数のスレッドから同時に変数が書き換えられる競合状態が発生しません。共有リソースへのアクセスにミューテックスなどのロック機構が不要となり、オーバーヘッドを抑えられます。
デメリット
- マルチコアの活用不可: 基本的に1つのスレッドで実行されるため、CPUの複数コアを同時に使い切ることはできません。
- ブロッキングの影響: 同期IOなどのブロッキング操作が含まれると、その間プログラム全体が停止してしまいます。
同期処理と非同期処理
PythonはGIL(グローバル・インタプリタ・ロック)の存在によりCPUバウンドな処理では性能上限に直面しがちですが、ネットワーク通信など待機時間が長いIOバウンドな処理においては、非同期処理を用いることで劇的なパフォーマンス向上が見込めます。
- 同期処理: 順次実行方式。前の処理が完了するまで次の処理は待機状態となります。
- 非同期処理: 実行中の処理の完了を待たずに、すぐさま次の処理へ移行します。結果はコールバックや状態監視を通じて後から受け取ります。
同期処理のコード例
import time
def request_data(server):
print(f"Fetching from {server}...")
time.sleep(2)
print(f"Completed fetching from {server}")
if __name__ == "__main__":
start = time.time()
for s in ["Server-A", "Server-B", "Server-C", "Server-D"]:
request_data(s)
print(f"Total time: {time.time() - start:.2f}s")
非同期処理のコード例
import time
import asyncio
async def request_data_async(server):
print(f"Fetching from {server}...")
await asyncio.sleep(2)
print(f"Completed fetching from {server}")
async def main():
start = time.time()
tasks = [request_data_async(s) for s in ["Server-A", "Server-B", "Server-C", "Server-D"]]
await asyncio.gather(*tasks)
print(f"Total time: {time.time() - start:.2f}s")
if __name__ == "__main__":
asyncio.run(main())
asyncioモジュールの活用
Python 3.4以降、標準ライブラリとして導入されたasyncioは、非同期IOをネイティブにサポートします。イベントループを中心としたプログラミングモデルを採用しており、コルーチンをイベントループに登録することで非同期実行を実現します。
主要キーワード
- event_loop: イベントループ。登録されたコルーチンを監視し、条件が満たされたタイミングで該当する処理を呼び出します。
- coroutine:
asyncキーワードで定義された非同期関数。呼び出すだけで実行されるのではなく、コルーチンオブジェクトを返します。 - task: コルーチンをラップして実行状態を管理するオブジェクト。
- async/await:
asyncでコルーチンを定義し、awaitで非同期処理の完了を待機(制御を一時放棄)します。
基本的な実行方法
Python 3.7以降では、asyncio.run()を用いることで簡潔にイベントループを起動できます。
import asyncio
async def delay_execution(seconds):
print(f"Waiting for {seconds}s...")
await asyncio.sleep(seconds)
print("Done")
asyncio.run(delay_execution(2))
タスクの生成と並行実行
複数の非同期処理を並行して実行するには、asyncio.create_task()を使用してタスクを生成し、イベントループにスケジュールします。
import asyncio
async def delay_execution(seconds):
print(f"Waiting for {seconds}s...")
await asyncio.sleep(seconds)
print("Done")
async def main():
task = asyncio.create_task(delay_execution(2))
await task
asyncio.run(main())
コールバックによる結果取得
タスク完了時に特定の処理を行いたい場合、add_done_callbackでコールバック関数を設定できます。
import asyncio
async def fetch_from_db(query):
print(f"Executing query: {query}")
await asyncio.sleep(3)
return f"Result of {query}"
def on_task_complete(task):
print(f"Callback received: {task.result()}")
async def main():
task = asyncio.create_task(fetch_from_db("SELECT * FROM users"))
task.add_done_callback(on_task_complete)
await task
asyncio.run(main())
複数タスクの待機
複数の対象に対して同時にリクエストを送り、全ての完了を待機する例です。asyncio.gatherを利用することで、待機時間を大幅に短縮できます。
import asyncio
import time
async def fetch_resource(endpoint):
print(f"Requesting {endpoint}...")
await asyncio.sleep(2)
return f"Data from {endpoint}"
async def main():
endpoints = ["API_1", "API_2", "API_3", "API_4"]
start = time.time()
tasks = [asyncio.create_task(fetch_resource(ep)) for ep in endpoints]
results = await asyncio.gather(*tasks)
for res in results:
print(res)
print(f"Total elapsed time: {time.time() - start:.2f}s")
asyncio.run(main())