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へ渡すキーワード引数 |
同じ優先度では追加順が維持されます。name、priority、timeoutなどのフィールドは
作成後に再代入できず、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で指定した時間を超過した |
TaskTimeoutErrorはTaskErrorのサブクラスです。
性能のヒント
TaskManagerのmax_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)
| File | Size | Uploaded | |
|---|---|---|---|
| asyncx_tools-1.0.0.tar.gz | 20.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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
|