A powerful async task processing kit based on RabbitMQ with Coroutine, Thread, and Process support.
Project description
async-task-kit
A powerful async task processing kit based on RabbitMQ with Coroutine, Thread, and Process support.
Installation
pip install async-task-kit
Features
- RabbitMQ Client: Robust connection pooling and delay/dead-letter queue support.
- Multiple Consumer Models: Support for Coroutine (asyncio), Thread, and Process consumers depending on your workload (I/O-bound vs CPU-bound).
- Extensible Processor: Easily define your task logic by inheriting
TaskProcessor. - Built-in Logger & EnvLoader: Useful utilities for production-ready applications.
Quick Start
1. Configuration (.env)
Use the built-in EnvLoader to manage your environment variables. Create a .env file:
RABBITMQ_URL=amqp://guest:guest@localhost/
TASK_IDS=demo_task
# You can also configure specific task settings using the {TASK_ID}_ prefix
DEMO_TASK_QUEUE_NAME=my_demo_queue
DEMO_TASK_CONCURRENCY=3
2. Define your Task Processor (demo_processor.py)
import logging
from async_task_kit import TaskProcessor
logger = logging.getLogger(__name__)
class DemoProcessor(TaskProcessor):
async def process(self, task: dict):
logger.info(f"Processing task: {task}")
# Return any truthy value (e.g., dict, object, True) for success and pass to callback.
# Return None or False to trigger retry.
return {"status": "ok", "processed_data": task}
async def callback(self, task: dict, result: any):
logger.info(f"Task completed with result: {result}")
3. Main Consumer Application (main.py)
A production-ready setup with signal handling for graceful shutdown.
import asyncio
import logging
import signal
from typing import List, Type
from demo_processor import DemoProcessor
# -------------- 只需要改这一行来切换并发模型 --------------
from async_task_kit import CoroutineConsumer as Consumer
# from async_task_kit import ThreadConsumer as Consumer
# from async_task_kit import ProcessConsumer as Consumer
# --------------------------------------------------------
from async_task_kit import TaskProcessor, EnvLoader, setup_logger
# Initialize logger
setup_logger()
logger = logging.getLogger(__name__)
# Register your processors
TASK_REGISTRY: dict[str, Type[TaskProcessor]] = {
"demo_task": DemoProcessor,
}
consumers: List[Consumer] = []
async def run_all_consumers(amqp_url: str, task_ids: List[str]):
tasks = []
for task_id in task_ids:
if task_id not in TASK_REGISTRY:
continue
processor_cls = TASK_REGISTRY[task_id]
processor = processor_cls(task_id=task_id)
consumer = Consumer(
amqp_url=amqp_url,
queue_name=processor.queue_name,
processor=processor,
concurrency=processor.concurrency,
)
consumers.append(consumer)
tasks.append(consumer.start())
logger.info(f"🚀 启动任务 [{task_id}] | queue={processor.queue_name} | 并发={processor.concurrency}")
await asyncio.gather(*tasks)
async def shutdown_all():
logger.info("🛑 优雅关闭所有消费者...")
for consumer in consumers:
await consumer.stop()
logger.info("✅ 所有消费者已关闭")
def handle_exit_signal(*args, **kwargs):
asyncio.create_task(shutdown_all())
async def main():
env = EnvLoader()
amqp_url = env.get("RABBITMQ_URL")
task_ids_str = env.get("TASK_IDS", "").strip()
if not task_ids_str:
logger.warning("⚠️ 未配置 TASK_IDS")
return
task_ids = [t.strip() for t in task_ids_str.split(",") if t.strip()]
valid_tasks = [t for t in task_ids if t in TASK_REGISTRY]
loop = asyncio.get_running_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
loop.add_signal_handler(sig, handle_exit_signal)
await run_all_consumers(amqp_url, valid_tasks)
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
logger.info("👋 服务已安全退出")
4. Publishing Tasks (publisher.py)
import asyncio
import logging
from async_task_kit import RabbitMQ, setup_logger
setup_logger()
logger = logging.getLogger(__name__)
async def publish():
rmq = RabbitMQ("amqp://guest:guest@localhost/")
await rmq.init()
await rmq.push("my_demo_queue", {"message": "Hello from async-task-kit!"})
logger.info("Task published successfully.")
await rmq.close()
if __name__ == "__main__":
asyncio.run(publish())
5. RabbitMQ Client API
RabbitMQ 提供连接池、push / pop、延迟重试队列,以及 queue_length 查询队列深度。
queue_length
from async_task_kit import RabbitMQ
rmq = RabbitMQ("amqp://guest:guest@localhost/")
await rmq.init()
length = await rmq.queue_length("my_demo_queue")
if length == RabbitMQ.QUEUE_LENGTH_UNAVAILABLE:
# -1:连不上或队列不存在,稍后重试
...
elif length == 0:
# 队列存在且为空
...
else:
# 当前消息数
...
await rmq.close()
| 返回值 | 含义 |
|---|---|
>= 0 |
队列真实消息数(0 = 空队列) |
-1 (QUEUE_LENGTH_UNAVAILABLE) |
连接失败、队列不存在或其它异常,应重试 |
实现要点:
- 使用
passive=True声明:只查询已存在的队列,不会自动创建空队列。 - 断连时会关闭旧连接池并自动重建(与
push/pop一致)。
pop 与队列声明
pop 取消息前会 declare_queue(durable=True)(默认),参数须与 push 创建队列时一致。
queue_length 使用 passive 模式,无需传 durable。
Consumer 空队列轮询
basic_get 空队列时 broker 立即 返回 GetEmpty(不会阻塞等新消息)。pop 使用 get(timeout=None, fail=False),不要用 timeout=0(会导致永远无法取到消息)。
空队列时的轮询节奏由 Consumer 的 poll_interval 控制(默认 1.0 秒):
consumer = CoroutineConsumer(
...,
poll_interval=2.0, # 空队列时每 2 秒 pop 一次
)
Changelog
各版本变更说明见 CHANGELOG.md。
License
MIT
Contact & Support
If you have any questions, suggestions, or need help with this library, feel free to reach out!
WeChat (微信): realwrtoff
Email: realwrtoff@gmail.com
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file async_task_kit-0.1.17.tar.gz.
File metadata
- Download URL: async_task_kit-0.1.17.tar.gz
- Upload date:
- Size: 13.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.11.2 {"installer":{"name":"uv","version":"0.11.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"macOS","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
13209021d3ac2fab8334cf3f9f766c0276572009b045ddc1a8791456ea6bd90d
|
|
| MD5 |
a0dde55755589650be48861408337cdc
|
|
| BLAKE2b-256 |
5cbe0ac256d9d698cecb24a8f0b9f3ee609388fefca737ce2ec873d7b1d20019
|
File details
Details for the file async_task_kit-0.1.17-py3-none-any.whl.
File metadata
- Download URL: async_task_kit-0.1.17-py3-none-any.whl
- Upload date:
- Size: 14.9 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: uv/0.11.2 {"installer":{"name":"uv","version":"0.11.2","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"macOS","version":null,"id":null,"libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
998d9b2e082ff8707b78c4997a8ff8bb8ba9335a3554722ddbd886ba6f31a0be
|
|
| MD5 |
9fb2e2f9cb0a9027b11367d58bc61ccb
|
|
| BLAKE2b-256 |
0360ad6884824d041d21ed93041416d5739f63be928354e9481a841975ec05ad
|