LSST Applications  21.0.0+04719a4bac,21.0.0-1-ga51b5d4+f5e6047307,21.0.0-11-g2b59f77+a9c1acf22d,21.0.0-11-ga42c5b2+86977b0b17,21.0.0-12-gf4ce030+76814010d2,21.0.0-13-g1721dae+760e7a6536,21.0.0-13-g3a573fe+768d78a30a,21.0.0-15-g5a7caf0+f21cbc5713,21.0.0-16-g0fb55c1+b60e2d390c,21.0.0-19-g4cded4ca+71a93a33c0,21.0.0-2-g103fe59+bb20972958,21.0.0-2-g45278ab+04719a4bac,21.0.0-2-g5242d73+3ad5d60fb1,21.0.0-2-g7f82c8f+8babb168e8,21.0.0-2-g8f08a60+06509c8b61,21.0.0-2-g8faa9b5+616205b9df,21.0.0-2-ga326454+8babb168e8,21.0.0-2-gde069b7+5e4aea9c2f,21.0.0-2-gecfae73+1d3a86e577,21.0.0-2-gfc62afb+3ad5d60fb1,21.0.0-25-g1d57be3cd+e73869a214,21.0.0-3-g357aad2+ed88757d29,21.0.0-3-g4a4ce7f+3ad5d60fb1,21.0.0-3-g4be5c26+3ad5d60fb1,21.0.0-3-g65f322c+e0b24896a3,21.0.0-3-g7d9da8d+616205b9df,21.0.0-3-ge02ed75+a9c1acf22d,21.0.0-4-g591bb35+a9c1acf22d,21.0.0-4-g65b4814+b60e2d390c,21.0.0-4-gccdca77+0de219a2bc,21.0.0-4-ge8a399c+6c55c39e83,21.0.0-5-gd00fb1e+05fce91b99,21.0.0-6-gc675373+3ad5d60fb1,21.0.0-64-g1122c245+4fb2b8f86e,21.0.0-7-g04766d7+cd19d05db2,21.0.0-7-gdf92d54+04719a4bac,21.0.0-8-g5674e7b+d1bd76f71f,master-gac4afde19b+a9c1acf22d,w.2021.13
LSST Data Management Base Package
ingestDriver.py
Go to the documentation of this file.
1 from lsst.ctrl.pool.pool import Pool, startPool, abortOnError
2 from lsst.ctrl.pool.parallel import BatchCmdLineTask
3 from lsst.pipe.base import Struct
4 from lsst.pipe.tasks.ingest import IngestTask, IngestError
5 
6 
8  """Parallel version of IngestTask"""
9  @classmethod
10  def batchWallTime(cls, time, parsedCmd, numCores):
11  return float(time)*len(parsedCmd.files)/numCores
12 
13  @classmethod
14  def _makeArgumentParser(cls, *args, **kwargs):
15  """Build an ArgumentParser
16 
17  Removes the batch-specific parts.
18  """
19  kwargs.pop("doBatch", False)
20  kwargs.pop("add_help", False)
21  return cls.ArgumentParserArgumentParser(*args, name="ingest", **kwargs)
22 
23  @classmethod
24  def parseAndRun(cls, *args, **kwargs):
25  """Run with a MPI process pool"""
26  pool = startPool()
27  config = cls.ConfigClassConfigClass()
28  parser = cls.ArgumentParserArgumentParser(name=cls._DefaultName_DefaultName)
29  args = parser.parse_args(config)
30  task = cls(config=args.config)
31  task.run(args)
32  pool.exit()
33 
34  def runFileWrapper(self, struct, args):
35  """Run ingest on one file
36 
37  This is a wrapper method for calling ``runFile``.
38 
39  Parameters
40  ----------
41  struct : `lsst.pipe.base.Struct`
42  Structure containing ``filename`` (`str`) and ``position`` (`int`).
43  args : `argparse.Namespace`
44  Parsed command-line arguments.
45 
46  Returns
47  -------
48  hduInfoList : `list` of `dict`
49  Parsed information from FITS HDUs, or ``None``.
50  """
51  filename = struct.filename
52  position = struct.position
53  try:
54  return self.runFilerunFile(filename, None, args, position)
55  except IngestError as exc:
56  self.loglog.warn(f"Unable to ingest {filename}: {exc}")
57  return None
58 
59  @abortOnError
60  def run(self, args):
61  """Run ingest
62 
63  We read and ingest the files in parallel, and then
64  stuff the registry database in serial.
65  """
66  # Parallel
67  pool = Pool(None)
68  filenameList = self.expandFilesexpandFiles(args.files)
69  dataList = [Struct(filename=filename, position=ii) for ii, filename in enumerate(filenameList)]
70  infoList = pool.map(self.runFileWrapperrunFileWrapper, dataList, args)
71 
72  # Serial
73  root = args.input
74  context = self.register.openRegistry(root, create=args.create, dryrun=args.dryrun)
75  with context as registry:
76  for hduInfoList in infoList:
77  if hduInfoList is None:
78  continue
79  for info in hduInfoList:
80  self.register.addRow(registry, info, dryrun=args.dryrun, create=args.create)
81 
82  def writeConfig(self, *args, **kwargs):
83  pass
84 
85  def writeMetadata(self, *args, **kwargs):
86  pass
def writeConfig(self, *args, **kwargs)
Definition: ingestDriver.py:82
def batchWallTime(cls, time, parsedCmd, numCores)
Return walltime request for batch job.
Definition: ingestDriver.py:10
def writeMetadata(self, *args, **kwargs)
Definition: ingestDriver.py:85
def expandFiles(self, fileNameList)
Expand a set of filenames and globs, returning a list of filenames.
Definition: ingest.py:557
def runFile(self, infile, registry, args, pos)
Examine and ingest a single file.
Definition: ingest.py:576
def startPool(comm=None, root=0, killSlaves=True)
Start a process pool.
Definition: pool.py:1216