def argos_group_translate(logger: Logger, o_set: set, num_workers=6) -> set:
if not o_set:
return o_set
task_queue = multiprocessing.Queue()
r_queue = multiprocessing.Queue()
results = set()
logger.info(f'Called translation for {len(o_set)} objects...')
amt = 0
for o in o_set:
if not o.name.invalid_lang_name and any([n is None for n in [o.name.name_en, o.name.name_es, o.name.name_it,
o.name.name_fr, o.name.name_de]]):
task_queue.put(o)
amt += 1
else:
results.add(o)
if amt == 0:
return o_set
logger.info(f'Starting translation for {amt} {o.__class__.__name__} items using {num_workers} workers.')
workers = [multiprocessing.Process(target=argos_target_worker, args=(task_queue, r_queue))
for _ in range(num_workers)]
for p in workers:
p.start()
expected_results = amt
collected = 0
while collected < expected_results:
try:
if collected % 1000 == 0:
logger.debug(f'Current tally: {collected}/{expected_results}')
result = r_queue.get(timeout=10)
if result is not None:
results.add(result)
collected += 1
else:
logger.warn("Main: Received unexpected None from r_queue")
except Empty:
logger.warn(f"Main: Timeout waiting for result after 10s, collected {collected}/{expected_results}")
break
for _ in range(len(workers)):
task_queue.put(None)
for p in workers:
p.join()
return results