API 参考

核心模块

异步任务管理器

一个纯Python实现的异步任务管理器,支持线程池和动态伸缩。

class fish_async_task.TaskManager(instance_key='default')[源代码]

基类:object

纯Python实现的异步任务管理器(线程池 + 动态伸缩)

参数:

instance_key (str)

返回类型:

TaskManager

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 对应不同的单例实例。

返回:

任务管理器实例(单例)

返回类型:

TaskManager

__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)[源代码]

提交任务到任务队列

此方法是线程安全的,可以在多个线程中并发调用。 任务会被添加到队列中,由工作线程异步执行。

参数:
  • func (Callable[[...], Any]) -- 要执行的任务函数

  • *args (Any) -- 任务函数的 positional 参数

  • block (bool) -- 如果队列已满,是否阻塞等待。默认为 False。

  • timeout (float | None) -- 阻塞等待的超时时间(秒)。仅在 block=True 时有效。 如果为 None,则无限等待。默认为 None。

  • **kwargs (Any) -- 任务函数的关键字参数

返回:

任务ID(UUID格式的字符串)

返回类型:

str

抛出:

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 资源的场景。

参数:

instance_key (str) -- 要销毁的实例键名,默认为 "default"

返回:

如果实例存在并被销毁则返回 True,否则返回 False

返回类型:

bool

示例

>>> # 销毁默认实例
>>> TaskManager.destroy_instance()
>>>
>>> # 销毁特定实例
>>> TaskManager.destroy_instance("order")

备注

  • 销毁后,再次创建相同 instance_key 的实例会重新初始化

  • 如果实例正在运行中,会先调用 shutdown() 清理资源

  • 此方法是线程安全的

exception fish_async_task.TaskQueueFullError[源代码]

基类:Exception

任务队列已满异常

TaskManager

class fish_async_task.TaskManager(instance_key='default')[源代码]

基类:object

纯Python实现的异步任务管理器(线程池 + 动态伸缩)

参数:

instance_key (str)

返回类型:

TaskManager

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 对应不同的单例实例。

返回:

任务管理器实例(单例)

返回类型:

TaskManager

__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)[源代码]

提交任务到任务队列

此方法是线程安全的,可以在多个线程中并发调用。 任务会被添加到队列中,由工作线程异步执行。

参数:
  • func (Callable[[...], Any]) -- 要执行的任务函数

  • *args (Any) -- 任务函数的 positional 参数

  • block (bool) -- 如果队列已满,是否阻塞等待。默认为 False。

  • timeout (float | None) -- 阻塞等待的超时时间(秒)。仅在 block=True 时有效。 如果为 None,则无限等待。默认为 None。

  • **kwargs (Any) -- 任务函数的关键字参数

返回:

任务ID(UUID格式的字符串)

返回类型:

str

抛出:

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 资源的场景。

参数:

instance_key (str) -- 要销毁的实例键名,默认为 "default"

返回:

如果实例存在并被销毁则返回 True,否则返回 False

返回类型:

bool

示例

>>> # 销毁默认实例
>>> TaskManager.destroy_instance()
>>>
>>> # 销毁特定实例
>>> TaskManager.destroy_instance("order")

备注

  • 销毁后,再次创建相同 instance_key 的实例会重新初始化

  • 如果实例正在运行中,会先调用 shutdown() 清理资源

  • 此方法是线程安全的

异常

exception fish_async_task.TaskQueueFullError[源代码]

任务队列已满异常

配置模块

配置管理模块

负责加载和验证任务管理器的配置项。

本模块提供了从环境变量加载配置的功能,支持以下配置项: - 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
__init__(logger)[源代码]

初始化配置加载器

参数:

logger (Logger) -- 日志记录器

load_int_config(env_key, default_value, config_name, min_value=1, max_value=None)[源代码]

加载并验证整数配置项

参数:
  • env_key (str) -- 环境变量键名

  • default_value (int) -- 默认值

  • config_name (str) -- 配置项名称(用于日志)

  • min_value (int) -- 最小值(默认为1)

  • max_value (int | None) -- 最大值(可选)

返回:

验证后的配置值

返回类型:

int

备注

如果环境变量不存在、不是有效整数或值超出允许范围, 将使用默认值并记录警告日志。

load_timeout_config(default_value)[源代码]

加载并验证任务超时配置

参数:

default_value (float | None) -- 默认超时值

返回:

任务超时时间(秒),如果为None则表示无超时限制

返回类型:

Optional[float]

load_adaptive_worker_config()[源代码]

加载自适应线程管理配置

返回:

自适应配置字典,包含以下键:
  • adaptive_worker_enabled: 是否启用自适应线程管理

  • cpu_threshold: CPU使用率阈值

  • queue_threshold_high: 扩容队列积压阈值

  • queue_threshold_low: 缩容队列空闲阈值

  • scale_up_cooldown: 扩容冷却期

  • scale_down_cooldown: 缩容冷却期

返回类型:

Dict[str, Any]

load_performance_config()[源代码]

加载性能优化配置

返回:

性能优化配置字典,包含以下键:
  • shard_count: 分片数量

  • batch_update_buffer_size: 批量更新缓冲区大小

  • batch_update_interval: 批量更新刷新间隔(秒)

  • enable_auto_cleanup: 是否启用自动清理

  • enable_batch_updates: 是否启用批量更新

  • enable_adaptive_scaling: 是否启用自适应扩展

返回类型:

Dict[str, Any]

class fish_async_task.config.HotReloadConfig(logger=None, reload_interval=60)[源代码]

基类:object

支持热重载的配置管理器

参数:
  • logger (Logger | None)

  • reload_interval (int)

__init__(logger=None, reload_interval=60)[源代码]

初始化热重载配置管理器

参数:
  • logger (Logger | None) -- 日志记录器

  • reload_interval (int) -- 重载间隔(秒)

返回类型:

None

logger: Logger
register_config(key, parser, default=None)[源代码]

注册配置项

参数:
  • key (str) -- 配置键名

  • parser (Callable[[], Any]) -- 配置解析函数

  • default (Any) -- 默认值

返回类型:

None

get(key, use_cache=True)[源代码]

获取配置值(支持热重载)

参数:
  • key (str) -- 配置键名

  • use_cache (bool) -- 是否使用缓存

返回:

配置值

返回类型:

Any

reload()[源代码]

手动触发配置重载

返回:

重载的配置数量

返回类型:

int

get_all()[源代码]

获取所有配置

返回:

所有配置项

返回类型:

Dict[str, Any]

clear_cache()[源代码]

清除配置缓存

返回类型:

None

fish_async_task.config.validate_config(min_value=None, max_value=None, default=None, allowed_values=None)[源代码]

配置验证装饰器

参数:
  • min_value (Any) -- 最小值

  • max_value (Any) -- 最大值

  • default (Any) -- 默认值

  • allowed_values (List[Any] | None) -- 允许的值列表

返回:

装饰器函数

返回类型:

Callable[[Callable[..., Any]], Callable[..., Any]]

任务状态

任务状态管理模块

负责任务状态的更新、查询和清理。

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
__init__()[源代码]

初始化读写锁

acquire_read(timeout=None)[源代码]

获取读锁

如果有写操作在等待,新的读操作会排队等待,防止写饥饿。

参数:

timeout (float | None) -- 超时时间(秒)。如果为 None,使用默认超时时间。 如果为负数或 0,立即尝试获取,不等待。

返回:

如果成功获取锁则返回 True,超时则返回 False

返回类型:

bool

抛出:

TimeoutError -- 等待超时后抛出异常

release_read()[源代码]

释放读锁

当所有读锁都被释放后,唤醒等待的写操作。

返回类型:

None

acquire_write(timeout=None)[源代码]

获取写锁

写锁是独占的,会等待所有读操作完成后才能获取。 使用写优先策略,防止写饥饿。

参数:

timeout (float | None) -- 超时时间(秒)。如果为 None,使用默认超时时间。 如果为负数或 0,立即尝试获取,不等待。

返回:

如果成功获取锁则返回 True,超时则返回 False

返回类型:

bool

抛出:

TimeoutError -- 等待超时后抛出异常

release_write()[源代码]

释放写锁

释放后,唤醒所有等待的读操作和写操作。

返回类型:

None

class fish_async_task.task_status.ReadWriteLockContext(lock, write=False, timeout=None)[源代码]

基类:object

读写锁上下文管理器

提供便捷的锁获取和释放方式。支持超时参数。

参数:
__init__(lock, write=False, timeout=None)[源代码]

初始化上下文管理器

参数:
  • lock (ReadWriteLock) -- 读写锁实例

  • write (bool) -- 是否为写操作(True=写锁,False=读锁)

  • timeout (float | None) -- 超时时间(秒),如果为 None 则使用默认超时

__enter__()[源代码]

获取锁并返回上下文管理器实例

根据 write 参数决定获取读锁或写锁。 读锁允许并发获取,写锁独占访问。

返回:

返回自身实例,用于 with 语句块

返回类型:

ReadWriteLockContext

抛出:

TimeoutError -- 获取锁超时

__exit__(*args)[源代码]

释放锁

释放之前获取的读锁或写锁。 释放后,其他读操作或写操作可以继续执行。

参数:

*args -- 接收异常信息参数(如果 with 语句中发生异常)

返回类型:

None

class fish_async_task.task_status.ShardedTaskStatusWithExpiry(shard_count, ttl)[源代码]

基类:object

分片任务状态存储(带过期时间管理)

使用分片锁减少锁竞争,每个分片内部使用优先级队列管理过期时间。 支持高并发查询和更新,以及高效的增量清理。

线程安全说明: - 每个分片有独立的锁,不同分片的操作可以并发执行 - 同一分片内的操作串行化,保证线程安全 - 清理操作支持增量清理,避免长时间阻塞

参数:
__init__(shard_count, ttl)[源代码]

初始化分片状态存储

参数:
  • shard_count (int) -- 分片数量,建议为2的幂次(8, 16, 32, 64)

  • ttl (int) -- 任务状态TTL(秒)

shards: List[Dict[str, TaskStatusDict]]
rw_locks: List[ReadWriteLock]
expiry_heaps: List[List[Tuple[float, str]]]
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

remove_status(task_id)[源代码]

移除任务状态(线程安全)

参数:

task_id (str) -- 任务ID

返回类型:

None

cleanup_expired(max_cleanup=None)[源代码]

清理过期任务(增量清理)

遍历所有分片,清理过期任务。支持增量清理,避免长时间阻塞。

参数:

max_cleanup (int | None) -- 最大清理数量,None表示清理所有过期任务

返回:

清理的任务数量

返回类型:

int

enforce_max_count(max_count)[源代码]

强制执行最大任务数量限制

当任务状态数量超过限制时,按时间顺序清理最旧的任务。 使用优化策略:优先尝试增量清理,仅在必要时获取所有锁。

参数:

max_count (int) -- 最大任务数量

返回:

清理的任务数量

返回类型:

int

get_all_statuses()[源代码]

获取所有任务状态(需要获取所有锁)

返回:

所有任务状态字典

返回类型:

Dict[str, TaskStatusDict]

clear_all()[源代码]

清空所有任务状态

返回类型:

None

get_total_count()[源代码]

获取总任务数量(不需要锁,仅用于统计)

返回:

总任务数量

返回类型:

int

resize_shards(new_shard_count)[源代码]

动态调整分片数量

此方法会重新分配所有任务状态到新的分片结构中。 由于需要重建所有数据结构,这可能是一个耗时操作, 建议在低负载时执行或在外部异步执行。

参数:

new_shard_count (int) -- 新的分片数量,必须为正整数

返回:

如果调整成功返回True,否则返回False

返回类型:

bool

警告

此操作会短暂阻塞所有状态更新和查询操作。 建议在系统初始化时设置合适的分片数量, 并尽量避免在生产环境中频繁调整。

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()

参数:
  • update_func (Callable[..., None])

  • batch_size (int)

  • flush_interval (float)

__init__(update_func, batch_size=100, flush_interval=0.1)[源代码]

初始化批量状态更新器

参数:
  • update_func (Callable[[...], None]) -- 实际的状态更新函数,接收 task_id, status, **kwargs

  • batch_size (int) -- 批量大小,达到此数量时立即刷新

  • flush_interval (float) -- 刷新间隔(秒),达到此时间间隔时刷新

update(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']) -- 任务状态

  • start_time (float | None) -- 任务开始时间(可选)

  • end_time (float | None) -- 任务结束时间(可选)

  • result (Any) -- 任务执行结果(可选)

  • error (str | None) -- 错误信息(可选)

  • submit_time (float | None) -- 任务提交时间(可选)

返回类型:

None

check_and_flush()[源代码]

手动检查并执行批量刷新

可以在外部定期调用此方法以触发定时刷新。

返回类型:

None

force_flush()[源代码]

强制刷新所有待处理的更新

返回:

刷新前队列中的更新数量

返回类型:

int

get_pending_count()[源代码]

获取当前待处理的更新数量

返回:

待处理的更新数量

返回类型:

int

shutdown()[源代码]

关闭更新器,刷新所有待处理的更新

返回:

刷新前队列中的更新数量

返回类型:

int

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

任务状态管理器

负责任务状态的存储、更新和查询。 使用分片锁和优先级队列优化性能,支持高并发操作。 支持批量状态更新,减少锁获取次数,提升写入性能。

线程安全说明: - 使用分片锁,不同分片的操作可以并发执行 - 同一分片内的操作串行化,保证线程安全 - 清理操作支持增量清理,避免长时间阻塞 - 批量更新器内部使用队列和锁,保证线程安全

参数:
  • logger (Logger)

  • task_status_ttl (int)

  • max_task_status_count (int)

  • shard_count (int | None)

  • batch_size (int | None)

  • batch_flush_interval (float | None)

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)[源代码]

初始化任务状态管理器

参数:
  • logger (Logger) -- 日志记录器

  • task_status_ttl (int) -- 任务状态TTL(秒)

  • max_task_status_count (int) -- 最大任务状态数量

  • shard_count (int | None) -- 分片数量,默认从环境变量 TASK_STATUS_SHARD_COUNT 读取,或使用16

  • batch_size (int | None) -- 批量更新大小(可选,默认使用类常量)

  • batch_flush_interval (float | None) -- 批量刷新间隔(秒)(可选,默认使用类常量)

enable_batch_update(enabled=True)[源代码]

启用或禁用批量更新

参数:

enabled (bool) -- 是否启用批量更新

返回类型:

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. 如果任务状态数量超过限制,清理最旧的任务

返回:

清理的任务数量

返回类型:

int

resize_shards(new_shard_count)[源代码]

动态调整分片数量

此方法会重新分配所有任务状态到新的分片结构中。 由于需要重建所有数据结构,这可能是一个耗时操作, 建议在低负载时执行或在外部异步执行。

参数:

new_shard_count (int) -- 新的分片数量,必须为正整数

返回:

如果调整成功返回True,否则返回False

返回类型:

bool

警告

此操作会短暂阻塞所有状态更新和查询操作。 建议在系统初始化时设置合适的分片数量, 并尽量避免在生产环境中频繁调整。

shutdown()[源代码]

关闭状态管理器,刷新所有待处理的批量更新

返回:

刷新前待处理的更新数量

返回类型:

int

get_pending_update_count()[源代码]

获取当前待处理的更新数量

返回:

待处理的更新数量

返回类型:

int

类型定义

类型定义模块

定义任务管理器相关的类型别名和类型定义。

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。

status: Literal['pending', 'running', 'completed', 'failed'] | None
submit_time: float | None
start_time: float | None
end_time: float | None
result: Any | None
error: str | None
worker_id: str | None
class fish_async_task.types.ShardedTaskStatusDict[源代码]

基类:TypedDict

分片任务状态字典,包含分片索引信息

shard_index: int
task_status: TaskStatusDict
class fish_async_task.types.BatchedUpdate[源代码]

基类:TypedDict

批量更新项

task_id: str
status: TaskStatusDict
class fish_async_task.types.ScalingMetrics[源代码]

基类:TypedDict

自适应扩展指标

current_workers: int
avg_task_time: float
cpu_usage: float | None
last_scale_up_time: float
last_scale_down_time: float
queue_size: int

工作线程

工作线程模块

负责工作线程的创建、管理和任务执行。 支持自适应线程管理,根据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: 缩容冷却期(秒),避免频繁缩容

参数:
  • min_workers (int)

  • max_workers (int)

  • cpu_threshold (float)

  • queue_threshold_high (int)

  • queue_threshold_low (int)

  • scale_up_cooldown (float)

  • scale_down_cooldown (float)

  • use_cpu_monitoring (bool)

__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使用率低于阈值

参数:
  • current_workers (int) -- 当前工作线程数

  • queue_size (int) -- 当前队列积压任务数

  • cpu_usage (float | None) -- 当前CPU使用率(0.0-1.0),可选

返回:

如果应该扩容则返回True

返回类型:

bool

should_scale_down(current_workers, queue_size, idle_time=None)[源代码]

判断是否应该缩减线程

缩容条件(需全部满足): 1. 当前线程数大于最小限制 2. 距离上次缩容超过冷却期 3. 队列为空且空闲时间超过阈值

参数:
  • current_workers (int) -- 当前工作线程数

  • queue_size (int) -- 当前队列积压任务数

  • idle_time (float | None) -- 队列空闲时间(秒),可选

返回:

如果应该缩容则返回True

返回类型:

bool

record_task_time(task_time)[源代码]

记录任务执行时间

参数:

task_time (float) -- 任务执行时间(秒)

返回类型:

None

get_avg_task_time()[源代码]

获取平均任务执行时间

返回:

平均任务执行时间(秒),如果没有数据则返回0

返回类型:

float

get_cpu_usage()[源代码]

获取当前CPU使用率

返回:

CPU使用率(0.0-1.0),如果CPU监控不可用则返回None

返回类型:

float

get_stats()[源代码]

获取管理器统计信息

返回:

统计信息字典

返回类型:

Dict[str, Any]

class fish_async_task.worker.CPUMonitor(sample_interval=0.5, sample_count=3)[源代码]

基类:object

CPU使用率监控器

使用psutil获取系统CPU使用率,支持采样平均。

参数:
  • sample_interval (float)

  • sample_count (int)

__init__(sample_interval=0.5, sample_count=3)[源代码]

初始化CPU监控器

参数:
  • sample_interval (float) -- 每次采样的时间间隔(秒)

  • sample_count (int) -- 采样次数,用于计算平均CPU使用率

get_cpu_usage()[源代码]

获取CPU使用率

如果psutil不可用,返回None。

返回:

CPU使用率(0.0-1.0),如果不可用则返回None

返回类型:

float

get_cpu_count()[源代码]

获取CPU核心数

如果psutil不可用或返回None,返回1。

返回:

CPU核心数

返回类型:

int

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)[源代码]

初始化任务执行器

参数:
  • logger (Logger) -- 日志记录器

  • task_timeout_getter (Callable[[], float | None]) -- 获取任务超时时间的函数(支持动态获取)

  • update_status_func (Callable[[...], None]) -- 状态更新函数

  • batch_size (int | None) -- 批量更新大小(可选,默认使用类常量)

  • batch_flush_interval (float | None) -- 批量刷新间隔(秒)(可选,默认使用类常量)

flush_pending_updates()[源代码]

强制刷新所有待处理的批量更新

返回:

刷新前队列中的更新数量

返回类型:

int

get_pending_update_count()[源代码]

获取当前待处理的批量更新数量

返回:

待处理的更新数量

返回类型:

int

execute_task(task)[源代码]

执行任务

参数:

task (Tuple[str, Callable[[...], Any], Tuple[Any, ...], Dict[str, Any]]) -- 任务元组,包含(task_id, func, args, kwargs)

返回类型:

None

add_cleanup_callback(callback)[源代码]

添加清理回调函数

当任务超时时,会调用所有注册的清理回调函数。

参数:

callback (Callable[[str], None]) -- 清理回调函数,接收任务ID作为参数

返回类型:

None

cleanup_timed_out_tasks(max_cleanup=100)[源代码]

清理超时的任务线程引用

注意:由于Python线程无法被强制终止,此方法主要用于清理跟踪信息, 实际的任务线程仍会在后台运行直到任务完成或进程退出。

参数:

max_cleanup (int) -- 最大清理数量

返回:

清理的任务数量

返回类型:

int

get_timed_out_task_count()[源代码]

获取当前跟踪的超时任务数量

返回:

超时任务数量

返回类型:

int

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保护下进行 - 使用退出事件确认机制确保线程完全退出

参数:
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]]]) -- 任务队列

  • worker_threads (List[Thread]) -- 工作线程列表

  • 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

scale_up_workers_if_needed()[源代码]

根据队列大小和CPU使用率动态扩展工作线程

当队列中的任务数量超过当前线程数时,自动创建新线程。 线程数量不会超过 max_workers 限制。 使用自适应策略,考虑CPU使用率和队列积压情况。

返回类型:

None

record_task_time(task_time)[源代码]

记录任务执行时间(用于自适应管理)

参数:

task_time (float) -- 任务执行时间(秒)

返回类型:

None

get_idle_time()[源代码]

获取当前队列空闲时间

返回:

队列空闲时间(秒),如果当前不空闲则返回None

返回类型:

float

should_scale_down()[源代码]

判断是否应该缩减线程(使用自适应策略)

返回:

如果应该缩容则返回True

返回类型:

bool

get_adaptive_stats()[源代码]

获取自适应管理器的统计信息

返回:

统计信息字典,如果自适应管理未启用则返回None

返回类型:

Dict[str, Any]

send_shutdown_signals()[源代码]

向所有工作线程发送退出信号

通过向任务队列中放入 None 值来通知工作线程退出。 如果队列已满,会尝试等待并重试。

返回类型:

None

wait_for_threads_exit(join_timeout)[源代码]

等待所有工作线程退出

在锁内创建线程副本以避免竞态条件,然后逐个等待线程退出。 使用退出事件确认机制,确保线程完全退出。 如果线程在超时时间内未退出,会记录警告但继续执行。

参数:

join_timeout (int) -- 线程join超时时间(秒)

返回类型:

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)[源代码]

更新任务状态(线程安全)

参数:
抛出:

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()[源代码]

获取当前任务数量

返回:

任务状态字典中的任务数量

返回类型:

int

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()[源代码]

获取所有任务状态(需要获取所有锁)

警告

此方法会按顺序获取所有分片的锁,可能阻塞较长时间。 仅在必要时使用(如关闭、统计)。

返回:

所有任务状态的字典

返回类型:

Dict[str, TaskStatusDict]

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
clear_all()[源代码]

清空所有任务状态(需要获取所有锁)

警告

此方法会按顺序获取所有分片的锁。

Performance:

O(n) 时间复杂度,n 为任务总数 线程安全:按顺序获取所有锁

示例

>>> store = ShardedTaskStatus()
>>> store.update_status("task-1", {"status": "completed"})
>>> store.clear_all()
>>> store.get_task_count()
0
返回类型:

None

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

参数:
返回类型:

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)[源代码]

清理过期任务(增量清理)

参数:

max_cleanup (int | None) -- 最大清理数量,None 表示清理所有过期任务 默认 None

返回:

清理的任务数量

返回类型:

int

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)。

参数:

max_count (int) -- 最大任务数量

返回:

删除的任务数量

返回类型:

int

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()[源代码]

获取当前任务数量

返回:

任务状态字典中的任务数量

返回类型:

int

Performance:

O(1) 时间复杂度

Thread-Safety:

线程安全(返回近似值,但非常准确)

示例

>>> store = TaskStatusWithExpiry()
>>> store.add_task("task-123", {"end_time": time.time()})
>>> store.get_task_count()
1
get_all_statuses()[源代码]

获取所有任务状态

返回:

所有任务状态的字典

返回类型:

Dict[str, TaskStatusDict]

备注

此方法会锁定优先级队列,避免在清理期间调用

class fish_async_task.performance.BatchedStatusUpdater(buffer_size=100, flush_interval=1.0, underlying_store=None)[源代码]

基类:object

批量状态更新器

将多个状态更新缓存到缓冲区,然后批量刷新到底层存储, 减少锁竞争和提高吞吐量。

参数:
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,会自动触发刷新。

参数:
抛出:

RuntimeError -- 如果更新器已关闭

返回类型:

None

Thread-Safety:

线程安全

示例

>>> updater = BatchedStatusUpdater()
>>> updater.queue_update("task-1", {"status": "running"})
flush()[源代码]

手动刷新缓冲区到底层存储

返回:

刷新的任务数量

返回类型:

int

Thread-Safety:

线程安全

示例

>>> updater = BatchedStatusUpdater()
>>> updater.queue_update("task-1", {"status": "running"})
>>> flushed = updater.flush()
>>> print(f"刷新了 {flushed} 个任务")
update_sync(task_id, status)[源代码]

同步更新(立即写入底层存储,不经过缓冲区)

用于需要立即更新的场景,例如关键状态变更。

参数:
抛出:

RuntimeError -- 如果更新器已关闭

返回类型:

None

Thread-Safety:

线程安全

示例

>>> updater = BatchedStatusUpdater()
>>> updater.update_sync("task-1", {"status": "completed"})
get_buffer_length()[源代码]

获取当前缓冲区长度

返回:

缓冲区中的任务数量

返回类型:

int

Thread-Safety:

线程安全(返回近似值,但非常准确)

示例

>>> updater = BatchedStatusUpdater()
>>> updater.queue_update("task-1", {"status": "running"})
>>> updater.get_buffer_length()
1
close()[源代码]

关闭更新器并刷新所有待处理的更新

关闭后不再接受新更新,但会确保所有已排队的更新都被刷新。

Thread-Safety:

线程安全

示例

>>> updater = BatchedStatusUpdater()
>>> updater.queue_update("task-1", {"status": "running"})
>>> updater.close()  # 刷新所有待处理更新并关闭
返回类型:

None

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 (int)

  • max_workers (int)

  • scale_up_cooldown (float)

  • scale_down_cooldown (float)

  • cpu_threshold (float)

  • queue_threshold (int)

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)[源代码]

初始化自适应工作线程管理器

参数:
  • min_workers (int) -- 最小工作线程数,默认 2

  • max_workers (int) -- 最大工作线程数,默认 10

  • scale_up_cooldown (float) -- 扩展冷却期(秒),默认 30.0

  • scale_down_cooldown (float) -- 缩减冷却期(秒),默认 60.0

  • cpu_threshold (float) -- CPU 使用率阈值(0-1),默认 0.8

  • queue_threshold (int) -- 队列大小阈值,默认 100

抛出:

ValueError -- 如果参数无效

返回类型:

None

should_scale_up(current_workers, queue_size, cpu_usage=None)[源代码]

判断是否应该扩展工作线程

参数:
  • current_workers (int) -- 当前工作线程数

  • queue_size (int) -- 当前队列大小

  • cpu_usage (float | None) -- 当前 CPU 使用率(0-1),如果为 None 则自动获取

返回:

(should_scale, reason) 元组 - should_scale: 是否应该扩展 - reason: 扩展原因的描述

返回类型:

tuple[bool, str]

Thread-Safety:

线程安全

示例

>>> manager = AdaptiveWorkerManager()
>>> should_scale, reason = manager.should_scale_up(
...     current_workers=5,
...     queue_size=100
... )
should_scale_down(current_workers, queue_size)[源代码]

判断是否应该缩减工作线程

参数:
  • current_workers (int) -- 当前工作线程数

  • queue_size (int) -- 当前队列大小

返回:

(should_scale, reason) 元组 - should_scale: 是否应该缩减 - reason: 缩减原因的描述

返回类型:

tuple[bool, str]

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

返回类型:

float

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 (int) -- 当前工作线程数

  • queue_size (int) -- 当前队列大小

返回:

包含扩展指标的字典,包括: - current_workers: 当前工作线程数 - avg_task_time: 平均任务执行时间 - cpu_usage: CPU 使用率 - last_scale_up_time: 最后一次扩展时间 - last_scale_down_time: 最后一次缩减时间 - queue_size: 队列大小

返回类型:

Dict[str, Any]

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)

__init__(max_history=1000)[源代码]

初始化性能指标收集器

参数:

max_history (int) -- 最大历史记录数

record_task_submitted()[源代码]

记录任务提交

返回类型:

None

record_task_completed(execution_time, queue_wait_time=0.0)[源代码]

记录任务完成

参数:
  • execution_time (float) -- 任务执行时间(秒)

  • queue_wait_time (float) -- 任务在队列中等待时间(秒)

返回类型:

None

record_task_failed(execution_time, queue_wait_time=0.0)[源代码]

记录任务失败

参数:
  • execution_time (float) -- 任务执行时间(秒)

  • queue_wait_time (float) -- 任务在队列中等待时间(秒)

返回类型:

None

record_task_cancelled()[源代码]

记录任务取消

返回类型:

None

record_cleanup(cleanup_time, cleaned_count)[源代码]

记录清理操作

参数:
  • cleanup_time (float) -- 清理耗时(秒)

  • cleaned_count (int) -- 清理的任务数量

返回类型:

None

get_metrics()[源代码]

获取当前性能指标

返回:

性能指标字典

返回类型:

Dict[str, Any]

get_percentiles(percentiles=None)[源代码]

获取执行时间百分位数

参数:

percentiles (List[int]) -- 百分位列表,默认 [50, 75, 90, 95, 99]

返回:

百分位数据

返回类型:

Dict[str, float]

reset()[源代码]

重置所有指标

返回类型:

None

class fish_async_task.performance.SystemHealthMonitor(logger=None)[源代码]

基类:object

系统健康状态监控器

参数:

logger (Logger)

HEALTH_STATUS_GREEN = 'green'
HEALTH_STATUS_YELLOW = 'yellow'
HEALTH_STATUS_RED = 'red'
__init__(logger=None)[源代码]

初始化系统健康监控器

参数:

logger (Logger) -- 日志记录器

register_health_check(name, value_func, threshold=None, status_warning=None, status_critical=None)[源代码]

注册健康检查项

参数:
  • name (str) -- 检查项名称

  • value_func (callable) -- 获取值的函数,接受metrics字典

  • threshold (float) -- 默认阈值

  • status_warning (float) -- 警告状态阈值

  • status_critical (float) -- 严重状态阈值

返回类型:

None

update_health_status(metrics)[源代码]

更新健康状态

参数:

metrics (Dict[str, Any]) -- 性能指标字典

返回:

健康状态信息

返回类型:

Dict[str, Any]

get_health_status()[源代码]

获取当前健康状态

返回:

健康状态信息

返回类型:

Dict[str, Any]

class fish_async_task.performance.TaskResourceManager(logger=None, max_tracked=10000)[源代码]

基类:object

任务资源管理器 - 跟踪和管理任务相关资源

参数:
MAX_TRACKED_RESOURCES = 10000
DEFAULT_CLEANUP_TIMEOUT = 2.0
__init__(logger=None, max_tracked=10000)[源代码]

初始化任务资源管理器

参数:
  • logger (Logger) -- 日志记录器

  • max_tracked (int) -- 最大跟踪资源数

start()[源代码]

启动资源清理线程

返回类型:

None

stop(timeout=2.0)[源代码]

停止资源清理线程

参数:

timeout (float) -- 等待超时时间(秒)

返回类型:

None

register_resource(task_id, resource_id, resource, cleanup_func=None)[源代码]

注册任务资源

参数:
  • task_id (str) -- 任务ID

  • resource_id (str) -- 资源唯一标识

  • resource (Any) -- 资源对象

  • cleanup_func (Callable[[], None] | None) -- 资源清理函数(可选)

返回类型:

None

unregister_resource(resource_id)[源代码]

注销资源

参数:

resource_id (str) -- 资源唯一标识

返回:

是否成功注销

返回类型:

bool

register_task(task_id)[源代码]

注册任务(用于跟踪)

参数:

task_id (str) -- 任务ID

返回类型:

None

cleanup_task_resources(task_id, timeout=2.0)[源代码]

清理任务的所有资源

参数:
  • task_id (str) -- 任务ID

  • timeout (float) -- 等待超时时间(秒)

返回:

清理的资源数量

返回类型:

int

force_cleanup_task(task_id)[源代码]

强制清理任务资源(不使用线程)

参数:

task_id (str) -- 任务ID

返回:

清理的资源数量

返回类型:

int

get_resource_count()[源代码]

获取当前跟踪的资源数量

返回:

资源数量

返回类型:

int

get_task_resource_count(task_id)[源代码]

获取任务的资源数量

参数:

task_id (str) -- 任务ID

返回:

资源数量

返回类型:

int

get_stats()[源代码]

获取资源管理统计信息

返回:

统计信息

返回类型:

Dict[str, Any]

class fish_async_task.performance.TimeoutTaskTracker(logger=None, max_tracked=1000, task_expiry=3600)[源代码]

基类:object

超时任务跟踪器 - 跟踪并管理超时任务

参数:
DEFAULT_TASK_EXPIRY = 3600
__init__(logger=None, max_tracked=1000, task_expiry=3600)[源代码]

初始化超时任务跟踪器

参数:
  • logger (Logger) -- 日志记录器

  • max_tracked (int) -- 最大跟踪任务数

  • task_expiry (int) -- 任务信息过期时间(秒)

track_timeout_task(task_id, thread, submit_time=None)[源代码]

跟踪超时任务

参数:
  • task_id (str) -- 任务ID

  • thread (Thread) -- 任务执行线程

  • submit_time (float) -- 任务提交时间

返回类型:

None

untrack_task(task_id)[源代码]

取消跟踪任务

参数:

task_id (str) -- 任务ID

返回:

是否成功取消跟踪

返回类型:

bool

get_tracked_count()[源代码]

获取跟踪的任务数量

返回:

任务数量

返回类型:

int

get_stats()[源代码]

获取跟踪器统计信息

返回:

统计信息

返回类型:

Dict[str, Any]

class fish_async_task.performance.PriorityTaskQueue(maxsize=1000)[源代码]

基类:object

优先级任务队列

参数:

maxsize (int)

__init__(maxsize=1000)[源代码]

初始化优先级任务队列

参数:

maxsize (int) -- 队列最大容量,0表示无限制

put(task, block=True, timeout=None)[源代码]

添加任务到队列

参数:
  • task (PrioritizedTask) -- 优先级任务

  • block (bool) -- 是否阻塞

  • timeout (float | None) -- 超时时间

抛出:

queue.Full -- 队列满且阻塞超时

返回类型:

None

get(block=True, timeout=None)[源代码]

获取最高优先级任务

参数:
  • block (bool) -- 是否阻塞

  • timeout (float | None) -- 超时时间

返回:

优先级任务

返回类型:

PrioritizedTask

抛出:

queue.Empty -- 队列空且阻塞超时

task_done()[源代码]

标记任务完成

返回类型:

None

qsize()[源代码]

获取队列大小

返回:

队列大小

返回类型:

int

empty()[源代码]

检查队列是否为空

返回:

队列是否为空

返回类型:

bool

full()[源代码]

检查队列是否已满

返回:

队列是否已满

返回类型:

bool

contains(task_id)[源代码]

检查任务是否在队列中

参数:

task_id (str) -- 任务ID

返回:

任务是否在队列中

返回类型:

bool

remove(task_id)[源代码]

从队列中移除任务

参数:

task_id (str) -- 任务ID

返回:

是否成功移除

返回类型:

bool

clear()[源代码]

清空队列

返回:

清空的任务数量

返回类型:

int

class fish_async_task.performance.PriorityTaskManager(logger=None, maxsize=1000)[源代码]

基类:object

优先级任务管理器

参数:
__init__(logger=None, maxsize=1000)[源代码]

初始化优先级任务管理器

参数:
  • logger (Logger) -- 日志记录器

  • maxsize (int) -- 队列最大容量

submit_task(func, priority=5, task_id=None, args=(), kwargs=None)[源代码]

提交带优先级的任务

参数:
  • func (Callable[[...], Any]) -- 任务函数

  • priority (int) -- 优先级(数字越小优先级越高)

  • task_id (str | None) -- 任务ID(可选,默认自动生成)

  • args (tuple) -- 位置参数

  • kwargs (dict) -- 关键字参数

返回:

任务ID

返回类型:

str

get_task(timeout=None)[源代码]

获取最高优先级任务

参数:

timeout (float | None) -- 超时时间

返回:

优先级任务

返回类型:

Optional[PrioritizedTask]

cancel_task(task_id)[源代码]

取消任务

参数:

task_id (str) -- 任务ID

返回:

是否成功取消

返回类型:

bool

get_queue_size()[源代码]

获取队列大小

返回:

队列大小

返回类型:

int

is_empty()[源代码]

检查队列是否为空

返回:

队列是否为空

返回类型:

bool

get_pending_count()[源代码]

获取待处理任务数量

返回:

待处理任务数量

返回类型:

int

class fish_async_task.performance.TaskDependencyManager(logger=None)[源代码]

基类:object

任务依赖管理器

参数:

logger (Logger)

__init__(logger=None)[源代码]

初始化任务依赖管理器

参数:

logger (Logger) -- 日志记录器

add_dependency(task_id, depends_on)[源代码]

添加任务依赖

参数:
  • task_id (str) -- 任务ID

  • depends_on (List[str]) -- 依赖的任务ID列表

返回类型:

None

mark_completed(task_id)[源代码]

标记任务完成

参数:

task_id (str) -- 任务ID

返回类型:

None

mark_failed(task_id)[源代码]

标记任务失败

参数:

task_id (str) -- 任务ID

返回类型:

None

is_ready(task_id)[源代码]

检查任务是否就绪(所有依赖已满足)

参数:

task_id (str) -- 任务ID

返回:

任务是否就绪

返回类型:

bool

get_ready_tasks()[源代码]

获取所有就绪的任务

返回:

就绪的任务ID列表

返回类型:

List[str]

has_circular_dependency(task_id)[源代码]

检查是否存在循环依赖

参数:

task_id (str) -- 任务ID

返回:

是否存在循环依赖

返回类型:

bool

clear()[源代码]

清除所有依赖信息

返回类型:

None

get_stats()[源代码]

获取依赖管理器统计信息

返回:

统计信息

返回类型:

Dict[str, Any]

优先级队列

class fish_async_task.performance.PriorityTaskManager(logger=None, maxsize=1000)[源代码]

基类:object

优先级任务管理器

参数:
__init__(logger=None, maxsize=1000)[源代码]

初始化优先级任务管理器

参数:
  • logger (Logger) -- 日志记录器

  • maxsize (int) -- 队列最大容量

submit_task(func, priority=5, task_id=None, args=(), kwargs=None)[源代码]

提交带优先级的任务

参数:
  • func (Callable[[...], Any]) -- 任务函数

  • priority (int) -- 优先级(数字越小优先级越高)

  • task_id (str | None) -- 任务ID(可选,默认自动生成)

  • args (tuple) -- 位置参数

  • kwargs (dict) -- 关键字参数

返回:

任务ID

返回类型:

str

get_task(timeout=None)[源代码]

获取最高优先级任务

参数:

timeout (float | None) -- 超时时间

返回:

优先级任务

返回类型:

Optional[PrioritizedTask]

cancel_task(task_id)[源代码]

取消任务

参数:

task_id (str) -- 任务ID

返回:

是否成功取消

返回类型:

bool

get_queue_size()[源代码]

获取队列大小

返回:

队列大小

返回类型:

int

is_empty()[源代码]

检查队列是否为空

返回:

队列是否为空

返回类型:

bool

get_pending_count()[源代码]

获取待处理任务数量

返回:

待处理任务数量

返回类型:

int

自适应伸缩

资源监控

分片状态管理

批量更新