Skip to main content

AsyncX Tools

AsyncX Toolsは、非同期タスクの優先度付き実行と、同期関数・非同期関数の相互変換を 小さなAPIで扱うPythonライブラリです。

主な機能

  • 優先度と追加順を考慮した非同期タスクの実行
  • 同時実行数の制限
  • タスク単位のタイムアウト
  • バッチ全体のエラー収集
  • 同期関数をブロッキングせずに呼び出すsync_to_async
  • 非同期関数を同期コードから呼び出すasync_to_sync
  • ContextVarの引き継ぎ
  • Python 3.10〜3.13と型情報(py.typed)のサポート

動作環境

  • Python 3.10以上
  • 実行時の外部依存なし

インストール

python -m pip install asyncx-tools

インストール後はasyncxパッケージから公開APIをインポートします。

import asyncx

print(asyncx.__version__)

クイックスタート

import asyncio

from asyncx import Task, TaskManager


async def fetch_value(value: str, delay: float) -> str:
    await asyncio.sleep(delay)
    return value


async def main() -> None:
    manager = TaskManager[str](max_concurrent_tasks=2)
    tasks = [
        Task("low", fetch_value, args=("low", 0.1), priority=1),
        Task("high", fetch_value, args=("high", 0.1), priority=10),
        Task("medium", fetch_value, args=("medium", 0.1), priority=5),
    ]

    results = await manager.run_tasks(tasks)
    print(results)


asyncio.run(main())

優先度は開始順を決めます。並行実行時の完了順は各タスクの処理時間によって変わります。

Task

Taskは実行対象と実行条件を保持する不変データクラスです。

Task(
    name,
    func,
    priority=0,
    timeout=None,
    args=(),
    kwargs={},
)
引数 説明
name 結果とエラーのキーになる空でない文字列
func Awaitableを返す呼び出し可能オブジェクト
priority 整数。値が大きいタスクから開始
timeout 0以上の有限な秒数。Noneは無制限
args funcへ渡す位置引数
kwargs funcへ渡すキーワード引数

同じ優先度では追加順が維持されます。nameprioritytimeoutなどのフィールドは 作成後に再代入できず、kwargsもコピーされた読み取り専用マッピングになります。

TaskManager

作成

from asyncx import TaskManager

manager = TaskManager[str](max_concurrent_tasks=10)

max_concurrent_tasksは1以上の整数です。同時に実行されるタスク数の上限になります。

メソッド 説明
await add_task(task) 優先度付きキューへ1件追加
await run_task(task) 1件を直ちに実行
await run_tasks(tasks=None, *, raise_on_error=True) キューまたは指定バッチを実行
get_results() 成功結果のコピーを取得
get_errors() エラーのコピーを取得
clear() 待機キュー、結果、エラーを消去

タスクを直接渡して実行

results = await manager.run_tasks(
    [
        Task("first", fetch_value, args=("A", 0.1)),
        Task("second", fetch_value, args=("B", 0.1)),
    ]
)

先にキューへ追加して実行

await manager.add_task(Task("first", fetch_value, args=("A", 0.1)))
await manager.add_task(Task("second", fetch_value, args=("B", 0.1)))

results = await manager.run_tasks()

run_tasks(tasks)へ明示的にタスクを渡すと、追加済みキューはそのバッチで置き換わります。 空リストを渡した場合も追加済みキューは破棄されます。

単一タスクの実行

result = await manager.run_task(Task("single", fetch_value, args=("value", 0.1)))

run_task()も同時実行数の制限を使用し、結果またはエラーをManagerへ記録します。

結果とエラー

results = manager.get_results()
errors = manager.get_errors()

両メソッドは内部状態のコピーを返します。返された辞書を変更してもManager内部には 影響しません。結果とエラーは、run_tasks()を開始するたびに初期化されます。

エラーを収集して処理を継続

既定ではすべてのタスクを完走した後、1件以上失敗していればTaskErrorを送出します。

import asyncio

from asyncx import Task, TaskManager


async def succeeds() -> str:
    return "ok"


async def fails() -> str:
    raise ValueError("failed")


async def main() -> None:
    manager = TaskManager[str]()
    results = await manager.run_tasks(
        [Task("success", succeeds), Task("failure", fails)],
        raise_on_error=False,
    )

    print(results)  # {"success": "ok"}
    print(manager.get_errors())  # {"failure": TaskError(...)}


asyncio.run(main())

個別タスクの例外はTaskErrorへ変換され、元の例外は__cause__に保持されます。 Managerが適用したタイムアウトはTaskTimeoutErrorになります。タスク本体が自ら送出した TimeoutErrorは、Managerのタイムアウトとは区別して通常のTaskErrorとして扱われます。

キャンセルと同時操作

  • run_tasks()をキャンセルすると、実行中のワーカーと兄弟タスクもキャンセルされます。
  • タスク自身がCancelledErrorを送出した場合も、兄弟タスクを残留させません。
  • 複数のrun_tasks()呼び出しは、同じManager内で順番に実行されます。
  • 実行中のadd_task()clear()は競合防止のためRuntimeErrorになります。
  • clear()は待機中のキュー、結果、エラーを消去します。

キャンセルは協調的です。タスク内にawaitがなくイベントループを占有する処理は、即座に 停止できません。

sync_to_async

同期関数をイベントループの外に移し、非同期関数として呼び出せるようにします。

sync_to_async(func, thread_sensitive=False, executor=None)
import asyncio
import time

from asyncx import sync_to_async


@sync_to_async
def blocking_io(value: int) -> int:
    time.sleep(0.2)
    return value * 2


async def main() -> None:
    results = await asyncio.gather(
        blocking_io(1),
        blocking_io(2),
        blocking_io(3),
    )
    print(results)  # [2, 4, 6]


asyncio.run(main())

通常の関数呼び出し形式でも利用できます。

def blocking_io_plain(value: int) -> int:
    time.sleep(0.2)
    return value * 2


async_blocking_io = sync_to_async(blocking_io_plain)

スレッド動作

設定 動作 主な用途
既定値 複数スレッドで並行実行 ファイル、ネットワークなどスレッドセーフなI/O
thread_sensitive=True 共通の専用スレッドで逐次実行 同一スレッドを要求する既存コード
executor=pool 指定したThreadPoolExecutorで実行 ワーカー数やライフサイクルを制御したい場合

カスタムExecutorはthread_sensitive=Falseの場合だけ指定できます。

import time
from concurrent.futures import ThreadPoolExecutor


def blocking_io_plain(value: int) -> int:
    time.sleep(0.2)
    return value * 2


async def use_custom_pool() -> int:
    with ThreadPoolExecutor(max_workers=4) as pool:
        async_func = sync_to_async(
            blocking_io_plain,
            thread_sensitive=False,
            executor=pool,
        )
        return await async_func(10)

呼び出し元のContextVarはワーカースレッドへ引き継がれます。呼び出し側のasyncioタスクを キャンセルしても、スレッド上ですでに開始した同期関数自体は停止できません。同期関数の 終了後、ワーカーは再利用されます。

スレッド化は、Pythonコード主体のCPU負荷を必ず高速化するものではありません。 主な対象はブロッキングI/Oです。

非同期関数やasync __call__を持つオブジェクトを渡した場合は二重変換せず、そのまま返します。

async_to_sync

非同期関数を同期コードから呼び出せる関数に変換します。

async_to_sync(func, *, timeout=None)
import asyncio

from asyncx import async_to_sync


async def fetch_number(value: int) -> int:
    await asyncio.sleep(0.1)
    return value


fetch_number_sync = async_to_sync(fetch_number, timeout=1.0)
result = fetch_number_sync(10)

timeoutには0以上の有限な秒数、またはNoneを指定できます。Managerとは独立した変換API なので、タイムアウト時はRuntimeErrorを送出します。非同期関数自身が送出した TimeoutErrorは変換せず、そのまま伝播します。

同期コンテキストでは、呼び出しスレッド上に新しいイベントループを作成します。実行中の イベントループがあるスレッドから呼ばれた場合は、デッドロックを避けるため別スレッドの イベントループを使用します。ただし呼び出し元スレッドは完了まで同期的に待機します。

非同期コード内ではasync_to_syncを使わず、対象関数を直接awaitする方が効率的です。 特定のイベントループに結び付いたFutureなどを別スレッドへ持ち込まないでください。

非同期関数だけでなく、async __call__を持つオブジェクトにも対応します。

例外

from asyncx import TaskError, TaskTimeoutError
例外 意味
TaskError タスク失敗、またはバッチ内に1件以上の失敗があった
TaskTimeoutError Task.timeoutで指定した時間を超過した

TaskTimeoutErrorTaskErrorのサブクラスです。

性能のヒント

  • TaskManagermax_concurrent_tasksは接続先やリソースの上限に合わせて調整してください。
  • 優先度は開始順を制御しますが、完了順までは保証しません。
  • スレッドセーフなI/O関数では、既定のsync_to_asyncが並行性を活用します。
  • スレッド固定が不要な処理にthread_sensitive=Trueを指定すると、全呼び出しが直列になります。
  • 非同期関数内ではtime.sleep()などを避け、asyncio.sleep()またはsync_to_asyncを使います。
  • CPU負荷の高いPython処理には、スレッドではなくプロセス分離も検討してください。

1.0.0へ移行する場合の注意

  • sync_to_asyncの既定値は並行実行です。スレッド固定が必要なら thread_sensitive=Trueを明示してください。
  • Taskの設定は作成後に変更できません。変更が必要な場合は新しいTaskを作成してください。
  • タイムアウトには有限値だけを指定できます。NaNと無限大は拒否されます。

その他の変更点はCHANGELOG.mdを参照してください。

開発

Windows PowerShellでは、リポジトリのルートで次を実行します。

python -m venv .venv
.\.venv\Scripts\Activate.ps1
python -m pip install -e ".[dev]"
python -m pytest
python -m mypy asyncx
python -m ruff check .
python -m ruff format --check .
python -m build --sdist --wheel
python -m twine check dist/*

CIではPython 3.10〜3.13のテスト、Ruff、mypy、カバレッジ、配布物のビルド、 Twine検査、wheelからのimportを確認します。

asyncx/    ライブラリ本体
tests/     テスト
examples/  実行例

ライセンス

MIT Licenseです。詳細はLICENSEを参照してください。

利用報告や、このライブラリへのリンクは必須ではありませんが歓迎します。

Release files for asyncx-tools 1.0.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for asyncx-tools 1.0.0
File Size Uploaded
asyncx_tools-1.0.0.tar.gz 20.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for asyncx-tools 1.0.0
File Interpreter ABI Platform
asyncx_tools-1.0.0-py3-none-any.whl Python 3 none any Details

Total release size: 33.0 kB

Release files / asyncx_tools-1.0.0.tar.gz

Download URL asyncx_tools-1.0.0.tar.gz
Size 20.7 kB
Tags Source
SHA-256 checksum
How to use checksums
63ebced786de636139cb9546d71b6198739b4d728b1c5f3fdf585fdbd0269003
BLAKE2b-256 checksum
How to use checksums
b6fed0e8c6bfe8a6146d1400baf3f4d4f460890f3da4ca68046f0533b86ca17e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.3

Release files / asyncx_tools-1.0.0-py3-none-any.whl

Download URL asyncx_tools-1.0.0-py3-none-any.whl
Size 12.3 kB
Tags Python 3
SHA-256 checksum
How to use checksums
9fd28eff85e30e2bcd59584fea90d0c7d39142a0fb3bcac5607649022a4b6ef1
BLAKE2b-256 checksum
How to use checksums
833d9ff6b8b39c95a533c380bfb73afd03d95e1104425060058bffaf7bb4dc37
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.3

Release history Release notifications | RSS feed

This release

1.0.0 This release

2 release files

0.2.0

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page