本文目录导读:

我来介绍几种Python动态调度空闲同步节点的方法:
基于队列的任务调度
import queue
import threading
import time
import random
from typing import List, Dict, Any
from dataclasses import dataclass
from enum import Enum
class NodeStatus(Enum):
IDLE = "idle"
BUSY = "busy"
OFFLINE = "offline"
@dataclass
class Node:
id: str
status: NodeStatus
capacity: int
current_load: int = 0
last_heartbeat: float = 0
class DynamicScheduler:
def __init__(self):
self.task_queue = queue.Queue()
self.nodes: Dict[str, Node] = {}
self.lock = threading.Lock()
def register_node(self, node_id: str, capacity: int):
"""注册新节点"""
with self.lock:
self.nodes[node_id] = Node(
id=node_id,
status=NodeStatus.IDLE,
capacity=capacity
)
def get_idle_nodes(self) -> List[Node]:
"""获取空闲节点列表"""
with self.lock:
return [node for node in self.nodes.values()
if node.status == NodeStatus.IDLE and
node.current_load < node.capacity]
def schedule_task(self, task: Any):
"""动态调度任务到空闲节点"""
idle_nodes = self.get_idle_nodes()
if not idle_nodes:
print("没有可用空闲节点,任务排队等待")
self.task_queue.put(task)
return
# 选择负载最轻的节点
target_node = min(idle_nodes, key=lambda n: n.current_load)
with self.lock:
target_node.status = NodeStatus.BUSY
target_node.current_load += 1
# 在单独线程执行任务
thread = threading.Thread(
target=self.execute_task,
args=(target_node, task)
)
thread.start()
def execute_task(self, node: Node, task: Any):
"""执行任务并更新节点状态"""
try:
print(f"节点 {node.id} 开始执行任务: {task}")
time.sleep(random.uniform(1, 3)) # 模拟任务执行
# 任务完成后释放节点
with self.lock:
node.status = NodeStatus.IDLE
node.current_load -= 1
node.last_heartbeat = time.time()
print(f"节点 {node.id} 完成任务,恢复空闲")
except Exception as e:
print(f"节点 {node.id} 执行失败: {e}")
with self.lock:
node.status = NodeStatus.OFFLINE
基于权重的负载均衡调度
class WeightedNode:
def __init__(self, id: str, weight: int = 1):
self.id = id
self.weight = weight
self.current_load = 0
self.is_idle = True
self.score = 0
class WeightedScheduler:
def __init__(self):
self.nodes: List[WeightedNode] = []
self.current_index = 0
def add_node(self, node_id: str, weight: int = 1):
"""添加带权重的节点"""
self.nodes.append(WeightedNode(node_id, weight))
def select_node_weighted_round_robin(self) -> WeightedNode:
"""加权轮询选择节点"""
if not self.nodes:
return None
total_weight = sum(node.weight for node in self.nodes if node.is_idle)
if total_weight == 0:
return None
# 选择权重最高的空闲节点
idle_nodes = [n for n in self.nodes if n.is_idle]
if not idle_nodes:
return None
# 基于权重选择
return max(idle_nodes, key=lambda n: n.weight - n.current_load)
def adaptive_schedule(self, tasks: List[Any]):
"""自适应调度"""
for task in tasks:
node = self.select_node_weighted_round_robin()
if node:
node.is_idle = False
node.current_load += 1
print(f"调度任务到节点 {node.id} (权重: {node.weight})")
# 执行任务
self._execute_on_node(node, task)
def _execute_on_node(self, node: WeightedNode, task: Any):
"""在节点上执行任务"""
# 模拟执行
time.sleep(1)
node.is_idle = True
node.current_load -= 1
基于资源监控的动态调度
import psutil
import asyncio
from dataclasses import dataclass, field
from typing import Optional
@dataclass
class ResourceInfo:
cpu_usage: float = 0.0
memory_usage: float = 0.0
network_latency: float = 0.0
is_healthy: bool = True
class ResourceAwareScheduler:
def __init__(self, threshold: float = 0.8):
self.nodes: Dict[str, ResourceInfo] = {}
self.threshold = threshold # 资源使用阈值
def monitor_node_resources(self, node_id: str) -> ResourceInfo:
"""监控节点资源使用情况(模拟)"""
return ResourceInfo(
cpu_usage=random.uniform(0.1, 0.9),
memory_usage=random.uniform(0.2, 0.8),
network_latency=random.uniform(1, 100)
)
def can_accept_task(self, node_id: str) -> bool:
"""检查节点是否能接受新任务"""
if node_id not in self.nodes:
return True
info = self.nodes[node_id]
return (info.cpu_usage < self.threshold and
info.memory_usage < self.threshold and
info.is_healthy)
async def dynamic_schedule_async(self, task: Any):
"""异步动态调度"""
available_nodes = []
for node_id in self.nodes:
if self.can_accept_task(node_id):
available_nodes.append(node_id)
if not available_nodes:
raise Exception("没有可用节点")
# 选择资源使用率最低的节点
target_node = min(available_nodes,
key=lambda n: self.nodes[n].cpu_usage)
return await self.execute_async(target_node, task)
async def execute_async(self, node_id: str, task: Any):
"""异步执行任务"""
print(f"在节点 {node_id} 上异步执行任务")
await asyncio.sleep(random.uniform(0.5, 2))
return f"Task {task} completed on {node_id}"
完整示例:动态节点池管理
import asyncio
from concurrent.futures import ThreadPoolExecutor
import heapq
from typing import Callable
class DynamicNodePool:
"""动态节点池"""
def __init__(self, min_nodes: int = 2, max_nodes: int = 10):
self.min_nodes = min_nodes
self.max_nodes = max_nodes
self.nodes = []
self.executor = ThreadPoolExecutor(max_workers=max_nodes)
self.pending_tasks = []
self.running_tasks = {}
async def initialize(self):
"""初始化节点池"""
for i in range(self.min_nodes):
node = await self.create_node(f"node-{i}")
self.nodes.append(node)
async def create_node(self, node_id: str) -> dict:
"""创建新节点"""
node = {
'id': node_id,
'status': 'idle',
'tasks_completed': 0
}
print(f"创建节点: {node_id}")
return node
def get_idle_node(self):
"""获取空闲节点"""
for node in self.nodes:
if node['status'] == 'idle':
return node
return None
def scale_up(self):
"""扩容"""
if len(self.nodes) < self.max_nodes:
new_node_id = f"node-{len(self.nodes)}"
new_node = self.create_node(new_node_id)
self.nodes.append(new_node)
return new_node
return None
def scale_down(self):
"""缩容"""
idle_nodes = [n for n in self.nodes if n['status'] == 'idle']
if len(idle_nodes) > self.min_nodes:
node_to_remove = idle_nodes[0]
self.nodes.remove(node_to_remove)
print(f"移除节点: {node_to_remove['id']}")
async def execute_task(self, task: Any, task_func: Callable):
"""执行任务"""
node = self.get_idle_node()
if node is None and len(self.nodes) < self.max_nodes:
node = await self.scale_up()
if node is None:
self.pending_tasks.append(task)
print("任务排队等待")
return
# 标记节点忙碌
node['status'] = 'busy'
try:
# 在线程池中执行
loop = asyncio.get_event_loop()
result = await loop.run_in_executor(
self.executor,
task_func,
task
)
node['tasks_completed'] += 1
print(f"节点 {node['id']} 完成任务: {task}")
return result
finally:
node['status'] = 'idle'
# 处理排队的任务
if self.pending_tasks:
next_task = self.pending_tasks.pop(0)
await self.execute_task(next_task, task_func)
# 使用示例
async def main():
pool = DynamicNodePool(min_nodes=3, max_nodes=8)
await pool.initialize()
def process_task(task):
"""任务处理函数"""
time.sleep(1) # 模拟处理时间
return f"Processed: {task}"
# 提交多个任务
tasks = [f"task-{i}" for i in range(10)]
for task in tasks:
result = await pool.execute_task(task, process_task)
print(result)
# 动态调整节点池
pool.scale_down()
if __name__ == "__main__":
asyncio.run(main())
高级特性:基于优先级的调度
from dataclasses import dataclass, field
from typing import List, Tuple
import bisect
@dataclass(order=True)
class PriorityTask:
priority: int
task_id: str = field(compare=False)
data: Any = field(compare=False)
class PriorityScheduler:
def __init__(self):
self.task_queue: List[PriorityTask] = []
self.node_queue: List[str] = []
def add_task(self, task: PriorityTask):
"""添加优先级任务"""
heapq.heappush(self.task_queue, (-task.priority, task))
def register_node(self, node_id: str):
"""注册节点"""
self.node_queue.append(node_id)
def schedule_next(self):
"""调度下一个最高优先级任务"""
if not self.task_queue or not self.node_queue:
return None, None
_, task = heapq.heappop(self.task_queue)
# 轮询选择节点
node = self.node_queue.pop(0)
self.node_queue.append(node)
return node, task
这些调度方法可以根据实际需求组合使用,实现动态、高效的节点调度系统。