Compare commits

...
Author SHA1 Message Date
Brendan Smithyman 96ce267b08 Merge branch 'master' into parallel 2015-08-24 11:12:09 -04:00
Brendan Smithyman ad99b24a82 Change to new import standard for IPython Parallel. 2015-08-24 11:09:17 -04:00
Brendan Smithyman 3f02ce9073 Store more useful information in jobs. 2015-07-02 14:25:45 -04:00
Brendan Smithyman 12e64f418c Change location of endpointName setting. 2015-07-02 14:25:15 -04:00
Brendan Smithyman c507955903 Source estimation code. 2015-06-29 18:45:44 -04:00
Brendan Smithyman 0cf783d993 Improvements to scheduler (now use slices).
Change name "dispatcher" to "problem".
2015-06-29 14:30:12 -04:00
Brendan Smithyman 0987897896 Rename. 2015-06-29 09:49:45 -04:00
Brendan Smithyman 0db348cdce Begin to restructure such that the objects in the Endpoint are Survey and Problem subclasses. 2015-06-23 10:48:22 -04:00
Brendan Smithyman 0cd0308d44 Fixed error when the 'mkl' package is not installed.' 2015-06-22 14:42:11 -04:00
Brendan Smithyman 522c6c943e Make changes to allow for better integration with Endpoint code. 2015-06-22 11:02:58 -04:00
Brendan Smithyman 4021638d2b Merge commit '6d3d8d78b601a11c52e90f87367f7b40e4d537a5' into parallel 2015-06-22 09:26:48 -04:00
Brendan Smithyman aef3958794 Updated "Endpoint" architecture. 2015-05-11 10:28:16 -04:00
Brendan Smithyman 064caa1427 Refactor to use Endpoint namespacing.
Now everything in the remote memory footprint sits inside an object called an 'Endpoint'. The interfaces to set up and hold fields, etc. are generic, and are initialized as callables in the remote namespace, so they should be able to be adapted to different data backends (Fields objects, HDF5, etc.).
2015-05-07 17:14:02 -04:00
Brendan Smithyman d5cb4adf1f Add support for bootstrapping remote namespace based on code from Problem or Dispatcher. 2015-05-07 14:01:46 -04:00
Brendan Smithyman dc7b102550 Rough proposal design for Dispatcher/Problem interface. Refactor at will. 2015-05-07 13:33:03 -04:00
Brendan Smithyman 75e3163a47 Added some prototypes for parallel constructs.
SuperReference acts like an IPython.parallel.Reference instance specialized for function calls. If it's scheduled by a call to 'apply' on a Load Balanced View, it will raise UnmetDependency errors until it's scheduled on a worker that is allowed by its 'rank' attribute. This can therefore point at functions that are known to be defined on particular workers.

Endpoint is a prototype layout for an object to hold multiple Problems, Fields, etc. on the remote workers. The idea is to clean up the remote namespace and reduce the use of globals().
2015-05-07 13:32:37 -04:00
Brendan Smithyman d94efc3b85 Merge branch 'parallel' of github.com:simpeg/simpeg into parallel 2015-05-06 12:23:30 -04:00
Brendan Smithyman 4765b541cd Change graph plotting style and enable doubleclick.
Double-clicking on a graph node now inserts a new IPython Notebook cell and sets its text to the job listing for that node.
2015-05-06 12:23:06 -04:00
Brendan Smithyman c982fb131a Move module-level functions to static methods for their respective classes. 2015-05-06 12:21:46 -04:00
Brendan Smithyman 989d7c12fd Add NetworkX dependency 2015-05-06 10:27:42 -04:00
Rowan Cockett 0ddecdc7ce Merge branch 'parallel' of https://github.com/simpeg/simpeg into parallel 2015-05-05 17:39:28 -07:00
Rowan Cockett ea2b75bd97 try adding pyzmq to make travis work! 2015-05-05 15:57:02 -07:00
Brendan Smithyman d7d665456f Merge branch 'parallel' of github.com:simpeg/simpeg into parallel 2015-05-05 17:40:29 -04:00
Brendan Smithyman 983df96750 Better sizing. 2015-05-05 17:39:54 -04:00
Rowan Cockett c494d8ad63 add IPython to travis dependencies. 2015-05-05 14:08:07 -07:00
Brendan Smithyman 97f832031a Created SystemGraph subclass of the networkx.DiGraph class that knows how to render itself in the IPython Notebook using d3.js. 2015-05-05 17:03:45 -04:00
Brendan Smithyman caf40af35c Merge branch 'master' into parallel 2015-05-04 14:01:46 -04:00
Brendan Smithyman 29195e6608 Incorporating the first parts of my parallel distributed tools.
DataWrappers.py contains the CommonReducer class, which enables pass-through math operations and function calls on dictionaries.
Parallel.py contains the job scheduler, remote interface (w/ MPI support) and several helper functions.
2015-05-01 17:04:59 -04:00
7 changed files with 919 additions and 1 deletions
+1 -1
View File
@@ -15,7 +15,7 @@ before_install:
# Install packages
install:
- conda install --yes pip python=$TRAVIS_PYTHON_VERSION numpy scipy matplotlib cython
- conda install --yes pip python=$TRAVIS_PYTHON_VERSION numpy scipy matplotlib cython ipython networkx pyzmq
- pip install nose-cov python-coveralls
# - pip install -r requirements.txt
- python setup.py install
+112
View File
@@ -0,0 +1,112 @@
### PROTOTYPE INTERFACE FOR PARALLEL DISPATCHER ###
from functools import wraps
def synchronize(fn):
@wraps(fn)
def wrapper(*args, **kwargs):
self = args[0]
pr = isinstance(getattr(self, '_dispatcher', None), ParallelDispatcher)
if pr:
print('Parallel stuff: (start) %(prob)s.%(fn)s'%{'prob': self.__class__.__name__, 'fn': fn.__name__})
result = fn(*args, **kwargs)
if pr:
print('Parallel stuff: ( end ) %(prob)s.%(fn)s'%{'prob': self.__class__.__name__, 'fn': fn.__name__})
return result
return wrapper
class BaseDispatcher(object):
def __init__(self, *args, **kwargs):
print('INIT: Dispatcher!')
def pair(self, problem):
self._prob = problem
print('PAIR: Dispatcher setup...')
class SerialDispatcher(BaseDispatcher):
def __init__(self, *args, **kwargs):
BaseDispatcher.__init__(self, *args, **kwargs)
print('INIT: Serial dispatcher...')
class ParallelDispatcher(BaseDispatcher):
remoteOnly = ['someotherattribute']
def __init__(self, *args, **kwargs):
BaseDispatcher.__init__(self, *args, **kwargs)
print('INIT: Parallel dispatcher...')
def pair(self, problem):
BaseDispatcher.pair(self, problem)
print('PAIR: Parallel dispatcher setup...')
def interceptSetattr(self, prob, name, value):
print('SET: Parallel dispatcher set %(prob)s.%(name)s = %(value)r'
%{'prob': prob.__class__.__name__, 'name': name, 'value': value})
if name in self.remoteOnly:
print('Setting remote state...')
else:
raise AttributeError('Set local copy!')
def interceptGetattr(self, prob, name):
print('GET: Parallel dispatcher get %(prob)s.%(name)s'
%{'prob': prob.__class__.__name__, 'name': name})
if name in self.remoteOnly:
return '***Value from remote state***'
else:
raise AttributeError('Attribute %s not in parallel namespace!'%(name,))
class StandinSurvey(object):
def pair(self, problem):
self._prob = problem
class StandinProblem(object):
def __init__(self):
print('INIT: Problem!')
self._dispatcher = SerialDispatcher()
def __setattr__(self, name, value):
d = getattr(self, '_dispatcher', None)
if isinstance(d, ParallelDispatcher):
try:
d.interceptSetattr(self, name, value)
except AttributeError:
super(self.__class__, self).__setattr__(name, value)
finally:
return
else:
super(self.__class__, self).__setattr__(name, value)
def __getattr__(self, name):
d = super(self.__class__, self).__getattribute__('_dispatcher')
if isinstance(d, ParallelDispatcher):
return d.interceptGetattr(self, name)
def pair(self, survey, dispatcher=None):
self._survey = survey
self._survey.pair(self)
if dispatcher is not None:
self._dispatcher = dispatcher
print('PAIR: Problem setup...')
self._dispatcher.pair(self)
@synchronize
def dosomething(self):
print('Doing something!')
+638
View File
@@ -0,0 +1,638 @@
from ipyparallel import Client, parallel, Reference, require, depend, interactive
from SimPEG.Utils import CommonReducer
import numpy as np
import networkx
DEFAULT_MPI = True
MPI_BELLWETHERS = ['PMI_SIZE', 'OMPI_UNIVERSE_SIZE']
class SuperReference(object):
'''
Object that can be called to return a reference, but
will only be schedulable on the correct worker(s) if
its 'lrank' parameter has been set.
'''
def __init__(self, ref, lrank=None):
if (lrank is None) or (type(lrank) is list):
self.rank = lrank
else:
self.rank = [lrank]
self.ref = ref
def __call__(self, *args, **kwargs):
from ipyparallel import depend
from ipyparallel.error import UnmetDependency
if (self.rank is not None) and (globals().get('rank', None) not in self.rank):
raise UnmetDependency('Global \'rank\' does not satisfy requirements')
return self.ref(*args, **kwargs)
class Endpoint(object):
'''
Object that holds the namespace of the SimPEG parallel
footprint on the remote workers.
'''
problemFactory = lambda: None # Callable for constructing system / problem
surveyFactory = lambda: None # Callable for constructing survey
localFields = {} # Dictionary for storing local fields
globalFields = {} # Dictionary for storing merged fields
localProblems = {} # Dictionary of local subsystem / problem objects
localSurveys = {} # Dictionary of local survey objects
functions = {} # Dictionary of callables to carry out modelling / etc.
fieldspec = None # Dictionary of callables to setup field storage objects
baseSystemConfig = {} # Base configuration for system
def setupLocalFields(self, whichfields=None):
# If no names are specified, clear all fields first
if whichfields is None:
self.localFields = {}
# If we have a 'fieldspec' object...
if getattr(self, 'fieldspec', None) is not None:
# ...either loop over the specified names, or all the fields...
for fn in (whichfields or self.fieldspec):
# ...and construct a new empty object per the 'fieldspec' constructor.
self.localFields[fn] = self.fieldspec[fn]()
def setupLocalSurveys(self, subConfigs):
# Loop over possible survey configurations (may differ in source terms, etc.)
for isub in subConfigs:
# For each 'isub' create a separate copy of the base configuration...
geom = self.baseSystemConfig['geom'].copy()
# ...and update with any differences...
geom.update(subConfigs[isub])
# ...then construct the Survey object and store it for later pairing.
self.localSurveys[isub] = self.surveyFactory(geom)
def setupLocalProblem(self, subConfig):
# Make a copy w/o the geometry information (which is used by Survey)
systemConfig = {key: self.baseSystemConfig[key] for key in self.baseSystemConfig if key not in ['geom']}
# Update with the configuration for this subproblem
systemConfig.update(subConfig)
# Create the local subproblem...
problem = self.problemFactory(systemConfig)
# ...and pair it with a corresponding survey for this 'isub' (e.g., frequency)...
problem.pair(self.localSurveys[subConfig['isub']])
# ...then store in the Endpoint for later access by the scheduler.
self.localProblems[subConfig['tag']] = problem
class SystemGraph(networkx.DiGraph):
'''
NetworkX Directed Graph subclass that knows about
job status information, and can return a representation
of itself for use in interactive debugging/testing.
'''
@staticmethod
def _codeStatus(data):
status = 0
if 'jobs' in data:
status = 1 * data['jobs'][-1].ready() + 1
if status > 1:
status += 1 * (not data['jobs'][-1].successful())
return status
def _codeGraph(self):
from networkx.readwrite import json_graph
G = networkx.DiGraph()
for e in self.edges_iter():
G.add_edge(e[0], e[1])
for n, data in self.nodes_iter(data=True):
G.add_node(n, status=self._codeStatus(data))
return json_graph.node_link_data(G)
def RenderHTML(self):
import pkg_resources
from IPython.core import display
import time
data = str(self._codeGraph())
uniqueID = hash(time.time())
formatstr = {
'uniqueID': 'Graph%s'%uniqueID,
'JSONData': data,
}
code = pkg_resources.resource_string('SimPEG', 'Resources/Parallel/SystemGraph.html')%formatstr
return display.HTML(data=code)._repr_html_()
try:
get_ipython().display_formatter.formatters['text/html'].for_type(SystemGraph, SystemGraph.RenderHTML)
except NameError:
pass
class SystemSolver(object):
def __init__(self, problem, schedule):
self.problem = problem
self.remote = problem.remote
self.schedule = schedule
def __call__(self, entry, isrcs):
# TODO: Replace with SuperReference instances
fnformat = '%s.functions["%s"]'
fnRef = Reference(fnformat%(self.remote.endpointName, self.schedule[entry]['solve']))
clearRef = Reference(fnformat%(self.remote.endpointName, self.schedule[entry]['clear']))
reduceLabels = self.schedule[entry]['reduce']
dview = self.problem.remote.dview
lview = self.problem.remote.lview
chunksPerWorker = getattr(self.problem, 'chunksPerWorker', 1)
G = SystemGraph()
mainNode = 'Beginning'
G.add_node(mainNode)
# Parse sources
# TODO: Get from Survey somehow?
nsrc = self.problem.nsrc
if isrcs is None:
isrcs = slice(None)
elif not isinstance(isrcs, slice):
raise Exception('Scheduler must run over slice or None!')
# TODO: Replace w/ hook into Endpoint classes
systemsOnWorkers = dview['%s.localProblems.keys()'%self.remote.endpointName]
ids = dview['rank']
tags = set()
for ltags in systemsOnWorkers:
tags = tags.union(set(ltags))
clearJobs = []
endNodes = {}
tailNodes = []
for tag in tags:
tagNode = 'Head: %d, %d'%tag
G.add_edge(mainNode, tagNode)
relIDs = []
for i in xrange(len(ids)):
systems = systemsOnWorkers[i]
rank = ids[i]
if tag in systems:
relIDs.append(rank)
systemJobs = []
endNodes[tag] = []
systemNodes = []
with lview.temp_flags(block=False):
iworks = 0
for work in self._subSlice(isrcs, int(round(chunksPerWorker*len(relIDs)))):
if work:
job = lview.apply(fnRef, Reference(self.remote.endpointName), tag, work)
systemJobs.append(job)
label = 'Compute: %d, %d, %d'%(tag[0], tag[1], iworks)
systemNodes.append(label)
G.add_node(label, jobs=[job], subslice=work, tag=tag)
G.add_edge(tagNode, label)
iworks += 1
if getattr(self.problem, 'ensembleClear', False): # True for ensemble ending, False for individual ending
tagNode = 'Wrap: %d, %d'%tag
for label in systemNodes:
G.add_edge(label, tagNode)
for rank in relIDs:
with lview.temp_flags(block=False, after=systemJobs):
# TODO: Remove dependency on self._hasSystemRank, once the SuperReferences
# are able to be used. They will automatically schedule only on the
# correct (allowed) systems.
job = lview.apply(clearRef, Reference(self.remote.endpointName), tag, rank)
clearJobs.append(job)
label = 'Wrap: %d, %d, %d'%(tag[0],tag[1], rank)
G.add_node(label, jobs=[job], tag=tag, rank=rank)
endNodes[tag].append(label)
G.add_edge(tagNode, label)
else:
for i, sjob in enumerate(systemJobs):
with lview.temp_flags(block=False, follow=sjob, after=sjob):
job = lview.apply(clearRef, Reference(self.remote.endpointName), tag)
clearJobs.append(job)
label = 'Wrap: %d, %d, %d'%(tag[0],tag[1],i)
G.add_node(label, jobs=[job])
endNodes[tag].append(label)
G.add_edge(systemNodes[i], label)
tagNode = 'Tail: %d, %d'%tag
for label in endNodes[tag]:
G.add_edge(label, tagNode)
tailNodes.append(tagNode)
endNode = 'End'
jobs = []
after = clearJobs
for label in reduceLabels:
job = self.problem.remote.reduceLB(Reference(self.remote.endpointName), label, after)
after = job
if job is not None:
jobs.append(job)
G.add_node(endNode, jobs=jobs)
for node in tailNodes:
G.add_edge(node, endNode)
return G
def wait(self, G):
self.problem.remote.lview.wait(G.node['End']['jobs'] if G.node['End']['jobs'] else (G.node[wn]['jobs'] for wn in (G.predecessors(tn)[0] for tn in G.predecessors('End'))))
# TODO: Hopefully obsoleted by SuperReference
@staticmethod
@interactive
def _hasSystemRank(endpoint, tag, wid):
global rank
return (tag in endpoint.localProblems) and (rank == wid)
@staticmethod
def _getChunks(problems, chunks=1):
nproblems = len(problems)
return (problems[i*nproblems // chunks: (i+1)*nproblems // chunks] for i in range(chunks))
@staticmethod
def _subSlice(insl, chunks=1):
start = insl.start or 0
nproblems = insl.stop - start
return [slice(start + i*nproblems/chunks, start + (i+1)*nproblems/chunks) for i in xrange(chunks)]
class RemoteInterface(object):
def __init__(self, profile=None, MPI=None, nThreads=1, bootstrap=None, endpointName='endpoint'):
# TODO: Add interface for namespace bootstrapping from
# the dispatcher / problem side
if profile is not None:
pupdate = {'profile': profile}
else:
pupdate = {}
pclient = Client(**pupdate)
if not self._cdSame(pclient):
print('Could not change all workers to the same directory as the client!')
dview = pclient[:]
dview.block = True
dview.clear()
remoteSetup = '''
import os'''
parMPISetup = '''
from mpi4py import MPI
comm = MPI.COMM_WORLD
rank = comm.Get_rank()'''
for command in remoteSetup.strip().split('\n'):
dview.execute(command.strip())
dview.scatter('rank', pclient.ids, flatten=True)
self.e0 = pclient[0]
self.e0.block = True
self.useMPI = False
MPI = DEFAULT_MPI if MPI is None else MPI
if MPI:
MPISafe = False
for var in MPI_BELLWETHERS:
MPISafe = MPISafe or all(dview['os.getenv("%s")'%(var,)])
if MPISafe:
for command in parMPISetup.strip().split('\n'):
dview.execute(command.strip())
ranks = dview['rank']
reorder = [ranks.index(i) for i in xrange(len(ranks))]
dview = pclient[reorder]
dview.block = True
dview.activate()
# Set up necessary parts for broadcast-based communication
self.e0 = pclient[reorder[0]]
self.e0.block = True
self.comm = Reference('comm')
self.useMPI = MPISafe
self.pclient = pclient
self.dview = dview
self.lview = pclient.load_balanced_view()
self.nThreads = nThreads
if bootstrap is not None:
for command in bootstrap.strip().split('\n'):
dview.execute(command.strip())
self.endpointName = endpointName
@property
def nThreads(self):
return self._nThreads
@nThreads.setter
def nThreads(self, value):
self._nThreads = value
self.dview.apply(self._adjustMKLVectorization, self._nThreads)
def __setitem__(self, key, item):
if self.useMPI:
self.e0[key] = item
code = 'if rank != 0: %(key)s = None\n%(key)s = comm.bcast(%(key)s, root=0)'
self.dview.execute(code%{'key': key})
else:
self.dview[key] = item
def __getitem__(self, key):
if self.useMPI:
code = 'temp_%(key)s = None\ntemp_%(key)s = comm.gather(%(key)s, root=%(root)d)'
self.dview.execute(code%{'key': key, 'root': 0})
item = self.e0['temp_%s'%(key,)]
self.e0.execute('del temp_%s'%(key,))
else:
item = self.dview[key]
return item
def reduceLB(self, endpoint, key, after=None):
repeat = lambda value: (value for i in xrange(len(self.pclient.ids)))
if self.useMPI:
with self.lview.temp_flags(block=False, after=after):
job = self.lview.map(self._reduceJob, xrange(len(self.pclient.ids)), repeat(0), repeat(endpoint), repeat(key))
return job
def reduce(self, key, axis=None):
if self.useMPI:
code = 'temp_%(key)s = comm.reduce(%(key)s, root=%(root)d)'
self.dview.execute(code%{'key': key, 'root': 0})
# if axis is not None:
# code = 'temp_%(key)s = temp_%(key)s.sum(axis=%(axis)d)'
# self.e0.execute(code%{'key': key, 'axis': axis})
item = self.e0['temp_%s'%(key,)]
self.dview.execute('del temp_%s'%(key,))
else:
item = reduce(np.add, self.dview[key])
return item
def reduceMul(self, key1, key2, axis=None):
if self.useMPI:
# Gather
code_reduce = 'temp_%(key)s = comm.reduce(%(key)s, root=%(root)d)'
self.dview.execute(code_reduce%{'key': key1, 'root': 0})
self.dview.execute(code_reduce%{'key': key2, 'root': 0})
# Multiply
code_mul = 'temp_%(key1)s%(key2)s = temp_%(key1)s * temp_%(key2)s'
self.e0.execute(code_mul%{'key1': key1, 'key2': key2})
# Potentially sum
if axis is not None:
code = 'temp_%(key1)s%(key2)s = temp_%(key1)s%(key2)s.sum(axis=%(axis)d)'
self.e0.execute(code%{'key1': key1, 'key2': key2, 'axis': axis})
# Pull
item = self.e0['temp_%(key1)s%(key2)s'%{'key1': key1, 'key2': key2}]
# Clear
self.dview.execute('del temp_%s'%(key1,))
self.dview.execute('del temp_%s'%(key2,))
self.e0.execute('del temp_%(key1)s%(key2)s'%{'key1': key1, 'key2': key2})
else:
item1 = reduce(np.add, self.dview[key1])
item2 = reduce(np.add, self.dview[key2])
item = item1 * item2
return item
def remoteMulE0(self, key1, key2, axis=None):
code_mul = 'temp_field = %(key1)s * %(key2)s'
self.e0.execute(code_mul%{'key1': key1, 'key2': key2})
if axis is not None:
code = 'temp_field = temp_field.sum(axis=%(axis)d)'
self.e0.execute(code%{'axis': axis})
item = self.e0['temp_field']
self.e0.execute('del temp_field')
return item
def remoteDifference(self, key1, key2, keyresult):
if self.useMPI:
root = 0
# Gather
code_reduce = 'temp_%(key)s = comm.reduce(%(key)s, root=%(root)d)'
self.dview.execute(code_reduce%{'key': key1, 'root': root})
self.dview.execute(code_reduce%{'key': key2, 'root': root})
# Difference
code_difference = '%(keyresult)s = temp_%(key1)s - temp_%(key2)s'
self.e0.execute(code_difference%{'key1': key1, 'key2': key2, 'keyresult': keyresult})
# Broadcast
code = 'if rank != 0: %(key)s = None\n%(key)s = comm.bcast(%(key)s, root=%(root)d)'
self.dview.execute(code%{'key': keyresult, 'root': root})
# Clear
self.e0.execute('del temp_%s'%(key1,))
self.e0.execute('del temp_%s'%(key2,))
else:
item1 = reduce(np.add, self.dview[key1])
item2 = reduce(np.add, self.dview[key2])
item = item1 - item2
self.dview[keyresult] = item
def remoteOpGatherFirst(self, op, key1, key2, keyresult):
if self.useMPI:
root = 0
# Gather
code_reduce = 'temp_%(key)s = comm.reduce(%(key)s, root=%(root)d)'
self.dview.execute(code_reduce%{'key': key1, 'root': root})
# Difference
code_difference = '%(keyresult)s = temp_%(key1)s %(op)s %(key2)s'
self.e0.execute(code_difference%{'op': op, 'key1': key1, 'key2': key2, 'keyresult': keyresult})
# Broadcast
code = 'if rank != 0: %(key)s = None\n%(key)s = comm.bcast(%(key)s, root=%(root)d)'
self.dview.execute(code%{'key': keyresult, 'root': root})
# Clear
self.e0.execute('del temp_%s'%(key1,))
else:
item1 = reduce(np.add, self.dview[key1])
item2 = self.e0[key2] # Assumes that any arbitrary worker has this information
item = eval('item1 %s item2'%(op,))
self.dview[keyresult] = item
def remoteDifferenceGatherFirst(self, *args):
self.remoteOpGatherFirst('-', *args)
def remoteSrcEstGatherFirst(self, keyresult, key1, key2, individual=False):
if self.useMPI:
root = 0
# # Gather
# code_reduce = 'temp_%(key)s = comm.reduce(%(key)s, root=%(root)d)'
# self.dview.execute(code_reduce%{'key': key1, 'root': root})
# SrcEst
if individual:
code_srcest = '%(keyresult)s = (%(key2)s.conj() * %(key1)s).sum(axis=1) / (%(key1)s.conj() * %(key1)s).sum(axis=1)'
else:
code_srcest = '%(keyresult)s = (%(key2)s.conj() * %(key1)s).sum() / (%(key1)s.conj() * %(key1)s).sum()'
self.e0.execute(code_srcest%{'key1': key1, 'key2': key2, 'keyresult': keyresult})
# Broadcast
code = 'if rank != %(root)d: %(key)s = None\n%(key)s = comm.bcast(%(key)s, root=%(root)d)'
self.dview.execute(code%{'key': keyresult, 'root': root})
else:
item1 = reduce(np.add, self.dview[key1])
item2 = self.e0[key2]
if individual:
item = (item2.conj() * item1).sum(axis=1) / (item1.conj() * item1).sum(axis=1)
else:
item = (item2.conj() * item1).sum() / (item1.conj() * item1).sum()
self.dview[keyresult] = item
def remoteApplySrc(self, keyData, keySrc):
code = '%(keyData)s = %(keySrc)s * %(keyData)s'
self.dview.execute(code%{'keyData': keyData, 'keySrc': keySrc})
# def normFromDifference(self, key):
# code = 'temp_norm%(key)s = (%(key)s * %(key)s.conj()).sum(0).sum(0)'
# self.e0.execute(code%{'key': key})
# code = 'temp_norm%(key)s = {key: np.sqrt(temp_norm%(key)s[key]).real for key in temp_norm%(key)s.keys()}'
# self.e0.execute(code%{'key': key})
# result = CommonReducer(self.e0['temp_norm%s'%(key,)])
# self.e0.execute('del temp_norm%s'%(key,))
# return result
def normFromDifference(self, key):
code = 'temp_norm = (%(key)s * %(key)s.conj()).sum(0).sum(0)'
self.e0.execute(code%{'key': key})
code = 'temp_norm = {key: np.sqrt(temp_norm[key]).real for key in temp_norm}'
self.e0.execute(code%{'key': key})
result = CommonReducer(self.e0['temp_norm'])
self.e0.execute('del temp_norm')
return result
@staticmethod
@interactive
def _reduceJob(worker, root, endpoint, key):
from ipyparallel.error import UnmetDependency
if not rank == worker:
raise UnmetDependency
# code = '%(endpoint)s.globalFields["%(key)s"] = comm.reduce(%(endpoint)s.localFields["%(key)s"], root=%(root)d)'
# exec(code%{'endpoint': endpoint, 'key': key, 'root': root})
if key not in endpoint.localFields:
endpoint.localFields[key] = endpoint.fieldspec[key]()
endpoint.globalFields[key] = comm.reduce(endpoint.localFields[key], root=root)
@staticmethod
def _adjustMKLVectorization(nt=1):
try:
import mkl
mkl.set_num_threads(nt)
except ImportError:
pass
@staticmethod
def _cdSame(rc):
import os
dview = rc[:]
home = os.getenv('HOME')
cwd = os.getcwd()
@interactive
def cdrel(relpath):
import os
home = os.getenv('HOME')
fullpath = os.path.join(home, relpath)
try:
os.chdir(fullpath)
except OSError:
return False
else:
return True
if cwd.find(home) == 0:
relpath = cwd[len(home)+1:]
return all(rc[:].apply_sync(cdrel, relpath))
@@ -0,0 +1,75 @@
<div id="%(uniqueID)s"></div>
<style>
.node {stroke: #fff; stroke-width: 1.5px;}
.link {stroke: #999; stroke-opacity: .6;}
</style>
<script>
var True = true;
var False = false;
var graph = %(JSONData)s;
require.config({paths: {d3: "http://d3js.org/d3.v3.min"}});
require(["d3"], function(d3) {
var width = 800, height = 450, radius = 8;
var color = d3.scale.category10();
var domain = [0, 1, 2, 3];
color.domain(domain);
var force = d3.layout.force()
.charge(-120)
.linkDistance(30)
.size([width, height]);
var svg = d3.select("#%(uniqueID)s").select("svg");
if (svg.empty()) {
svg = d3.select("#%(uniqueID)s").append("svg")
.attr("width", width)
.attr("height", height);
}
force.nodes(graph.nodes)
.links(graph.links)
.start();
var link = svg.selectAll(".link")
.data(graph.links)
.enter().append("line")
.attr("class", "link");
var node = svg.selectAll(".node")
.data(graph.nodes)
.enter().append("circle")
.attr("class", "node")
.attr("r", radius)
.style("fill", function(d) {
return color(d.status);
})
.call(force.drag);
node.append("title")
.text(function(d) { return d.id; });
node.on("dblclick", function() {
n = d3.select(this);
name = n.text();
graph = IPython.notebook.get_selected_cell().get_text();
cell = IPython.notebook.insert_cell_below();
cell.set_text(graph + ".node['" + name + "'].get('jobs', [])");
})
force.on("tick", function() {
link.attr("x1", function(d) { return d.source.x; })
.attr("y1", function(d) { return d.source.y; })
.attr("x2", function(d) { return d.target.x; })
.attr("y2", function(d) { return d.target.y; });
node.attr("cx", function(d) { return d.x; })
.attr("cy", function(d) { return d.y; });
});
});
</script>
+91
View File
@@ -0,0 +1,91 @@
class CommonReducer(dict):
'''
Object based on 'dict' that implements the binary addition (obj1 + obj1) and
accumulation (obj += obj2). These operations pass through to the entries in
the commonReducer.
Instances of commonReducer are also callable, with the syntax:
cr(key, value)
this is equivalent to cr += {key: value}.
'''
DISALLOWED = ['__getinitargs__', '__getnewargs__', '__getstate__', '__setstate__']
def __init__(self, *args, **kwargs):
dict.__init__(self, *args, **kwargs)
def __add__(self, other):
result = CommonReducer(self)
for key in other.keys():
if key in result:
result[key] = self[key] + other[key]
else:
result[key] = other[key]
return result
def __iadd__(self, other):
for key in other.keys():
if key in self:
self[key] += other[key]
else:
self[key] = other[key]
return self
def __mul__(self, other):
result = CommonReducer()
for key in other.keys():
if key in self:
result[key] = self[key] * other[key]
return result
def __sub__(self, other):
result = CommonReducer()
for key in other.keys():
if key in self:
result[key] = self[key] - other[key]
return result
def __div__(self, other):
result = CommonReducer()
for key in other.keys():
if key in self:
result[key] = self[key] / other[key]
return result
def __getattr__(self, attr):
if not attr in self.DISALLOWED and all((getattr(self[key], attr, None) is not None for key in self)):
if any((callable(getattr(self[key], attr)) for key in self)):
def wrapperFunction(*args, **kwargs):
innerresult = CommonReducer({key: getattr(self[key], attr, None)(*args, **kwargs) for key in self})
if not all((innerresult[key] is None for key in innerresult)):
return innerresult
result = wrapperFunction
else:
return CommonReducer({key: getattr(self[key], attr) for key in self})
else:
raise AttributeError('\'CommonReducer\' object has no attribute \'%s\', and it could not be satisfied through cascade lookup'%attr)
return result
def copy(self):
return CommonReducer(self)
def __call__(self, key, result):
if key in self:
self[key] += result
else:
self[key] = result
+1
View File
@@ -5,6 +5,7 @@ from curvutils import volTetra, faceInfo, indexCube
from interputils import interpmat
from ipythonutils import easyAnimate as animate
from CounterUtils import *
from DataWrappers import *
import ModelBuilder
import SolverUtils
+1
View File
@@ -13,6 +13,7 @@ import InvProblem
import Optimization
import Directives
import Inversion
import Parallel
import Tests
__version__ = '0.1.3'