Skip to main content

A flexible task scheduling and server management library

Project description

Task Weaver

Task Weaver 是一个强大的分布式任务管理库,专门用于处理GPU和API资源的任务调度和执行。它提供了灵活的任务队列管理、服务器资源分配和任务执行监控功能。

特性

  • 支持多种任务类型(GPU任务和API任务)
  • 智能的服务器资源分配
  • 任务优先级管理
  • 实时任务状态监控
  • 可靠的错误处理和恢复机制
  • 支持异步并发执行
  • 服务器健康检查和自动重连

安装

使用 pip 安装:

pip install task-weaver

快速开始

基本使用

from task_weaver import (
    task_manager,
    server_manager,
    task_catalog,
    TaskPriority,
    ResourceType
)

# 1. 注册服务器
server_manager.register_server(
    ip="http://192.168.1.100:8000",
    server_name="gpu-1",
    description="GPU Server 1",
    tier=1,
    available_task_types=["image_generation"],
    server_type=ResourceType.GPU,
    max_concurrency=2  # 可选:单服务器并发槽位,默认 1
)

# 2. 定义并注册任务处理器
async def process_image(server, task_info, **params):
    # 实现你的任务处理逻辑
    result = await your_processing_logic(params)
    return result

task_catalog.add_task(
    task_type="image_generation",
    executor=process_image,
    required_resources=ResourceType.GPU,
    max_concurrency=10  # 可选:任务类型最大并发,None 为不限制
)

# 3. 创建并执行任务
async def main():
    task = await task_manager.create_task(
        task_type="image_generation",
        params={
            "prompt": "A beautiful sunset",
            "steps": 30
        },
        priority=TaskPriority.HIGH
    )
    
    # 获取任务状态
    task_info = task_manager.get_task_info(task.task_info.task_id)

服务器管理

# 添加服务器到运行队列
success, message = server_manager.add_running_server(ip="http://192.168.1.100:8000")

# 从运行队列中移除服务器
success, message = server_manager.remove_running_server(ip="http://192.168.1.100:8000")

# 检查服务器状态
server = server_manager.get_server_by_identifier(ip="http://192.168.1.100:8000")

高级功能

任务优先级

Task Weaver 支持三种任务优先级:

from task_weaver import TaskPriority

# 创建不同优先级的任务
high_priority_task = await task_manager.create_task(
    task_type="image_generation",
    params=params,
    priority=TaskPriority.HIGH
)

medium_priority_task = await task_manager.create_task(
    task_type="image_generation",
    params=params,
    priority=TaskPriority.MEDIUM
)

low_priority_task = await task_manager.create_task(
    task_type="image_generation",
    params=params,
    priority=TaskPriority.LOW
)

任务状态监控

from task_weaver import TaskStatus

# 获取任务信息
task_info = task_manager.get_task_info(task_id)

# 检查任务状态
if task_info.status == TaskStatus.FINISH:
    print("任务完成:", task_info.result)
elif task_info.status == TaskStatus.FAIL:
    print("任务失败:", task_info.error)

同一任务类型下按子业务限流

task_catalog.add_task_definition(
    task_name="Cloud Image",
    task_type="cloud_image_v1",
    executor=cloud_image_executor,
    required_resource=ResourceType.API,
    # 任务类型总并发(可选)
    max_concurrency=20,
    # 子业务并发规则(可选)
    subtask_key="provider",
    subtask_concurrency={
        "gemini": 10,
        "baidu": 2,
        "*": 5,  # 未配置的 provider 走兜底并发
    },
)

# 提交任务时带上 provider 参数
task = await task_manager.create_task(
    task_type="cloud_image_v1",
    priority=TaskPriority.MEDIUM,
    provider="gemini",
    prompt="a cat in watercolor style",
)
await task_manager.add_task(task)

配置要求

  • Python 3.10 或更高版本
  • 异步支持 (asyncio)
  • 网络连接(用于服务器通信)

许可证

本项目采用 MIT 许可证。详见 LICENSE 文件。

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

task_weaver_hprt-0.3.9.tar.gz (22.5 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

task_weaver_hprt-0.3.9-py3-none-any.whl (28.0 kB view details)

Uploaded Python 3

File details

Details for the file task_weaver_hprt-0.3.9.tar.gz.

File metadata

  • Download URL: task_weaver_hprt-0.3.9.tar.gz
  • Upload date:
  • Size: 22.5 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.10.20

File hashes

Hashes for task_weaver_hprt-0.3.9.tar.gz
Algorithm Hash digest
SHA256 c2967796b88fa5d93d426a913165c9553a8a217cb7c3d1f8fe282da81c05667d
MD5 05ac8623e1daf535564964f528628d1b
BLAKE2b-256 ed97bc27e77ba3f92725ab649204362d291ce5ee183b77b3bb88058616d03390

See more details on using hashes here.

File details

Details for the file task_weaver_hprt-0.3.9-py3-none-any.whl.

File metadata

File hashes

Hashes for task_weaver_hprt-0.3.9-py3-none-any.whl
Algorithm Hash digest
SHA256 6f324d52abfbc5c5270b4ce4998965e2d330be40c64c56848255e43e4a2fad4d
MD5 ed238aaad43d6e7e345072abeafefc7b
BLAKE2b-256 ddf8c4880e833b1778e465a1bf7c6e04c24b100f1765972b7d2b1642b2ab7bef

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page