You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.13 18:07:55