executor = cf.ProcessPoolExecutor(max_workers=prefetch) // pylint: disable=redefined-variable-type
else:
raise ValueError("target should be one of ["threads", "mpc"]")
futures = []
for batch in batch_generator:
futures.append(executor.submit(self._exec_all_actions, batch, True))
// wait until all batches have been processed
_ = [future.result() for future in futures]
else:
self._run_seq(batch_generator)
return self
After Change
else:
raise ValueError("target should be one of ["threads", "mpc"]")
self.prefetch_queue = q.Queue(maxsize=prefetch)
service_executor = cf.ThreadPoolExecutor(max_workers=2)
service_executor.submit(self._put_batches_into_queue, batch_generator)
future = service_executor.submit(self._run_batches_from_queue)
// wait until all batches have been processed