66
77try :
88 from dask .distributed import Client
9- from dask_jobqueue import JobQueueCluster
9+ from dask_jobqueue . core import JobQueueCluster
1010
1111 dask_available = True
1212except ImportError :
@@ -17,8 +17,8 @@ def serial_execution(func: Callable, entries: List):
1717 return [func (args ) for args in tqdm .tqdm (entries , total = len (entries ))]
1818
1919
20- def parallel_execution_processes (func : Callable , entries : List , njobs : int ):
21- pool = mp .Pool (njobs )
20+ def parallel_execution_processes (func : Callable , entries : List , process_count : int ):
21+ pool = mp .Pool (process_count )
2222 results = [entry for entry in tqdm .tqdm (pool .imap (func , entries ), total = len (entries ))]
2323 pool .close ()
2424
@@ -75,7 +75,7 @@ def run(self, func: Callable, entries: List):
7575 if self .mode == RunnerMode .SERIAL :
7676 return serial_execution (func , entries )
7777 elif self .mode == RunnerMode .PROCESSES :
78- return parallel_execution_processes (func , entries , njobs = self .process_count )
78+ return parallel_execution_processes (func , entries , process_count = self .process_count )
7979 elif self .mode == RunnerMode .DASK_LOCAL :
8080 return parallel_execution_dask_local (func , entries , process_count = self .process_count , process_threads_count = self .thread_count )
8181 if self .mode == RunnerMode .DASK_JOB_QUEUE_CLUSTER :
0 commit comments