Ray集群中Head节点(调度器所在节点)未被分配任何任务/Actor的问题求助
Ray集群中Head节点(调度器所在节点)未被分配任何任务/Actor的问题求助
大家好,我目前碰到了一个Ray集群的棘手问题:作为调度器的Head节点(我用本地机器充当)始终没有被分配任何任务或Actor,想请各位帮忙分析下原因。
先跟大家说下我的集群启动方式:
Head节点(调度器)启动命令
nohup ray start --head \ --node-ip-address=<HOST_IP> \ --port=6379 \ --resources='{"is_worker": 1}' --dashboard-host=localhost &
Worker节点启动命令
nohup ray start \ --address='<SCHEDULER_IP>:6379' \ --resources='{"is_worker": 1}' --node-ip-address=<WORKER_IP> &
为了能在集群资源变化(比如新增实例)时自动分配Actor,我写了一个DynamicPoolManager类,代码如下:
class DynamicPoolManager: def __init__(self, actor_per_cpu, mongouri, embeddings): self.workers = [] self.pool = None self.mongouri = mongouri self.embeddings = embeddings self.actor_per_cpu = actor_per_cpu self.scale() # Initial scale def scale(self): # 1. Clean up dead actors first # self.workers = [w for w in self.workers if self._is_alive(w)] # 2. Get current capacity # Use 'AvailableResources' if you want to be polite to other apps, # or 'cluster_resources' to claim the whole cluster. total_cpus = int(ray.available_resources().get("CPU", 0)) # Get the ID of the node the scheduler is currently running on current_node_id = ray.get_runtime_context().get_node_id() # Find out how many CPUs are on THIS specific node nodes = ray.nodes() scheduler_node_cpus = 0 for node in nodes: if node["NodeID"] == current_node_id: scheduler_node_cpus = node["Resources"].get("CPU", 0) break # Your actual target is the cluster total MINUS the head node's CPUs # total_cpus -= int(scheduler_node_cpus) num_actors = int(total_cpus * self.actor_per_cpu) # 3. Scale Up if len(self.workers) < num_actors: new_count = num_actors - len(self.workers) log(f"Resources added. Creating {new_count} actors") self.workers.extend([worker.Actor.remote(self.mongouri, self.embeddings) for _ in range(new_count)]) # 4. Scale Down (Crucial for Autoscaler) elif len(self.workers) > num_actors: to_remove = len(self.workers) - num_actors log(f"Resources removed. Removing {to_remove} actors") for _ in range(to_remove): victim = self.workers.pop() victim.shutdown.remote() # 5. Update the pool reference self.pool = ActorPool(self.workers)
对应的Worker Actor代码是这样的:
@ray.remote(resources={"is_worker": 1}) class Actor: def __init__(self, mongo_uri, embeddings): # INIT LOGIC log("Actor ready") def run(): # TASK
我查资料看到Ray的Head节点默认是可以作为Worker节点参与任务分配的,但实际测试下来,不管我有没有在启动命令和Actor注解里指定is_worker这个自定义资源,Head节点都不会被分配任何任务或者Actor。我的集群架构是本地机器做Head,两台云VM实例作为Worker,用Tailscale搭建了虚拟网络确保所有节点互通。
麻烦各位帮忙看看哪里出问题了,谢谢!
备注:内容来源于stack exchange,提问作者cuneyttyler
相关产品推荐
相关产品推荐

