Spaces:
Running
on
Zero
Running
on
Zero
File size: 6,623 Bytes
c8df52d |
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 |
import threading
import time
import uuid
from typing import Dict, List, Optional, Callable, Any
from dataclasses import dataclass
from enum import Enum
class TaskStatus(Enum):
QUEUED = "queued"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
@dataclass
class Task:
task_id: str
status: TaskStatus
created_at: float
started_at: Optional[float] = None
completed_at: Optional[float] = None
result: Optional[Any] = None
error: Optional[str] = None
function: Optional[Callable] = None
args: tuple = ()
kwargs: dict = None
def __post_init__(self):
if self.kwargs is None:
self.kwargs = {}
class ConcurrencyManager:
def __init__(self, max_concurrent: int = 2):
"""
并发控制管理器
Args:
max_concurrent: 最大并发数量
"""
self.max_concurrent = max_concurrent
self.running_tasks: Dict[str, Task] = {}
self.queued_tasks: List[Task] = []
self.completed_tasks: Dict[str, Task] = {}
self.lock = threading.RLock()
self.worker_threads: List[threading.Thread] = []
self.shutdown_event = threading.Event()
# 启动工作线程
self._start_workers()
def _start_workers(self):
"""启动工作线程"""
for i in range(self.max_concurrent):
worker = threading.Thread(target=self._worker_loop, daemon=True)
worker.start()
self.worker_threads.append(worker)
def _worker_loop(self):
"""工作线程主循环"""
while not self.shutdown_event.is_set():
try:
task = self._get_next_task()
if task:
self._execute_task(task)
else:
# 没有任务时短暂休眠
time.sleep(0.1)
except Exception as e:
print(f"Worker thread error: {e}")
time.sleep(1)
def _get_next_task(self) -> Optional[Task]:
"""获取下一个要执行的任务"""
with self.lock:
if self.queued_tasks:
return self.queued_tasks.pop(0)
return None
def _execute_task(self, task: Task):
"""执行任务"""
try:
with self.lock:
task.status = TaskStatus.RUNNING
task.started_at = time.time()
self.running_tasks[task.task_id] = task
# 执行任务
if task.function:
result = task.function(*task.args, **task.kwargs)
task.result = result
# 标记完成
with self.lock:
task.status = TaskStatus.COMPLETED
task.completed_at = time.time()
self.completed_tasks[task.task_id] = task
if task.task_id in self.running_tasks:
del self.running_tasks[task.task_id]
except Exception as e:
# 标记失败
with self.lock:
task.status = TaskStatus.FAILED
task.completed_at = time.time()
task.error = str(e)
self.completed_tasks[task.task_id] = task
if task.task_id in self.running_tasks:
del self.running_tasks[task.task_id]
def submit_task(self, func: Callable, *args, **kwargs) -> str:
"""
提交任务
Args:
func: 要执行的函数
*args: 函数参数
**kwargs: 函数关键字参数
Returns:
task_id: 任务ID
"""
task_id = str(uuid.uuid4())
task = Task(
task_id=task_id,
status=TaskStatus.QUEUED,
created_at=time.time(),
function=func,
args=args,
kwargs=kwargs
)
with self.lock:
self.queued_tasks.append(task)
return task_id
def get_task_status(self, task_id: str) -> Optional[Task]:
"""获取任务状态"""
with self.lock:
if task_id in self.running_tasks:
return self.running_tasks[task_id]
elif task_id in self.completed_tasks:
return self.completed_tasks[task_id]
else:
# 检查队列中的任务
for task in self.queued_tasks:
if task.task_id == task_id:
return task
return None
def get_queue_status(self) -> Dict[str, Any]:
"""获取队列状态"""
with self.lock:
return {
"max_concurrent": self.max_concurrent,
"running_count": len(self.running_tasks),
"queued_count": len(self.queued_tasks),
"completed_count": len(self.completed_tasks),
"running_tasks": [task.task_id for task in self.running_tasks.values()],
"queued_tasks": [task.task_id for task in self.queued_tasks],
}
def wait_for_task(self, task_id: str, timeout: Optional[float] = None) -> Task:
"""
等待任务完成
Args:
task_id: 任务ID
timeout: 超时时间(秒),None表示无限等待
Returns:
Task: 完成的任务
"""
start_time = time.time()
while True:
task = self.get_task_status(task_id)
if task and task.status in [TaskStatus.COMPLETED, TaskStatus.FAILED]:
return task
if timeout and (time.time() - start_time) > timeout:
raise TimeoutError(f"Task {task_id} timed out after {timeout} seconds")
time.sleep(0.1)
def cleanup_old_tasks(self, max_age_hours: int = 24):
"""清理旧任务"""
current_time = time.time()
max_age_seconds = max_age_hours * 3600
with self.lock:
# 清理已完成的任务
old_tasks = [
task_id for task_id, task in self.completed_tasks.items()
if current_time - task.completed_at > max_age_seconds
]
for task_id in old_tasks:
del self.completed_tasks[task_id]
def shutdown(self):
"""关闭管理器"""
self.shutdown_event.set()
for worker in self.worker_threads:
worker.join(timeout=5)
|