プログラム 技術

Pythonで非同期処理を実装する

これまでC++、Kotlinといった言語で非同期処理について記載してきましたが、今回はPythonでの非同期処理についてとなります
今回も特にライブラリを追加することなく実装しています

名称バージョン
Python3.12
関連記事
Kotlinで非同期処理を実装する - ナストンのまとめ
Kotlinで非同期処理を実装する

今回はKotlinでの非同期処理についてになります。非同期処理は【kotlinx-coroutines-core】モジュ ...

関連記事へ

関連記事
C++で非同期処理を実装する - ナストンのまとめ
C++で非同期処理を実装する

以前にKotlinで非同期処理について記載しましたが今回はC++で非同期処理について記載したいと思います 名称バージョン ...

関連記事へ

asyncio (async/await)を使用するパターン

  • 用途 : 大量のI/O待ちが発生する処理を効率的にさばきたい
  • 利点
    • シングルスレッドで動作するため、スレッド/プロセスの生成コストがゼロ
    • ロックやデッドロックの心配が不要(共有メモリの競合が起きない)
    • 数千〜数万の同時接続を少ないメモリで処理可能
  • 適した場面
    • Webスクレイピング(数百〜数千のURLを同時取得)
    • APIサーバー
    • チャットアプリ、リアルタイム通知
    • 大量のファイルダウンロード/アップロード
 # --- 非同期関数の定義 ---
async def async_io_task(task_name: str, duration: float) -> str:
    """非同期I/Oタスク: asyncio.sleep で待機"""
    print(f"  🔄 {task_name}: 開始 (待機 {duration}秒)")
    await asyncio.sleep(duration)  # ← ここで制御を他タスクへ渡す
    print(f"  ✅ {task_name}: 完了")
    return f"{task_name} の結果"


# --- 逐次実行 vs 並行実行の比較 ---
async def sequential():
    """逐次実行: 一つずつ順番に処理"""
    print("\n--- 逐次実行 ---")
    start = time.perf_counter()
    result1 = await async_io_task("タスクA", 1.0)
    result2 = await async_io_task("タスクB", 1.5)
    result3 = await async_io_task("タスクC", 0.5)
    elapsed = time.perf_counter() - start
    print(f"  結果: {[result1, result2, result3]}")
    print(f"  ⏱  逐次実行時間: {elapsed:.3f} 秒 (≒ 3.0秒)")


async def concurrent_run():
    """並行実行: gather で同時に処理"""
    print("\n--- 並行実行 (asyncio.gather) ---")
    start = time.perf_counter()
    results = await asyncio.gather(
        async_io_task("タスクA", 1.0),
        async_io_task("タスクB", 1.5),
        async_io_task("タスクC", 0.5),
    )
    elapsed = time.perf_counter() - start
    print(f"  結果: {results}")
    print(f"  ⏱  並行実行時間: {elapsed:.3f} 秒 (≒ 1.5秒)")


async def main_async():
    await sequential()
    await concurrent_run()


# --- 実行 ---
asyncio.run(main_async())

実行結果は以下の通りになります

threadingを使用するパターン

  • 用途 : I/Oバウンドで既存の同期ライブラリを使いたい
  • 利点
    • 待ちの間に別スレッドが動けるため、待ち時間を有効活用
    • asynco未対応の同期ライブラリをそのまま使用出来る
    • スレッド間でメモリ空間を共有するため、データのやり取りが容易
  • 適した場面
    • 少数〜中程度の並行タスク(数十スレッド程度)
    • 既存の同期ライブラリ(requests等)を並行利用したい場合
    • ファイルの読み書きやネットワーク通信の並行化
    • guiアプリのバックグラウンド
threads = []
for name, dur in [("タスクA", 1.0), ("タスクB", 1.5), ("タスクC", 0.5)]:
    t = threading.Thread(target=worker, args=(name, dur))
    threads.append(t)
    t.start()  # スレッド開始

for t in threads:
    t.join()  # 全スレッドの完了を待つ

elapsed = time.perf_counter() - start
print(f"  結果: {results}")
print(f"  ⏱  並行実行時間: {elapsed:.3f} 秒 (≒ 1.5秒)")

# --- Lock による排他制御 ---
print("\n--- Lock による排他制御 ---")
counter = 0
lock = threading.Lock()
def increment_without_lock():
    nonlocal counter
    for _ in range(100_000):
        counter += 1  # ← 競合が発生する可能性あり
def increment_with_lock():
    nonlocal counter
    for _ in range(100_000):
        with lock:  # ← ロックで保護
            counter += 1

# ロックなし
counter = 0
threads = [threading.Thread(target=increment_without_lock) for _ in range(5)]
for t in threads:
    t.start()
for t in threads:
    t.join()
print(f"  ロックなし: counter = {counter:>10,} (期待値: 500,000)")

# ロックあり
counter = 0
threads = [threading.Thread(target=increment_with_lock) for _ in range(5)]
for t in threads:
    t.start()
for t in threads:
    t.join()
print(f"  ロックあり: counter = {counter:>10,} (期待値: 500,000)")

multiprocessingを使用するパターン

  • 用途 : CPUバウンドで全コアを使いたい
  • 利点
    • プロセスごとに独立したGILを持つため、CPUバウンド処マルチコアで並実行される
    • プロセスが分離されているため、1つのプロセスがクラッシュしても他に影響しない
    • cpuコア数に比例した高速化が期待できる
  • 適した場面
    • 画像処理
    • 動画エンコード
    • 大規模な数値計算
    • 科学計算
    • 機械学習のデータ前処理
target_n = 5000  # N番目の素数を求める

# --- 逐次実行 ---
print(f"\n--- 逐次実行 (素数計算 × 4回, 各{target_n}番目) ---")
start = time.perf_counter()
sequential_results = [heavy_cpu_task(target_n) for _ in range(4)]
elapsed_seq = time.perf_counter() - start
print(f"  結果: {sequential_results}")
print(f"  ⏱  逐次実行時間: {elapsed_seq:.3f} 秒")

# --- Pool.map による並列実行 ---
print(f"\n--- 並列実行 (multiprocessing.Pool) ---")
start = time.perf_counter()
with multiprocessing.Pool(processes=4) as pool:
    parallel_results = pool.map(heavy_cpu_task, [target_n] * 4)
elapsed_par = time.perf_counter() - start
print(f"  結果: {parallel_results}")
print(f"  ⏱  並列実行時間: {elapsed_par:.3f} 秒")

if elapsed_seq > 0:
    print(f"  📈 速度向上: {elapsed_seq / elapsed_par:.1f}倍")

concurrent.futuresを使用するパターン

  • 用途 : threading / multiprocessing を統一的なAPIで簡潔に記述したい
  • 利点
    • ThreadPoolExecutr(I/Oバウンド)とProcessPoolExecutr(CPUバウンド同じインターフェイス)で使える
    • Futueオブジェクトで非同期に結果を取得して、as_completdで完了順に処理できる
    • wih文でリソース管理が自動化される
    • mpで一括処理が簡単に書ける
  • 適した場面
    • threading/multiprocessing を簡潔に書きたい場合全般
    • 複数タスクの結果を完了順に処理したい場合
    • プロトタイピングやスレッド
print("\n--- ThreadPoolExecutor (I/Oバウンド) ---")
start = time.perf_counter()

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    # submit + as_completed パターン: 完了順に結果を取得
    futures = {
        executor.submit(simulate_io, f"タスク{name}", dur): name
        for name, dur in [("A", 1.0), ("B", 1.5), ("C", 0.5)]
    }
    for future in concurrent.futures.as_completed(futures):
        name = futures[future]
        result = future.result()
        print(f"  📦 {name} 完了: {result}")

elapsed = time.perf_counter() - start
print(f"  ⏱  実行時間: {elapsed:.3f} 秒")

# --- ProcessPoolExecutor (CPUバウンド) ---
print("\n--- ProcessPoolExecutor (CPUバウンド) ---")
target_n = 5000

# map パターン: 入力リストに対して一括処理
start = time.perf_counter()
with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(heavy_cpu_task, [target_n] * 4))
elapsed = time.perf_counter() - start
print(f"  結果: {results}")
print(f"  ⏱  実行時間: {elapsed:.3f} 秒")

callbackを使用するパターン

  • 用途 : 処理の完了時に通知を受け取り、後続の処理をイベント駆動で実行したい
  • 利点
    • 処理の完了を待たずに次の処理へ進められる(ノンブロッキング)
    • concurrent.futursのadd_done_callbak`で手軽に使える
    • イベント駆動型アーキテクチャに自然にフィットする
  • 適した場面
    • 処理完了通知(ファイルのダウンロード完了後にUIを更新など)
    • guiアプリのイベントハンドリング
    • ログ記録や監視など、メインフローを止めたくない後処理
    • 既存のイベント駆動フレームワークとの連携
# --- 手動コールバック (threading ベース) ---
print("\n--- 手動コールバック (threading ベース) ---")

def async_operation(task_name: str, duration: float, on_complete):
    """非同期に処理を実行し、完了時にコールバックを呼ぶ"""
    def _run():
        result = simulate_io(task_name, duration)
        on_complete(task_name, result)  # ← 完了時にコールバック実行
    thread = threading.Thread(target=_run)
    thread.start()
    return thread

completion_results = []
completion_event = threading.Event()
expected_count = 3

def on_task_complete(name: str, result: str):
    """コールバック関数: タスク完了時に呼ばれる"""
    print(f"  📬 コールバック受信: {name} → {result}")
    completion_results.append(result)
    if len(completion_results) >= expected_count:
        completion_event.set()

threads = []
for name, dur in [("タスクA", 1.0), ("タスクB", 1.5), ("タスクC", 0.5)]:
    t = async_operation(name, dur, on_task_complete)
    threads.append(t)

completion_event.wait()  # 全タスク完了を待つ
for t in threads:
    t.join()
print(f"  全結果: {completion_results}")

# --- concurrent.futures の add_done_callback ---
print("\n--- concurrent.futures の add_done_callback ---")

def on_future_done(future: concurrent.futures.Future):
    """Futureの完了時に呼ばれるコールバック"""
    result = future.result()
    print(f"  📬 Future コールバック: {result}")

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    for name, dur in [("タスクX", 0.8), ("タスクY", 0.5), ("タスクZ", 1.2)]:
        future = executor.submit(simulate_io, name, dur)
        future.add_done_callback(on_future_done)  # ← コールバック登録

print("  全 Future 完了")

今回紹介した方法で実務に適したものを選んでみてください

会社紹介

私が所属しているアドバンスド・ソリューション株式会社(以下、ADS)は一緒に働く仲間を募集しています

会社概要
「技術」×「知恵」=顧客課題の解決・新しい価値の創造

この方程式の実現はADSが大切にしている考えで、技術を磨き続けるgeekさと、顧客を思うloveがあってこそ実現できる世界観だと思っています
この『love & geek』の精神さえあれば、得意不得意はno problem!
技術はピカイチだけど顧客折衝はちょっと苦手。OKです。技術はまだ未熟だけど顧客と知恵を出し合って要件定義するのは大好き。OKです
凸凹な社員の集まり、色んなカラーや柄の個性が集まっているからこそ、常に新しいソリューションが生まれています

ミッション
私たちは、テクノロジーを活用し、業務や事業の生産性向上と企業進化を支援します

ホームページ
アドバンスド・ソリューション株式会社|ADS Co., Ltd.
アドバンスド・ソリューション株式会社|ADS Co., Ltd.

Microsoft 365/SharePoint/Power Platform/Azure による DX コンサル・シス ...

サイトへ移動

PR

-プログラム, 技術
-