Разрешить нескольким работникам Ray заполнять общую очередь, доступ к которой можно получить из основного потока. ⇐ Python
-
Anonymous
Разрешить нескольким работникам Ray заполнять общую очередь, доступ к которой можно получить из основного потока.
У меня есть класс DataGenerator
@ray.remote класс DataGenerator: Защиту генерировать_непрерывно (сам): пока правда: время.сон(5) данные = случайный.ранд() # Мне нужно, чтобы данные были помещены в очередь, общую для всех экземпляров DataGenerator Из основного скрипта я создаю экземпляры многих из них
queue = # Некоторая общая очередь # Создайте дескрипторы генератора и запустите непрерывный сбор для всех handles = [DataGenerator.remote() для i в диапазоне (10)] ray.wait([handle.generate_continuously.remote() для дескриптора в дескрипторах]) # Постоянно извлекать очередь и сохранять результаты локально все_данные = [] пока правда: данные = очередь.pop_all() all_data.extend() # делаем что-то ресурсоемкое с all_data Мне нужно, чтобы каждый дескриптор помещал результат в общую очередь, к которой я могу неоднократно обращаться из основного сценария.
Что я пробовал Это самое близкое к желаемому результату:
@ray.remote класс DataGenerator: защита генерировать (сам): время.сон(5) данные = случайный.ранд() возвращать данные Н_РУЧКИ = 10 generator_handles = [DataGenerator.remote() для i в диапазоне (N_HANDLES)] # Сопоставьте индекс дескриптора генератора со ссылкой на объект, создаваемый удаленно handle_idx_to_ref = {idx:generator_handles[idx].remote.generate() для idx в диапазоне(N_HANDLES)} все_данные = [] пока правда: для idx ссылка на handle_idx_to_ref.items(): Ready_id, not_ready_id = ray.wait(, timeout=0) если готов_ид: all_data.extend(ray.get([ready_id])) # Начать заново генерацию для этого работника handle_idx_to_ref[idx] = генератор_handles[idx].remote.generate() # Сделайте что-нибудь ресурсоемкое с all_data Это почти хорошо, но если операции с интенсивными вычислениями занимают слишком много времени, некоторый DataGenerator может завершиться и не запуститься снова до следующей итерации. Как я могу улучшить этот код?
У меня есть класс DataGenerator
@ray.remote класс DataGenerator: Защиту генерировать_непрерывно (сам): пока правда: время.сон(5) данные = случайный.ранд() # Мне нужно, чтобы данные были помещены в очередь, общую для всех экземпляров DataGenerator Из основного скрипта я создаю экземпляры многих из них
queue = # Некоторая общая очередь # Создайте дескрипторы генератора и запустите непрерывный сбор для всех handles = [DataGenerator.remote() для i в диапазоне (10)] ray.wait([handle.generate_continuously.remote() для дескриптора в дескрипторах]) # Постоянно извлекать очередь и сохранять результаты локально все_данные = [] пока правда: данные = очередь.pop_all() all_data.extend() # делаем что-то ресурсоемкое с all_data Мне нужно, чтобы каждый дескриптор помещал результат в общую очередь, к которой я могу неоднократно обращаться из основного сценария.
Что я пробовал Это самое близкое к желаемому результату:
@ray.remote класс DataGenerator: защита генерировать (сам): время.сон(5) данные = случайный.ранд() возвращать данные Н_РУЧКИ = 10 generator_handles = [DataGenerator.remote() для i в диапазоне (N_HANDLES)] # Сопоставьте индекс дескриптора генератора со ссылкой на объект, создаваемый удаленно handle_idx_to_ref = {idx:generator_handles[idx].remote.generate() для idx в диапазоне(N_HANDLES)} все_данные = [] пока правда: для idx ссылка на handle_idx_to_ref.items(): Ready_id, not_ready_id = ray.wait(, timeout=0) если готов_ид: all_data.extend(ray.get([ready_id])) # Начать заново генерацию для этого работника handle_idx_to_ref[idx] = генератор_handles[idx].remote.generate() # Сделайте что-нибудь ресурсоемкое с all_data Это почти хорошо, но если операции с интенсивными вычислениями занимают слишком много времени, некоторый DataGenerator может завершиться и не запуститься снова до следующей итерации. Как я могу улучшить этот код?