API 参考¶
核心模块¶
异步任务管理器
一个纯Python实现的异步任务管理器,支持线程池和动态伸缩。
- class fish_async_task.TaskManager(instance_key='default')[源代码]¶
基类:
object纯Python实现的异步任务管理器(线程池 + 动态伸缩)
- 参数:
instance_key (str)
- 返回类型:
- DEFAULT_QUEUE_SIZE = 1000¶
- DEFAULT_MIN_WORKERS = 1¶
- DEFAULT_IDLE_TIMEOUT = 60¶
- DEFAULT_TASK_STATUS_TTL = 3600¶
- DEFAULT_MAX_TASK_STATUS_COUNT = 10000¶
- DEFAULT_CLEANUP_INTERVAL = 300¶
- DEFAULT_THREAD_JOIN_TIMEOUT = 2¶
- DEFAULT_TASK_TIMEOUT = None¶
- DEFAULT_MIN_MAX_WORKERS = 4¶
- CPU_MULTIPLIER = 4¶
- static __new__(cls, instance_key='default')[源代码]¶
单例模式实现
使用双重检查锁定模式确保线程安全。 支持通过 instance_key 创建多个不同的单例实例。
使用场景: - 默认情况下,使用 "default" 作为 instance_key,所有调用返回同一个实例 - 如果需要多个独立的任务管理器实例(例如:不同业务模块使用不同的管理器),
可以使用不同的 instance_key,每个 key 对应一个独立的单例实例
- 示例:
# 获取默认实例 manager1 = TaskManager() manager2 = TaskManager() # manager1 和 manager2 是同一个实例
# 获取不同业务模块的独立实例 order_manager = TaskManager(instance_key="order") payment_manager = TaskManager(instance_key="payment") # 独立的实例
- 参数:
instance_key (str) -- 实例键名,默认为 "default"。不同的 key 对应不同的单例实例。
- 返回:
任务管理器实例(单例)
- 返回类型:
- __init__(instance_key='default')[源代码]¶
初始化方法(单例模式下此方法不会重复执行)
注意:由于使用单例模式,实际的初始化在 __new__ 中通过 _init_task_manager 完成。此方法存在是为了符合Python对象 创建规范,主要用于验证 instance_key 的一致性。
单例模式实现说明: - Python 的对象创建流程:__new__ -> __init__ - 在单例模式中,__new__ 负责创建或返回已有实例 - 如果实例已存在,Python 仍会调用 __init__,但此时实例已经初始化完成 - 此方法会验证 instance_key 是否匹配,避免混淆
- 参数:
instance_key (str) -- 实例键名,应该与创建实例时使用的 key 一致
- 返回类型:
None
- submit_task(func, *args, block=False, timeout=None, **kwargs)[源代码]¶
提交任务到任务队列
此方法是线程安全的,可以在多个线程中并发调用。 任务会被添加到队列中,由工作线程异步执行。
- 参数:
- 返回:
任务ID(UUID格式的字符串)
- 返回类型:
- 抛出:
TaskQueueFullError -- 当队列已满且 block=False 时抛出
备注
任务函数应该是线程安全的
如果 block=False 且队列已满,会抛出 TaskQueueFullError 异常
如果 block=True,会阻塞等待直到队列有空间或超时
提交后任务状态为"pending",执行时变为"running",完成后变为"completed"或"failed"
示例
>>> def my_task(name: str, value: int) -> str: ... return f"Task {name} completed with value {value}" >>> >>> manager = TaskManager() >>> # 非阻塞模式提交任务 >>> task_id = manager.submit_task(my_task, "task1", value=100) >>> >>> # 阻塞模式提交任务(最多等待10秒) >>> task_id = manager.submit_task( ... my_task, "task2", value=200, block=True, timeout=10.0 ... )
- get_task_status(task_id)[源代码]¶
获取任务状态
此方法是线程安全的,可以在多个线程中并发调用。
- 参数:
task_id (str) -- 任务ID(由 submit_task 返回的字符串)
- 返回:
- 任务状态字典,包含以下字段:
status: 任务状态("pending"、"running"、"completed"、"failed")
submit_time: 提交时间(Unix时间戳,可选)
start_time: 开始执行时间(Unix时间戳,可选)
end_time: 结束时间(Unix时间戳,可选)
result: 任务执行结果(仅当status为"completed"时存在)
error: 错误信息(仅当status为"failed"时存在)
如果任务不存在或已被清理,则返回None
- 返回类型:
Optional[TaskStatusDict]
示例
>>> task_id = manager.submit_task(my_task, "task1") >>> status = manager.get_task_status(task_id) >>> if status: ... print(f"任务状态: {status['status']}") ... if status['status'] == 'completed': ... print(f"结果: {status.get('result')}")
- clear_task_status(task_id=None)[源代码]¶
清除指定任务状态或所有任务状态
此方法是线程安全的,可以在多个线程中并发调用。
- 参数:
task_id (str | None) -- 要清除的任务ID。如果为None,则清除所有任务状态。
- 返回类型:
None
备注
清除任务状态不会影响正在执行的任务
清除后,get_task_status将返回None
- shutdown()[源代码]¶
关闭任务管理器
优雅关闭流程: 1. 清除运行标志,停止接受新任务 2. 取消所有活动的 Timer 3. 发送退出信号给所有工作线程 4. 等待所有工作线程退出 5. 等待清理线程退出 6. 清理所有资源(线程列表、任务状态等)
注意:如果多次调用,只有第一次调用会生效。
- 返回类型:
None
- classmethod destroy_instance(instance_key='default')[源代码]¶
销毁指定的单例实例
清理并移除指定的单例实例,释放相关资源。 适用于需要完全释放 TaskManager 资源的场景。
示例
>>> # 销毁默认实例 >>> TaskManager.destroy_instance() >>> >>> # 销毁特定实例 >>> TaskManager.destroy_instance("order")
备注
销毁后,再次创建相同 instance_key 的实例会重新初始化
如果实例正在运行中,会先调用 shutdown() 清理资源
此方法是线程安全的
TaskManager¶
- class fish_async_task.TaskManager(instance_key='default')[源代码]¶
基类:
object纯Python实现的异步任务管理器(线程池 + 动态伸缩)
- 参数:
instance_key (str)
- 返回类型:
- DEFAULT_QUEUE_SIZE = 1000¶
- DEFAULT_MIN_WORKERS = 1¶
- DEFAULT_IDLE_TIMEOUT = 60¶
- DEFAULT_TASK_STATUS_TTL = 3600¶
- DEFAULT_MAX_TASK_STATUS_COUNT = 10000¶
- DEFAULT_CLEANUP_INTERVAL = 300¶
- DEFAULT_THREAD_JOIN_TIMEOUT = 2¶
- DEFAULT_TASK_TIMEOUT = None¶
- DEFAULT_MIN_MAX_WORKERS = 4¶
- CPU_MULTIPLIER = 4¶
- static __new__(cls, instance_key='default')[源代码]¶
单例模式实现
使用双重检查锁定模式确保线程安全。 支持通过 instance_key 创建多个不同的单例实例。
使用场景: - 默认情况下,使用 "default" 作为 instance_key,所有调用返回同一个实例 - 如果需要多个独立的任务管理器实例(例如:不同业务模块使用不同的管理器),
可以使用不同的 instance_key,每个 key 对应一个独立的单例实例
- 示例:
# 获取默认实例 manager1 = TaskManager() manager2 = TaskManager() # manager1 和 manager2 是同一个实例
# 获取不同业务模块的独立实例 order_manager = TaskManager(instance_key="order") payment_manager = TaskManager(instance_key="payment") # 独立的实例
- 参数:
instance_key (str) -- 实例键名,默认为 "default"。不同的 key 对应不同的单例实例。
- 返回:
任务管理器实例(单例)
- 返回类型:
- __init__(instance_key='default')[源代码]¶
初始化方法(单例模式下此方法不会重复执行)
注意:由于使用单例模式,实际的初始化在 __new__ 中通过 _init_task_manager 完成。此方法存在是为了符合Python对象 创建规范,主要用于验证 instance_key 的一致性。
单例模式实现说明: - Python 的对象创建流程:__new__ -> __init__ - 在单例模式中,__new__ 负责创建或返回已有实例 - 如果实例已存在,Python 仍会调用 __init__,但此时实例已经初始化完成 - 此方法会验证 instance_key 是否匹配,避免混淆
- 参数:
instance_key (str) -- 实例键名,应该与创建实例时使用的 key 一致
- 返回类型:
None
- submit_task(func, *args, block=False, timeout=None, **kwargs)[源代码]¶
提交任务到任务队列
此方法是线程安全的,可以在多个线程中并发调用。 任务会被添加到队列中,由工作线程异步执行。
- 参数:
- 返回:
任务ID(UUID格式的字符串)
- 返回类型:
- 抛出:
TaskQueueFullError -- 当队列已满且 block=False 时抛出
备注
任务函数应该是线程安全的
如果 block=False 且队列已满,会抛出 TaskQueueFullError 异常
如果 block=True,会阻塞等待直到队列有空间或超时
提交后任务状态为"pending",执行时变为"running",完成后变为"completed"或"failed"
示例
>>> def my_task(name: str, value: int) -> str: ... return f"Task {name} completed with value {value}" >>> >>> manager = TaskManager() >>> # 非阻塞模式提交任务 >>> task_id = manager.submit_task(my_task, "task1", value=100) >>> >>> # 阻塞模式提交任务(最多等待10秒) >>> task_id = manager.submit_task( ... my_task, "task2", value=200, block=True, timeout=10.0 ... )
- get_task_status(task_id)[源代码]¶
获取任务状态
此方法是线程安全的,可以在多个线程中并发调用。
- 参数:
task_id (str) -- 任务ID(由 submit_task 返回的字符串)
- 返回:
- 任务状态字典,包含以下字段:
status: 任务状态("pending"、"running"、"completed"、"failed")
submit_time: 提交时间(Unix时间戳,可选)
start_time: 开始执行时间(Unix时间戳,可选)
end_time: 结束时间(Unix时间戳,可选)
result: 任务执行结果(仅当status为"completed"时存在)
error: 错误信息(仅当status为"failed"时存在)
如果任务不存在或已被清理,则返回None
- 返回类型:
Optional[TaskStatusDict]
示例
>>> task_id = manager.submit_task(my_task, "task1") >>> status = manager.get_task_status(task_id) >>> if status: ... print(f"任务状态: {status['status']}") ... if status['status'] == 'completed': ... print(f"结果: {status.get('result')}")
- clear_task_status(task_id=None)[源代码]¶
清除指定任务状态或所有任务状态
此方法是线程安全的,可以在多个线程中并发调用。
- 参数:
task_id (str | None) -- 要清除的任务ID。如果为None,则清除所有任务状态。
- 返回类型:
None
备注
清除任务状态不会影响正在执行的任务
清除后,get_task_status将返回None
- shutdown()[源代码]¶
关闭任务管理器
优雅关闭流程: 1. 清除运行标志,停止接受新任务 2. 取消所有活动的 Timer 3. 发送退出信号给所有工作线程 4. 等待所有工作线程退出 5. 等待清理线程退出 6. 清理所有资源(线程列表、任务状态等)
注意:如果多次调用,只有第一次调用会生效。
- 返回类型:
None
- classmethod destroy_instance(instance_key='default')[源代码]¶
销毁指定的单例实例
清理并移除指定的单例实例,释放相关资源。 适用于需要完全释放 TaskManager 资源的场景。
示例
>>> # 销毁默认实例 >>> TaskManager.destroy_instance() >>> >>> # 销毁特定实例 >>> TaskManager.destroy_instance("order")
备注
销毁后,再次创建相同 instance_key 的实例会重新初始化
如果实例正在运行中,会先调用 shutdown() 清理资源
此方法是线程安全的
异常¶
配置模块¶
配置管理模块
负责加载和验证任务管理器的配置项。
本模块提供了从环境变量加载配置的功能,支持以下配置项: - TASK_STATUS_TTL: 任务状态TTL(秒),默认3600 - MAX_TASK_STATUS_COUNT: 最大任务状态数量,默认10000 - TASK_CLEANUP_INTERVAL: 清理间隔(秒),默认300 - TASK_TIMEOUT: 任务超时时间(秒),默认无限制 - ADAPTIVE_WORKER_ENABLED: 是否启用自适应线程管理,默认True - WORKER_CPU_THRESHOLD: CPU使用率阈值,默认0.8 - WORKER_QUEUE_THRESHOLD_HIGH: 扩容队列积压阈值,默认100 - WORKER_QUEUE_THRESHOLD_LOW: 缩容队列空闲阈值,默认10 - WORKER_SCALE_UP_COOLDOWN: 扩容冷却期(秒),默认5.0 - WORKER_SCALE_DOWN_COOLDOWN: 缩容冷却期(秒),默认30.0
性能优化配置项: - SHARD_COUNT: 分片数量,默认16(建议为2的幂次) - BATCH_UPDATE_BUFFER_SIZE: 批量更新缓冲区大小,默认100 - BATCH_UPDATE_INTERVAL: 批量更新刷新间隔(秒),默认0.1 - ENABLE_AUTO_CLEANUP: 是否启用自动清理,默认True - ENABLE_BATCH_UPDATES: 是否启用批量更新,默认False - ENABLE_ADAPTIVE_SCALING: 是否启用自适应扩展,默认False
所有配置项都会进行验证,无效值会被拒绝并使用默认值。
- class fish_async_task.config.ConfigLoader(logger)[源代码]¶
基类:
object配置加载器
- 参数:
logger (Logger)
- MAX_TTL = 86400¶
- MAX_TASK_STATUS_COUNT = 1000000¶
- MAX_CLEANUP_INTERVAL = 3600¶
- MAX_TASK_TIMEOUT = 86400¶
- MAX_CPU_THRESHOLD = 1.0¶
- MAX_QUEUE_THRESHOLD = 10000¶
- MAX_COOLDOWN = 3600¶
- DEFAULT_ADAPTIVE_WORKER_ENABLED = True¶
- DEFAULT_CPU_THRESHOLD = 0.8¶
- DEFAULT_QUEUE_THRESHOLD_HIGH = 100¶
- DEFAULT_QUEUE_THRESHOLD_LOW = 10¶
- DEFAULT_SCALE_UP_COOLDOWN = 5.0¶
- DEFAULT_SCALE_DOWN_COOLDOWN = 30.0¶
- DEFAULT_SHARD_COUNT = 16¶
- DEFAULT_BATCH_UPDATE_BUFFER_SIZE = 100¶
- DEFAULT_BATCH_UPDATE_INTERVAL = 0.1¶
- DEFAULT_ENABLE_AUTO_CLEANUP = True¶
- DEFAULT_ENABLE_BATCH_UPDATES = False¶
- DEFAULT_ENABLE_ADAPTIVE_SCALING = False¶
- MAX_SHARD_COUNT = 1024¶
- MAX_BATCH_UPDATE_BUFFER_SIZE = 10000¶
- MAX_BATCH_UPDATE_INTERVAL = 60.0¶
- load_int_config(env_key, default_value, config_name, min_value=1, max_value=None)[源代码]¶
加载并验证整数配置项
- 参数:
- 返回:
验证后的配置值
- 返回类型:
备注
如果环境变量不存在、不是有效整数或值超出允许范围, 将使用默认值并记录警告日志。
- class fish_async_task.config.HotReloadConfig(logger=None, reload_interval=60)[源代码]¶
基类:
object支持热重载的配置管理器
任务状态¶
任务状态管理模块
负责任务状态的更新、查询和清理。
- class fish_async_task.task_status.ReadWriteLock[源代码]¶
基类:
object写优先的读写锁实现
允许多个读操作并发执行,写操作独占访问。 使用写优先策略防止写饥饿:当有写操作等待时,新的读操作会排队。
基于 threading.Condition 实现,支持上下文管理器。 支持超时机制,防止在高并发场景下死锁。
使用场景: - 读多写少的场景,读操作可以并发执行 - 写操作需要独占访问,与所有读操作互斥 - 需要防止写饥饿的场景 - 需要防止死锁的场景(使用超时)
- 示例:
lock = ReadWriteLock()
# 读操作 with ReadWriteLockContext(lock, write=False):
# 多个线程可以同时进入这里 data = self._data
# 写操作 with ReadWriteLockContext(lock, write=True):
# 同一时间只有一个线程可以进入这里 self._data = new_data
# 带超时的读操作 if lock.acquire_read(timeout=5.0):
- try:
data = self._data
- finally:
lock.release_read()
- DEFAULT_TIMEOUT = 30.0¶
- acquire_read(timeout=None)[源代码]¶
获取读锁
如果有写操作在等待,新的读操作会排队等待,防止写饥饿。
- 参数:
timeout (float | None) -- 超时时间(秒)。如果为 None,使用默认超时时间。 如果为负数或 0,立即尝试获取,不等待。
- 返回:
如果成功获取锁则返回 True,超时则返回 False
- 返回类型:
- 抛出:
TimeoutError -- 等待超时后抛出异常
- acquire_write(timeout=None)[源代码]¶
获取写锁
写锁是独占的,会等待所有读操作完成后才能获取。 使用写优先策略,防止写饥饿。
- 参数:
timeout (float | None) -- 超时时间(秒)。如果为 None,使用默认超时时间。 如果为负数或 0,立即尝试获取,不等待。
- 返回:
如果成功获取锁则返回 True,超时则返回 False
- 返回类型:
- 抛出:
TimeoutError -- 等待超时后抛出异常
- class fish_async_task.task_status.ReadWriteLockContext(lock, write=False, timeout=None)[源代码]¶
基类:
object读写锁上下文管理器
提供便捷的锁获取和释放方式。支持超时参数。
- 参数:
lock (ReadWriteLock)
write (bool)
timeout (float | None)
- __init__(lock, write=False, timeout=None)[源代码]¶
初始化上下文管理器
- 参数:
lock (ReadWriteLock) -- 读写锁实例
write (bool) -- 是否为写操作(True=写锁,False=读锁)
timeout (float | None) -- 超时时间(秒),如果为 None 则使用默认超时
- __enter__()[源代码]¶
获取锁并返回上下文管理器实例
根据 write 参数决定获取读锁或写锁。 读锁允许并发获取,写锁独占访问。
- 返回:
返回自身实例,用于 with 语句块
- 返回类型:
- 抛出:
TimeoutError -- 获取锁超时
- class fish_async_task.task_status.ShardedTaskStatusWithExpiry(shard_count, ttl)[源代码]¶
基类:
object分片任务状态存储(带过期时间管理)
使用分片锁减少锁竞争,每个分片内部使用优先级队列管理过期时间。 支持高并发查询和更新,以及高效的增量清理。
线程安全说明: - 每个分片有独立的锁,不同分片的操作可以并发执行 - 同一分片内的操作串行化,保证线程安全 - 清理操作支持增量清理,避免长时间阻塞
- shards: List[Dict[str, TaskStatusDict]]¶
- rw_locks: List[ReadWriteLock]¶
- get_status(task_id)[源代码]¶
获取任务状态(线程安全)
- 参数:
task_id (str) -- 任务ID
- 返回:
任务状态字典,如果任务不存在则返回None
- 返回类型:
Optional[TaskStatusDict]
- update_status(task_id, status, current_status=None)[源代码]¶
更新任务状态(线程安全)
- 参数:
task_id (str) -- 任务ID
status (TaskStatusDict) -- 新的任务状态字典
current_status (TaskStatusDict | None) -- 当前状态(如果已知,避免重复查询)
- 返回类型:
None
- enforce_max_count(max_count)[源代码]¶
强制执行最大任务数量限制
当任务状态数量超过限制时,按时间顺序清理最旧的任务。 使用优化策略:优先尝试增量清理,仅在必要时获取所有锁。
- get_all_statuses()[源代码]¶
获取所有任务状态(需要获取所有锁)
- 返回:
所有任务状态字典
- 返回类型:
Dict[str, TaskStatusDict]
- class fish_async_task.task_status.BatchedStatusUpdater(update_func, batch_size=100, flush_interval=0.1)[源代码]¶
基类:
object批量状态更新器
收集多个状态更新,批量提交到分片存储,减少锁获取次数,提升写入性能。 支持按批量大小和刷新间隔触发批量提交。
使用场景: - 高并发写入场景,减少状态更新的锁竞争 - 需要批量处理大量状态更新的场景
- 示例:
- def update_func(task_id, status, **kwargs):
# 实际的更新逻辑 pass
- updater = BatchedStatusUpdater(
update_func=update_func, batch_size=100, flush_interval=0.1
)
# 添加更新到批量队列 updater.update("task_1", "running") updater.update("task_2", "completed", result="result")
# shutdown时自动刷新所有待处理的更新 updater.shutdown()
- class fish_async_task.task_status.TaskStatusManager(logger, task_status_ttl, max_task_status_count, shard_count=None, batch_size=None, batch_flush_interval=None)[源代码]¶
基类:
object任务状态管理器
负责任务状态的存储、更新和查询。 使用分片锁和优先级队列优化性能,支持高并发操作。 支持批量状态更新,减少锁获取次数,提升写入性能。
线程安全说明: - 使用分片锁,不同分片的操作可以并发执行 - 同一分片内的操作串行化,保证线程安全 - 清理操作支持增量清理,避免长时间阻塞 - 批量更新器内部使用队列和锁,保证线程安全
- 参数:
- DEFAULT_SHARD_COUNT = 16¶
- DEFAULT_MAX_CLEANUP_PER_BATCH = 100¶
- DEFAULT_BATCH_SIZE = 100¶
- DEFAULT_BATCH_FLUSH_INTERVAL = 0.1¶
- __init__(logger, task_status_ttl, max_task_status_count, shard_count=None, batch_size=None, batch_flush_interval=None)[源代码]¶
初始化任务状态管理器
- update_task_status(task_id, status, start_time=None, end_time=None, result=None, error=None, submit_time=None)[源代码]¶
更新任务状态(线程安全)
- 参数:
task_id (str) -- 任务ID
status (Literal['pending', 'running', 'completed', 'failed']) -- 任务状态(pending, running, completed, failed)
start_time (float | None) -- 任务开始时间(可选)
end_time (float | None) -- 任务结束时间(可选)
result (Any) -- 任务执行结果(可选)
error (str | None) -- 错误信息(可选)
submit_time (float | None) -- 任务提交时间(可选,仅用于pending状态)
- 返回类型:
None
备注
此方法会保留已存在的 start_time,除非明确提供新的 start_time。
- get_task_status(task_id)[源代码]¶
获取任务状态
- 参数:
task_id (str) -- 任务ID
- 返回:
任务状态字典,如果任务不存在则返回None
- 返回类型:
Optional[TaskStatusDict]
- clear_task_status(task_id=None)[源代码]¶
清除指定任务状态或所有任务状态
- 参数:
task_id (str | None) -- 要清除的任务ID。如果为None,则清除所有任务状态。
- 返回类型:
None
- cleanup_old_task_status()[源代码]¶
清理过期的任务状态(增量清理)
清理策略: 1. 清理已完成或失败且超过TTL的任务(增量清理,每次最多100个) 2. 如果任务状态数量超过限制,清理最旧的任务
- 返回:
清理的任务数量
- 返回类型:
类型定义¶
类型定义模块
定义任务管理器相关的类型别名和类型定义。
- class fish_async_task.types.TaskStatusDict[源代码]¶
基类:
TypedDict任务状态字典类型定义
所有字段都是可选的(total=False),因为不同状态的任务包含的字段不同: - pending: status, submit_time - running: status, submit_time, start_time - completed: status, submit_time, start_time, end_time, result - failed: status, submit_time, start_time, end_time, error
所有可选字段都使用 Optional 注解,明确标识该字段可能为 None。
- class fish_async_task.types.ShardedTaskStatusDict[源代码]¶
基类:
TypedDict分片任务状态字典,包含分片索引信息
- task_status: TaskStatusDict¶
- class fish_async_task.types.BatchedUpdate[源代码]¶
基类:
TypedDict批量更新项
- status: TaskStatusDict¶
工作线程¶
工作线程模块
负责工作线程的创建、管理和任务执行。 支持自适应线程管理,根据CPU使用率和队列积压动态调整线程数量。
- class fish_async_task.worker.AdaptiveWorkerManager(min_workers, max_workers, cpu_threshold=0.8, queue_threshold_high=100, queue_threshold_low=10, scale_up_cooldown=5.0, scale_down_cooldown=30.0, use_cpu_monitoring=True)[源代码]¶
基类:
object自适应工作线程管理器
基于CPU使用率、队列积压和任务执行时间动态调整线程数量。 支持CPU监控,在CPU使用率过高时避免过度扩容。
使用场景: - 任务负载波动较大的场景 - 需要根据系统负载自动调整资源的场景 - CPU密集型和I/O密集型任务混合的场景
配置说明: - min_workers / max_workers: 线程数边界 - cpu_threshold: CPU使用率阈值,超过此值时避免扩容 - queue_threshold_high: 队列积压阈值,达到此值时触发扩容 - queue_threshold_low: 队列空闲阈值,低于此值时触发缩容 - scale_up_cooldown: 扩容冷却期(秒),避免频繁扩容 - scale_down_cooldown: 缩容冷却期(秒),避免频繁缩容
- 参数:
- __init__(min_workers, max_workers, cpu_threshold=0.8, queue_threshold_high=100, queue_threshold_low=10, scale_up_cooldown=5.0, scale_down_cooldown=30.0, use_cpu_monitoring=True)[源代码]¶
初始化自适应工作线程管理器
- 参数:
min_workers (int) -- 最小工作线程数
max_workers (int) -- 最大工作线程数
cpu_threshold (float) -- CPU使用率阈值(0.0-1.0),超过此值时避免扩容
queue_threshold_high (int) -- 队列积压高阈值,达到此值时触发扩容
queue_threshold_low (int) -- 队列空闲低阈值,低于此值且空闲超时后触发缩容
scale_up_cooldown (float) -- 扩容冷却期(秒)
scale_down_cooldown (float) -- 缩容冷却期(秒)
use_cpu_monitoring (bool) -- 是否启用CPU监控
- should_scale_up(current_workers, queue_size, cpu_usage=None)[源代码]¶
判断是否应该扩展线程
扩容条件(需全部满足): 1. 当前线程数小于最大限制 2. 距离上次扩容超过冷却期 3. 队列积压超过高阈值 或 CPU使用率低于阈值
- class fish_async_task.worker.CPUMonitor(sample_interval=0.5, sample_count=3)[源代码]¶
基类:
objectCPU使用率监控器
使用psutil获取系统CPU使用率,支持采样平均。
- class fish_async_task.worker.TaskExecutor(logger, task_timeout_getter, update_status_func, batch_size=None, batch_flush_interval=None)[源代码]¶
基类:
object任务执行器
负责任务的实际执行,包括超时控制和状态更新。 每个任务在独立的工作线程中执行,支持超时机制。
线程安全说明: - execute_task方法可以在多个线程中并发调用 - 任务执行是独立的,不会相互影响 - 超时机制使用daemon线程实现,超时后任务线程仍在后台运行 - 通过 cleanup_timed_out_tasks 方法可以清理超时的任务线程引用
- 参数:
- MAX_TRACKED_TIMEOUTS = 1000¶
- TIMEOUT_TASK_EXPIRY = 3600¶
- BATCH_SIZE = 100¶
- BATCH_FLUSH_INTERVAL = 0.1¶
- __init__(logger, task_timeout_getter, update_status_func, batch_size=None, batch_flush_interval=None)[源代码]¶
初始化任务执行器
- class fish_async_task.worker.WorkerManager(logger, task_queue, worker_threads, threads_lock, running_event, min_workers, max_workers, idle_timeout, task_timeout, execute_task_func, adaptive_worker_enabled=True, cpu_threshold=None, queue_threshold_high=None, queue_threshold_low=None, scale_up_cooldown=None, scale_down_cooldown=None, use_cpu_monitoring=True)[源代码]¶
基类:
object工作线程管理器
负责管理工作线程的生命周期,包括创建、扩展和回收。 支持动态线程池,根据任务队列大小和CPU使用率自动调整线程数量。 使用自适应策略,在负载高时扩容,空闲时缩容。
线程安全说明: - 所有对worker_threads列表的操作都在threads_lock保护下进行 - 线程退出时会从列表中安全移除,避免竞态条件 - 自适应扩缩容操作在threads_lock保护下进行 - 使用退出事件确认机制确保线程完全退出
- 参数:
logger (Logger)
task_queue (queue.Queue[TaskTuple])
threads_lock (allocate_lock)
running_event (Event)
min_workers (int)
max_workers (int)
idle_timeout (int)
task_timeout (float | None)
execute_task_func (Callable[[Tuple[str, Callable[[...], Any], Tuple[Any, ...], Dict[str, Any]]], None])
adaptive_worker_enabled (bool)
cpu_threshold (float | None)
queue_threshold_high (int | None)
queue_threshold_low (int | None)
scale_up_cooldown (float | None)
scale_down_cooldown (float | None)
use_cpu_monitoring (bool)
- QUEUE_GET_TIMEOUT = 1¶
- QUEUE_PUT_TIMEOUT = 1¶
- DEFAULT_CPU_THRESHOLD = 0.8¶
- DEFAULT_QUEUE_THRESHOLD_HIGH = 100¶
- DEFAULT_QUEUE_THRESHOLD_LOW = 10¶
- DEFAULT_SCALE_UP_COOLDOWN = 5.0¶
- DEFAULT_SCALE_DOWN_COOLDOWN = 30.0¶
- __init__(logger, task_queue, worker_threads, threads_lock, running_event, min_workers, max_workers, idle_timeout, task_timeout, execute_task_func, adaptive_worker_enabled=True, cpu_threshold=None, queue_threshold_high=None, queue_threshold_low=None, scale_up_cooldown=None, scale_down_cooldown=None, use_cpu_monitoring=True)[源代码]¶
初始化工作线程管理器
- 参数:
logger (Logger) -- 日志记录器
task_queue (Queue[Tuple[str, Callable[[...], Any], Tuple[Any, ...], Dict[str, Any]]]) -- 任务队列
threads_lock (allocate_lock) -- 线程锁
running_event (Event) -- 运行事件
min_workers (int) -- 最小工作线程数
max_workers (int) -- 最大工作线程数
idle_timeout (int) -- 空闲超时时间(秒)
task_timeout (float | None) -- 任务超时时间(秒)
execute_task_func (Callable[[Tuple[str, Callable[[...], Any], Tuple[Any, ...], Dict[str, Any]]], None]) -- 任务执行函数
adaptive_worker_enabled (bool) -- 是否启用自适应扩缩容
cpu_threshold (float | None) -- CPU使用率阈值(可选)
queue_threshold_high (int | None) -- 扩容队列积压阈值(可选)
queue_threshold_low (int | None) -- 缩容队列空闲阈值(可选)
scale_up_cooldown (float | None) -- 扩容冷却期(秒)(可选)
scale_down_cooldown (float | None) -- 缩容冷却期(秒)(可选)
use_cpu_monitoring (bool) -- 是否启用CPU监控
- start_initial_workers()[源代码]¶
启动初始工作线程
根据 min_workers 配置启动最小数量的工作线程。 这些线程会持续运行,不会被空闲超时机制回收。
- 返回类型:
None
高性能扩展¶
FishAsyncTask 性能优化模块
本模块提供性能优化功能,包括: - 分片任务状态存储(ShardedTaskStatus) - 优先级队列清理(TaskStatusWithExpiry) - 批量状态更新(BatchedStatusUpdater) - 自适应工作线程管理(AdaptiveWorkerManager) - 性能监控(PerformanceMetrics, SystemHealthMonitor) - 任务资源管理(TaskResourceManager) - 任务优先级队列(PriorityTaskQueue) - 任务取消管理(TaskCancellationManager)
所有优化遵循轻量化原则,核心功能仅使用 Python 标准库。
- class fish_async_task.performance.ShardedTaskStatus(shard_count=16)[源代码]¶
基类:
object分片任务状态存储
将任务状态分散到多个独立分片,每个分片有独立锁, 支持 10-15 倍的并发查询性能提升。
- 参数:
shard_count (int)
- shard_count¶
分片数量
- logger¶
日志记录器
- __init__(shard_count=16)[源代码]¶
初始化分片任务状态存储
- 参数:
shard_count (int) -- 分片数量,必须为正整数,建议为 2 的幂次 默认 16,在并发性和内存开销之间取得平衡
- 抛出:
ValueError -- 如果 shard_count < 1 或 shard_count > 1024
- 返回类型:
None
- get_status(task_id)[源代码]¶
获取任务状态(线程安全)
- 参数:
task_id (str) -- 任务 ID
- 返回:
任务状态字典,如果不存在返回 None
- 返回类型:
TaskStatusDict | None
- Performance:
O(1) 时间复杂度 线程安全:仅锁定单个分片
示例
>>> store = ShardedTaskStatus() >>> store.update_status("task-123", {"status": "completed", "result": "success"}) >>> status = store.get_status("task-123") >>> status["result"] 'success'
- update_status(task_id, status)[源代码]¶
更新任务状态(线程安全)
- 参数:
task_id (str) -- 任务 ID
status (TaskStatusDict) -- 新的任务状态字典
- 抛出:
TypeError -- 如果 status 不是 TaskStatusDict 类型
- 返回类型:
None
- Performance:
O(1) 时间复杂度 线程安全:仅锁定单个分片
示例
>>> store = ShardedTaskStatus() >>> store.update_status("task-123", {"status": "running"}) >>> store.get_status("task-123")["status"] 'running'
- remove_status(task_id)[源代码]¶
移除任务状态(线程安全)
- 参数:
task_id (str) -- 要移除的任务 ID
- 返回类型:
None
- Performance:
O(1) 时间复杂度 线程安全:仅锁定单个分片
示例
>>> store = ShardedTaskStatus() >>> store.update_status("task-123", {"status": "completed"}) >>> store.remove_status("task-123") >>> store.get_status("task-123") None
- get_task_count()[源代码]¶
获取当前任务数量
- 返回:
任务状态字典中的任务数量
- 返回类型:
- Performance:
O(n) 时间复杂度,n 为分片数量 线程安全(返回近似值,但非常准确)
示例
>>> store = ShardedTaskStatus() >>> store.get_task_count() 0 >>> store.update_status("task-1", {"status": "completed"}) >>> store.get_task_count() 1
- get_all_statuses()[源代码]¶
获取所有任务状态(需要获取所有锁)
警告
此方法会按顺序获取所有分片的锁,可能阻塞较长时间。 仅在必要时使用(如关闭、统计)。
- 返回:
所有任务状态的字典
- 返回类型:
- Performance:
O(n) 时间复杂度,n 为任务总数 线程安全:按顺序获取所有锁,避免死锁
示例
>>> store = ShardedTaskStatus() >>> store.update_status("task-1", {"status": "completed"}) >>> store.update_status("task-2", {"status": "running"}) >>> all_statuses = store.get_all_statuses() >>> len(all_statuses) 2
- class fish_async_task.performance.TaskStatusWithExpiry(ttl=300)[源代码]¶
基类:
object带过期时间的任务状态存储
使用优先级队列(最小堆)跟踪任务过期时间, 支持高效的增量清理操作。
- 参数:
ttl (int)
- ttl¶
任务状态生存时间(秒)
- logger¶
日志记录器
- __init__(ttl=300)[源代码]¶
初始化带过期时间的任务状态存储
- 参数:
ttl (int) -- 任务状态生存时间(秒),默认 300(5 分钟)
- 返回类型:
None
备注
清理操作会移除超过 TTL 的任务状态
- add_task(task_id, status)[源代码]¶
添加任务状态
- 参数:
task_id (str) -- 任务 ID
status (TaskStatusDict) -- 任务状态字典,必须包含 end_time 字段
- 返回类型:
None
- Behavior:
将任务添加到 status_dict
如果有 end_time,计算过期时间并添加到优先级队列
- 抛出:
ValueError -- 如果 status 不包含 end_time
- 参数:
task_id (str)
status (TaskStatusDict)
- 返回类型:
None
示例
>>> store = TaskStatusWithExpiry(ttl=300) >>> store.add_task("task-123", { ... "task_id": "task-123", ... "status": "completed", ... "end_time": time.time() ... })
- get_task(task_id)[源代码]¶
获取任务状态
- 参数:
task_id (str) -- 任务 ID
- 返回:
任务状态字典,如果不存在返回 None
- 返回类型:
TaskStatusDict | None
备注
此方法不锁定优先级队列(只读操作)
- cleanup_expired(max_cleanup=None)[源代码]¶
清理过期任务(增量清理)
- Performance:
O(k log n) 时间复杂度,k 为过期任务数量 通常 k << n,因此远快于全量扫描 O(n)
- Thread-Safety:
线程安全,使用 heap_lock 保护
备注
增量清理:每次最多清理 max_cleanup 个任务, 避免长时间阻塞其他操作
示例
>>> store = TaskStatusWithExpiry(ttl=300) >>> # 添加过期任务... >>> cleaned_count = store.cleanup_expired(max_cleanup=100) >>> print(f"清理了 {cleaned_count} 个过期任务")
- enforce_max_count(max_count)[源代码]¶
强制执行最大任务数量限制
当任务数量超过 max_count 时,删除最旧的任务(按 submit_time 或 start_time)。
- Performance:
O(n log n) 时间复杂度,n 为任务总数
- Thread-Safety:
线程安全,使用 heap_lock 保护
示例
>>> store = TaskStatusWithExpiry(ttl=300) >>> # 添加大量任务... >>> removed_count = store.enforce_max_count(max_count=10000) >>> print(f"移除了 {removed_count} 个旧任务")
- get_task_count()[源代码]¶
获取当前任务数量
- 返回:
任务状态字典中的任务数量
- 返回类型:
- Performance:
O(1) 时间复杂度
- Thread-Safety:
线程安全(返回近似值,但非常准确)
示例
>>> store = TaskStatusWithExpiry() >>> store.add_task("task-123", {"end_time": time.time()}) >>> store.get_task_count() 1
- class fish_async_task.performance.BatchedStatusUpdater(buffer_size=100, flush_interval=1.0, underlying_store=None)[源代码]¶
基类:
object批量状态更新器
将多个状态更新缓存到缓冲区,然后批量刷新到底层存储, 减少锁竞争和提高吞吐量。
- 参数:
buffer_size (int)
flush_interval (float)
underlying_store (Dict[str, TaskStatusDict] | None)
- buffer_size¶
触发自动刷新的缓冲区大小
- flush_interval¶
触发自动刷新的时间间隔(秒)
- underlying_store¶
底层任务状态存储(可选,用于测试)
示例
>>> store = {} >>> updater = BatchedStatusUpdater( ... buffer_size=100, ... flush_interval=1.0, ... underlying_store=store ... ) >>> updater.queue_update("task-1", {"status": "running"}) >>> updater.flush() # 手动刷新 1 >>> updater.close() # 关闭并刷新所有待处理更新
- __init__(buffer_size=100, flush_interval=1.0, underlying_store=None)[源代码]¶
初始化批量状态更新器
- 参数:
buffer_size (int) -- 触发自动刷新的缓冲区大小,默认 100
flush_interval (float) -- 触发自动刷新的时间间隔(秒),默认 1.0
underlying_store (Dict[str, TaskStatusDict] | None) -- 底层任务状态存储(可选,用于测试)
- 抛出:
ValueError -- 如果 buffer_size < 1 或 flush_interval <= 0
- 返回类型:
None
- queue_update(task_id, status)[源代码]¶
将状态更新排队到缓冲区
如果缓冲区达到 buffer_size,会自动触发刷新。
- 参数:
task_id (str) -- 任务 ID
status (TaskStatusDict) -- 任务状态字典
- 抛出:
RuntimeError -- 如果更新器已关闭
- 返回类型:
None
- Thread-Safety:
线程安全
示例
>>> updater = BatchedStatusUpdater() >>> updater.queue_update("task-1", {"status": "running"})
- flush()[源代码]¶
手动刷新缓冲区到底层存储
- 返回:
刷新的任务数量
- 返回类型:
- Thread-Safety:
线程安全
示例
>>> updater = BatchedStatusUpdater() >>> updater.queue_update("task-1", {"status": "running"}) >>> flushed = updater.flush() >>> print(f"刷新了 {flushed} 个任务")
- update_sync(task_id, status)[源代码]¶
同步更新(立即写入底层存储,不经过缓冲区)
用于需要立即更新的场景,例如关键状态变更。
- 参数:
task_id (str) -- 任务 ID
status (TaskStatusDict) -- 任务状态字典
- 抛出:
RuntimeError -- 如果更新器已关闭
- 返回类型:
None
- Thread-Safety:
线程安全
示例
>>> updater = BatchedStatusUpdater() >>> updater.update_sync("task-1", {"status": "completed"})
- class fish_async_task.performance.AdaptiveWorkerManager(min_workers=2, max_workers=10, scale_up_cooldown=30.0, scale_down_cooldown=60.0, cpu_threshold=0.8, queue_threshold=100)[源代码]¶
基类:
object自适应工作线程管理器
根据系统 CPU 使用率、任务队列大小和任务执行时间, 智能调整工作线程数量以优化性能和资源利用率。
- 参数:
- min_workers¶
最小工作线程数
- max_workers¶
最大工作线程数
- scale_up_cooldown¶
扩展冷却期(秒)
- scale_down_cooldown¶
缩减冷却期(秒)
- cpu_threshold¶
CPU 使用率阈值(0-1)
- queue_threshold¶
队列大小阈值
示例
>>> manager = AdaptiveWorkerManager( ... min_workers=2, ... max_workers=10, ... cpu_threshold=0.8 ... ) >>> should_scale, reason = manager.should_scale_up( ... current_workers=5, ... queue_size=100 ... ) >>> if should_scale: ... # 增加工作线程 ... pass
- __init__(min_workers=2, max_workers=10, scale_up_cooldown=30.0, scale_down_cooldown=60.0, cpu_threshold=0.8, queue_threshold=100)[源代码]¶
初始化自适应工作线程管理器
- should_scale_up(current_workers, queue_size, cpu_usage=None)[源代码]¶
判断是否应该扩展工作线程
- 参数:
- 返回:
(should_scale, reason) 元组 - should_scale: 是否应该扩展 - reason: 扩展原因的描述
- 返回类型:
- Thread-Safety:
线程安全
示例
>>> manager = AdaptiveWorkerManager() >>> should_scale, reason = manager.should_scale_up( ... current_workers=5, ... queue_size=100 ... )
- should_scale_down(current_workers, queue_size)[源代码]¶
判断是否应该缩减工作线程
- 参数:
- 返回:
(should_scale, reason) 元组 - should_scale: 是否应该缩减 - reason: 缩减原因的描述
- 返回类型:
- Thread-Safety:
线程安全
示例
>>> manager = AdaptiveWorkerManager() >>> should_scale, reason = manager.should_scale_down( ... current_workers=8, ... queue_size=5 ... )
- record_task_time(execution_time)[源代码]¶
记录任务执行时间
- 参数:
execution_time (float) -- 任务执行时间(秒)
- 返回类型:
None
- Thread-Safety:
线程安全
示例
>>> manager = AdaptiveWorkerManager() >>> manager.record_task_time(0.5) # 记录 500ms 的任务
- get_avg_task_time()[源代码]¶
获取平均任务执行时间
- 返回:
平均任务执行时间(秒),如果没有记录则返回 0.0
- 返回类型:
- Thread-Safety:
线程安全(返回近似值,但非常准确)
示例
>>> manager = AdaptiveWorkerManager() >>> manager.record_task_time(0.1) >>> manager.record_task_time(0.2) >>> manager.get_avg_task_time() 0.15
- get_cpu_usage()[源代码]¶
获取当前 CPU 使用率
- 返回:
CPU 使用率(0-1),如果 psutil 不可用则返回 None
- 返回类型:
float | None
备注
需要 psutil 可选依赖。如果不可用,返回 None。
示例
>>> manager = AdaptiveWorkerManager() >>> cpu = manager.get_cpu_usage() >>> if cpu is not None: ... print(f"CPU 使用率: {cpu:.1%}")
- get_scaling_metrics(current_workers=0, queue_size=0)[源代码]¶
获取扩展指标
- 参数:
- 返回:
包含扩展指标的字典,包括: - current_workers: 当前工作线程数 - avg_task_time: 平均任务执行时间 - cpu_usage: CPU 使用率 - last_scale_up_time: 最后一次扩展时间 - last_scale_down_time: 最后一次缩减时间 - queue_size: 队列大小
- 返回类型:
- Thread-Safety:
线程安全
示例
>>> manager = AdaptiveWorkerManager() >>> metrics = manager.get_scaling_metrics( ... current_workers=5, ... queue_size=10 ... ) >>> print(f"平均任务时间: {metrics['avg_task_time']:.3f}s")
- class fish_async_task.performance.PerformanceMetrics(max_history=1000)[源代码]¶
基类:
object性能指标收集器
- 参数:
max_history (int)
- class fish_async_task.performance.SystemHealthMonitor(logger=None)[源代码]¶
基类:
object系统健康状态监控器
- 参数:
logger (Logger)
- HEALTH_STATUS_GREEN = 'green'¶
- HEALTH_STATUS_YELLOW = 'yellow'¶
- HEALTH_STATUS_RED = 'red'¶
- register_health_check(name, value_func, threshold=None, status_warning=None, status_critical=None)[源代码]¶
注册健康检查项
- class fish_async_task.performance.TaskResourceManager(logger=None, max_tracked=10000)[源代码]¶
基类:
object任务资源管理器 - 跟踪和管理任务相关资源
- MAX_TRACKED_RESOURCES = 10000¶
- DEFAULT_CLEANUP_TIMEOUT = 2.0¶
- class fish_async_task.performance.TimeoutTaskTracker(logger=None, max_tracked=1000, task_expiry=3600)[源代码]¶
基类:
object超时任务跟踪器 - 跟踪并管理超时任务
- DEFAULT_TASK_EXPIRY = 3600¶
- class fish_async_task.performance.PriorityTaskQueue(maxsize=1000)[源代码]¶
基类:
object优先级任务队列
- 参数:
maxsize (int)
- put(task, block=True, timeout=None)[源代码]¶
添加任务到队列
- 参数:
- 抛出:
queue.Full -- 队列满且阻塞超时
- 返回类型:
None
- get(block=True, timeout=None)[源代码]¶
获取最高优先级任务
- 参数:
- 返回:
优先级任务
- 返回类型:
PrioritizedTask
- 抛出:
queue.Empty -- 队列空且阻塞超时
- class fish_async_task.performance.PriorityTaskManager(logger=None, maxsize=1000)[源代码]¶
基类:
object优先级任务管理器
- class fish_async_task.performance.TaskDependencyManager(logger=None)[源代码]¶
基类:
object任务依赖管理器
- 参数:
logger (Logger)
优先级队列¶
- class fish_async_task.performance.PriorityTaskManager(logger=None, maxsize=1000)[源代码]¶
基类:
object优先级任务管理器