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

如何实现多Ray Worker填充主线程可访问公共队列并优化代码

解决方案

方案一:使用Ray官方共享队列(推荐)

Ray提供了ray.util.queue.Queue,专门用于分布式场景下的多Actor/任务共享数据。让每个DataGenerator实例持续生成数据并写入共享队列,主线程独立从队列读取数据并处理,生成和计算完全异步,不会互相阻塞。

完整代码实现

import ray
import time
import random
from ray.util.queue import Queue

ray.init()

@ray.remote
class DataGenerator:
    def __init__(self, shared_queue):
        self.shared_queue = shared_queue

    def generate_continuously(self):
        while True:
            time.sleep(5)
            data = random.rand()
            # 将生成的数据放入共享队列
            self.shared_queue.put(data)

# 创建共享队列
queue = Queue(maxsize=100)  # 根据需求设置队列最大容量

# 实例化10个生成器并启动持续生成任务
handles = [DataGenerator.remote(queue) for _ in range(10)]
# 启动所有生成任务(非阻塞,生成器会在后台持续运行)
[handle.generate_continuously.remote() for handle in handles]

all_data = []
while True:
    # 从队列批量获取数据(可设置超时避免无限阻塞)
    batch = []
    while not queue.empty():
        batch.append(queue.get())
    
    if batch:
        all_data.extend(batch)
        # 执行计算密集型操作
        print(f"处理了{len(batch)}条数据,累计数据量:{len(all_data)}")
        # 这里替换成你的计算逻辑
        time.sleep(10)  # 模拟耗时计算
    else:
        time.sleep(0.1)  # 队列空时短暂休眠,避免空轮询占用CPU

方案优势

  • 生成器持续独立运行,不受主线程计算耗时影响,不会出现生成任务等待重启的情况
  • 队列天然解耦生成和处理逻辑,代码更简洁易维护
  • 支持批量读写,适合处理高频生成的数据

方案二:改进原有轮询逻辑

如果不想使用队列,也可以优化原有代码,让主线程先批量处理所有已完成的生成任务,再执行计算操作,确保生成任务能立刻重启,不会被计算阻塞。

改进后代码

import ray
import time
import random

ray.init()

@ray.remote
class DataGenerator:
    def generate(self):
        time.sleep(5)
        data = random.rand()
        return data

N_HANDLES = 10
generator_handles = [DataGenerator.remote() for _ in range(N_HANDLES)]
# 初始化所有生成任务
pending_refs = [handle.generate.remote() for handle in generator_handles]

all_data = []
while True:
    # 监听所有未完成的任务,一旦有就绪的就立即处理
    ready_refs, pending_refs = ray.wait(pending_refs, num_returns=len(pending_refs), timeout=0)
    
    if ready_refs:
        # 获取所有就绪的数据
        new_data = ray.get(ready_refs)
        all_data.extend(new_data)
        # 为每个完成的任务重启生成任务,并加入待处理列表
        for idx in range(len(ready_refs)):
            pending_refs.append(generator_handles[idx].generate.remote())
    
    # 执行计算密集型操作
    if all_data:
        print(f"当前累计数据量:{len(all_data)},开始计算...")
        # 模拟耗时计算
        time.sleep(10)
    else:
        time.sleep(0.1)

改进点说明

  • 使用ray.wait一次性监听所有待处理任务,批量获取就绪结果,避免逐个轮询的低效
  • 处理完就绪任务后立刻重启对应的生成任务,确保生成器不会空闲等待
  • 计算操作放在轮询处理完成之后,生成任务的重启不受计算耗时影响

内容的提问来源于stack exchange,提问作者Federico Taschin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 18:55:12