Python脚本如何动态调度空闲同步节点

wen python案例 31

本文目录导读:

Python脚本如何动态调度空闲同步节点

  1. 基于队列的任务调度
  2. 基于权重的负载均衡调度
  3. 基于资源监控的动态调度
  4. 完整示例:动态节点池管理
  5. 高级特性:基于优先级的调度

我来介绍几种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

这些调度方法可以根据实际需求组合使用,实现动态、高效的节点调度系统。

抱歉,评论功能暂时关闭!