LSSTApplications  19.0.0-14-gb0260a2+72efe9b372,20.0.0+7927753e06,20.0.0+8829bf0056,20.0.0+995114c5d2,20.0.0+b6f4b2abd1,20.0.0+bddc4f4cbe,20.0.0-1-g253301a+8829bf0056,20.0.0-1-g2b7511a+0d71a2d77f,20.0.0-1-g5b95a8c+7461dd0434,20.0.0-12-g321c96ea+23efe4bbff,20.0.0-16-gfab17e72e+fdf35455f6,20.0.0-2-g0070d88+ba3ffc8f0b,20.0.0-2-g4dae9ad+ee58a624b3,20.0.0-2-g61b8584+5d3db074ba,20.0.0-2-gb780d76+d529cf1a41,20.0.0-2-ged6426c+226a441f5f,20.0.0-2-gf072044+8829bf0056,20.0.0-2-gf1f7952+ee58a624b3,20.0.0-20-geae50cf+e37fec0aee,20.0.0-25-g3dcad98+544a109665,20.0.0-25-g5eafb0f+ee58a624b3,20.0.0-27-g64178ef+f1f297b00a,20.0.0-3-g4cc78c6+e0676b0dc8,20.0.0-3-g8f21e14+4fd2c12c9a,20.0.0-3-gbd60e8c+187b78b4b8,20.0.0-3-gbecbe05+48431fa087,20.0.0-38-ge4adf513+a12e1f8e37,20.0.0-4-g97dc21a+544a109665,20.0.0-4-gb4befbc+087873070b,20.0.0-4-gf910f65+5d3db074ba,20.0.0-5-gdfe0fee+199202a608,20.0.0-5-gfbfe500+d529cf1a41,20.0.0-6-g64f541c+d529cf1a41,20.0.0-6-g9a5b7a1+a1cd37312e,20.0.0-68-ga3f3dda+5fca18c6a4,20.0.0-9-g4aef684+e18322736b,w.2020.45
LSSTDataManagementBasePackage
demoTask.py
Go to the documentation of this file.
1 import math
2 import collections
3 import operator
4 from lsst.ctrl.pool.parallel import BatchPoolTask
5 from lsst.ctrl.pool.pool import Pool
6 from lsst.pipe.base import ArgumentParser
7 from lsst.pex.config import Config
8 
9 __all__ = ["DemoTask", ]
10 
11 
13  """Task for demonstrating the BatchPoolTask functionality"""
14  ConfigClass = Config
15  _DefaultName = "demo"
16 
17  @classmethod
18  def _makeArgumentParser(cls, *args, **kwargs):
19  kwargs.pop('doBatch', False) # Unused
20  parser = ArgumentParser(name="demo", *args, **kwargs)
21  parser.add_id_argument("--id", datasetType="raw", level="visit",
22  help="data ID, e.g. --id visit=12345")
23  return parser
24 
25  @classmethod
26  def batchWallTime(cls, time, parsedCmd, numCores):
27  """Return walltime request for batch job
28 
29  Subclasses should override if the walltime should be calculated
30  differently (e.g., addition of some serial time).
31 
32  @param time: Requested time per iteration
33  @param parsedCmd: Results of argument parsing
34  @param numCores: Number of cores
35  """
36  numTargets = [sum(1 for ccdRef in visitRef.subItems("ccd") if ccdRef.datasetExists("raw")) for
37  visitRef in parsedCmd.id.refList]
38  return time*sum(math.ceil(tt/numCores) for tt in numTargets)
39 
40  def runDataRef(self, visitRef):
41  """Main entry-point
42 
43  Only the master node runs this method. It will dispatch jobs to the
44  slave nodes.
45  """
46  pool = Pool("test")
47 
48  # Less overhead to transfer the butler once rather than in each dataRef
49  dataIdList = dict([(ccdRef.get("ccdExposureId"), ccdRef.dataId)
50  for ccdRef in visitRef.subItems("ccd") if ccdRef.datasetExists("raw")])
51  dataIdList = collections.OrderedDict(sorted(dataIdList.items()))
52 
53  with self.logOperation("master"):
54  total = pool.reduce(operator.add, self.run, list(dataIdList.values()),
55  butler=visitRef.getButler())
56  self.log.info("Total number of pixels read: %d" % (total,))
57 
58  def run(self, cache, dataId, butler=None):
59  """Read image and return number of pixels
60 
61  Only the slave nodes run this method.
62  """
63  assert butler is not None
64  with self.logOperation("read %s" % (dataId,)):
65  raw = butler.get("raw", dataId, immediate=True)
66  dims = raw.getDimensions()
67  num = dims.getX()*dims.getY()
68  self.log.info("Read %d pixels for %s" % (num, dataId,))
69  return num
70 
71  def _getConfigName(self):
72  return None
73 
74  def _getMetadataName(self):
75  return None
76 
77  def _getEupsVersionsName(self):
78  return None
lsst.ctrl.pool.test.demoTask.DemoTask.batchWallTime
def batchWallTime(cls, time, parsedCmd, numCores)
Definition: demoTask.py:26
lsst::log.log.logContinued.info
def info(fmt, *args)
Definition: logContinued.py:201
lsst.ctrl.pool.test.demoTask.DemoTask
Definition: demoTask.py:12
lsst.pipe.base.argumentParser.ArgumentParser
Definition: argumentParser.py:408
lsst.ctrl.pool.parallel
Definition: parallel.py:1
lsst.ctrl.pool.test.demoTask.DemoTask.runDataRef
def runDataRef(self, visitRef)
Definition: demoTask.py:40
lsst.ctrl.pool.test.demoTask.DemoTask.run
def run(self, cache, dataId, butler=None)
Definition: demoTask.py:58
lsst.pex.config
Definition: __init__.py:1
lsst.ctrl.pool.parallel.BatchCmdLineTask.logOperation
def logOperation(self, operation, catch=False, trace=True)
Provide a context manager for logging an operation.
Definition: parallel.py:502
lsst.pipe.base.task.Task.log
log
Definition: task.py:161
lsst.ctrl.pool.pool
Definition: pool.py:1
lsst.ctrl.pool.pool.Pool
Definition: pool.py:1205
list
daf::base::PropertyList * list
Definition: fits.cc:913
lsst.ctrl.pool.parallel.BatchPoolTask
Definition: parallel.py:527
lsst.pipe.base
Definition: __init__.py:1