LSST Applications  21.0.0-172-gfb10e10a+18fedfabac,22.0.0+297cba6710,22.0.0+80564b0ff1,22.0.0+8d77f4f51a,22.0.0+a28f4c53b1,22.0.0+dcf3732eb2,22.0.1-1-g7d6de66+2a20fdde0d,22.0.1-1-g8e32f31+297cba6710,22.0.1-1-geca5380+7fa3b7d9b6,22.0.1-12-g44dc1dc+2a20fdde0d,22.0.1-15-g6a90155+515f58c32b,22.0.1-16-g9282f48+790f5f2caa,22.0.1-2-g92698f7+dcf3732eb2,22.0.1-2-ga9b0f51+7fa3b7d9b6,22.0.1-2-gd1925c9+bf4f0e694f,22.0.1-24-g1ad7a390+a9625a72a8,22.0.1-25-g5bf6245+3ad8ecd50b,22.0.1-25-gb120d7b+8b5510f75f,22.0.1-27-g97737f7+2a20fdde0d,22.0.1-32-gf62ce7b1+aa4237961e,22.0.1-4-g0b3f228+2a20fdde0d,22.0.1-4-g243d05b+871c1b8305,22.0.1-4-g3a563be+32dcf1063f,22.0.1-4-g44f2e3d+9e4ab0f4fa,22.0.1-42-gca6935d93+ba5e5ca3eb,22.0.1-5-g15c806e+85460ae5f3,22.0.1-5-g58711c4+611d128589,22.0.1-5-g75bb458+99c117b92f,22.0.1-6-g1c63a23+7fa3b7d9b6,22.0.1-6-g50866e6+84ff5a128b,22.0.1-6-g8d3140d+720564cf76,22.0.1-6-gd805d02+cc5644f571,22.0.1-8-ge5750ce+85460ae5f3,master-g6e05de7fdc+babf819c66,master-g99da0e417a+8d77f4f51a,w.2021.48
LSST Data Management Base Package
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.logOperationlogOperation("master"):
54  total = pool.reduce(operator.add, self.runrun, 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.logOperationlogOperation("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
def logOperation(self, operation, catch=False, trace=True)
Provide a context manager for logging an operation.
Definition: parallel.py:502
def batchWallTime(cls, time, parsedCmd, numCores)
Definition: demoTask.py:26
def runDataRef(self, visitRef)
Definition: demoTask.py:40
def run(self, cache, dataId, butler=None)
Definition: demoTask.py:58
daf::base::PropertyList * list
Definition: fits.cc:913