发布订阅注解方式
Project description
版本说明
0.1.0
- 在方法上添加@task就可以直接开启一个线程执行任务
- 常用参数:
- 生产者:
role="producer",topic="自定义" - 消费者:
role="consumer",subscribe="自定义",publish_after="自定义"
- 生产者:
- 说明:
role: 用于区分上下游producer/consumertopic: 发布的主题subscribe: 对哪些主题感兴趣publish_after: 消费完后继续发布新主题
- 特别的:方法的返回值可以进行传递,若没有返回值则不会发布新主题(只针对消费者和消费者之间)
0.2.x
- dotask暴露
topic_manager=>可以更好的管理生产消费者 - 迁移
logger至util包下 - 添加本地调用shell工具函数->可以很好的适配@task消费者
max_concurrent:控制并发执行的shell数量- 引入
dotask.Shell调用Shell(max_concurrent=?).local_shell_execute(cmd,callback)来使用
0.3.x
- 添加令牌桶,通过
from dotask.util import TokenBucket使用,可以更好的限制生产消费速率- 参数说明
- TokenBucket(capacity=1, rate=0),capacity:令牌桶容量,rate:每秒生成速率
- get_token(block=True),block:决定令牌桶无令牌时是阻塞还是返回false
- 参数说明
topic_manager内部维护unfinished_tasks队列,记录所有消费失败的数据及其异常- 添加定时器
Timer,可以配合unfinished_tasks再次处理未消费完成的数据
快速开始
安装
pip install dotask
示例
- 简单的生产发布模式
from dotask import task
from dotask.util import logger
import random
if __name__ == '__main__':
@task(role="producer",topic="scan")
def a():
return random.randint(1,10)
@task(role="consumer",subscribe="scan",publish_after="sayHi")
def b(data):
if data >5:
return data
logger.warning(f"本次生成数字:{data},不会继续发布sayHi主题")
@task(role="consumer",subscribe="sayHi")
def c(data):
logger.debug(f"c触发")
logger.info(f"消费:{data}")
//=====================================================================
2026-02-08 17:30:37 - INFO - 发布主题-[scan]:5
2026-02-08 17:30:37 - WARNING - 本次生成数字:5,不会继续发布sayHi主题
2026-02-08 17:30:38 - INFO - 发布主题-[scan]:8
2026-02-08 17:30:38 - INFO - 发布主题-[sayHi]:8
2026-02-08 17:30:38 - DEBUG - c触发
2026-02-08 17:30:38 - INFO - 消费:8
2026-02-08 17:30:38 - INFO - 发布主题-[scan]:6
2026-02-08 17:30:38 - INFO - 发布主题-[sayHi]:6
2026-02-08 17:30:38 - DEBUG - c触发
2026-02-08 17:30:38 - INFO - 消费:6
//======================================================================
- 令牌桶+定时器使用
from dotask import task,topic_manager
from dotask.util import logger,TokenBucket,Timer
import random
token_bucket = TokenBucket(capacity=1, rate=1)
def error_inspect():
logger.warning(f"未完成任务:{topic_manager.unfinished_tasks.qsize()}")
item = topic_manager.unfinished_tasks.get()
logger.warning(f"其一:{item}")
logger.warning(f"重新发布失败数据,{item.get("unfinished_data")}")
topic_manager.publish(item.get("occurred"),item.get("unfinished_data"))
if __name__ == '__main__':
Timer(interval=2,task=error_inspect).start()
@task(role="producer",topic="scan")
def a():
if token_bucket.get_token(block=True):
return random.randint(1,10)
else:
logger.warning("令牌桶无可用令牌,本次生成数据失败")
return -1
@task(role="consumer",subscribe="scan",publish_after="sayHi")
def b(data):
if data >5:
return data
logger.warning(f"本次生成数字:{data},不会继续发布sayHi主题")
@task(role="consumer",subscribe="sayHi")
def c(data):
logger.debug(f"c触发,接收数据{data}")
raise Exception("处理失败")
//===========================================================================
2026-02-11 05:33:49 - INFO - Timer started,runs every 2 seconds
2026-02-11 05:33:49 - INFO - publish topic-[scan]:1
2026-02-11 05:33:49 - WARNING - 本次生成数字:1,不会继续发布sayHi主题
2026-02-11 05:33:50 - INFO - publish topic-[scan]:3
2026-02-11 05:33:50 - WARNING - 本次生成数字:3,不会继续发布sayHi主题
2026-02-11 05:33:50 - INFO - publish topic-[scan]:1
2026-02-11 05:33:50 - WARNING - 本次生成数字:1,不会继续发布sayHi主题
2026-02-11 05:33:51 - WARNING - 未完成任务:0
2026-02-11 05:33:51 - INFO - publish topic-[scan]:8
2026-02-11 05:33:51 - INFO - publish topic-[sayHi]:8
2026-02-11 05:33:51 - DEBUG - c触发,接收数据8
2026-02-11 05:33:51 - ERROR - Consumer [None] failed to process message: 处理失败
2026-02-11 05:33:51 - WARNING - 其一:{'unfinished_data': 8, 'occurred': 'sayHi', 'error': Exception('处理失败')}
2026-02-11 05:33:51 - WARNING - 重新发布失败数据,8
2026-02-11 05:33:51 - INFO - publish topic-[sayHi]:8
2026-02-11 05:33:51 - DEBUG - c触发,接收数据8
2026-02-11 05:33:51 - ERROR - Consumer [None] failed to process message: 处理失败
2026-02-11 05:33:51 - INFO - publish topic-[scan]:7
2026-02-11 05:33:51 - INFO - publish topic-[sayHi]:7
2026-02-11 05:33:51 - DEBUG - c触发,接收数据7
//==============================================================================
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
dotask-0.3.2.tar.gz
(14.1 kB
view details)
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
dotask-0.3.2-py3-none-any.whl
(13.6 kB
view details)
File details
Details for the file dotask-0.3.2.tar.gz.
File metadata
- Download URL: dotask-0.3.2.tar.gz
- Upload date:
- Size: 14.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.2
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
8c83d89a688e85ee9f60ff31375c359043f6ceb858cd0ead86214b3529fd835b
|
|
| MD5 |
c5b1a9aaaa020dc634e61d9e06c7ab04
|
|
| BLAKE2b-256 |
52409c41c04dec81b6e9a66929aa83bd8468554f5298ec96f6c90997b9912758
|
File details
Details for the file dotask-0.3.2-py3-none-any.whl.
File metadata
- Download URL: dotask-0.3.2-py3-none-any.whl
- Upload date:
- Size: 13.6 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.2
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b3ee23ac1ec14ee9a963eee39cfa7a4e8b3251c8464cd8440e4d48247839b6e1
|
|
| MD5 |
9f76bfcff15ecf2cfd54eb94112e4dcc
|
|
| BLAKE2b-256 |
823d78afc02a24f0cd68a2e8362fb1e9700baa842fdfe6f1272a06fd0e1446e8
|