如何实现多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
相关产品推荐
相关产品推荐

