polishing
This commit is contained in:
parent
e9bf02a52c
commit
f604314207
|
@ -949,21 +949,20 @@ class DADF5():
|
|||
|
||||
|
||||
def _job(self,group,func,datasets,args,lock):
|
||||
def _read(group,datasets,lock):
|
||||
"""
|
||||
Execute job for _add_generic_pointwise
|
||||
"""
|
||||
try:
|
||||
datasets_in = {}
|
||||
lock.acquire()
|
||||
with h5py.File(self.fname,'r') as f:
|
||||
for k,v in datasets.items():
|
||||
loc = f[group+'/'+v]
|
||||
datasets_in[k]={'data':loc[()],
|
||||
'label':v,
|
||||
'meta':{k2:v2.decode() for k2,v2 in loc.attrs.items()}}
|
||||
for arg,label in datasets.items():
|
||||
loc = f[group+'/'+label]
|
||||
datasets_in[arg]={'data' :loc[()],
|
||||
'label':label,
|
||||
'meta': {k:v.decode() for k,v in loc.attrs.items()}}
|
||||
lock.release()
|
||||
return datasets_in
|
||||
|
||||
try:
|
||||
d = _read(group,datasets,lock)
|
||||
r = func(**d,**args)
|
||||
r = func(**datasets_in,**args)
|
||||
return [group,r]
|
||||
except Exception as err:
|
||||
print('Error during calculation: {}.'.format(err))
|
||||
|
@ -971,18 +970,28 @@ class DADF5():
|
|||
|
||||
|
||||
def _add_generic_pointwise(self,func,datasets,args={}):
|
||||
"""
|
||||
General function to add pointwise data.
|
||||
|
||||
env = Environment()
|
||||
N_threads = int(env.options['DAMASK_NUM_THREADS'])
|
||||
Parameters
|
||||
----------
|
||||
func : function
|
||||
Callback function that calculates a new dataset from one or more datasets per HDF5 group.
|
||||
datasets : dictionary
|
||||
Details of the datasets to be used: label (in HDF5 file) and arg (argument to which the data is parsed in func).
|
||||
args : dictionary, optional
|
||||
Arguments parsed to func.
|
||||
"""
|
||||
N_threads = int(Environment().options['DAMASK_NUM_THREADS'])
|
||||
pool = multiprocessing.Pool(N_threads)
|
||||
m = multiprocessing.Manager()
|
||||
lock = m.Lock()
|
||||
lock = multiprocessing.Manager().Lock()
|
||||
|
||||
groups = self.groups_with_datasets(datasets.values())
|
||||
default_arg = partial(self._job,func=func,datasets=datasets,args=args,lock=lock)
|
||||
util.progressBar(iteration=0,total=len(groups)-1)
|
||||
|
||||
util.progressBar(iteration=0,total=len(groups))
|
||||
for i,result in enumerate(pool.imap_unordered(default_arg,groups)):
|
||||
util.progressBar(iteration=i,total=len(groups)-1)
|
||||
util.progressBar(iteration=i+1,total=len(groups))
|
||||
if not result: continue
|
||||
lock.acquire()
|
||||
with h5py.File(self.fname, 'a') as f:
|
||||
|
@ -993,6 +1002,7 @@ class DADF5():
|
|||
except OSError as err:
|
||||
print('Could not add dataset: {}.'.format(err))
|
||||
lock.release()
|
||||
|
||||
pool.close()
|
||||
pool.join()
|
||||
|
||||
|
|
Loading…
Reference in New Issue