Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions SConstruct
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
# -*- python -*-
from lsst.sconsUtils import scripts
scripts.BasicSConstruct("ImageProcessingPipelines",
versionModuleName='python/desc/ImageProcessingPipelines/version.py')
3 changes: 3 additions & 0 deletions bin.src/SConscript
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
# -*- python -*-
from lsst.sconsUtils import scripts
scripts.BasicSConscript.shebang()
66 changes: 66 additions & 0 deletions bin.src/compute_overlaps.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
#!/usr/bin/env python
"""
Script to compute overlaps of sensor-visits with a skymap.
"""
import argparse
import multiprocessing
import numpy as np
import pandas as pd
import lsst.daf.persistence as dp
from desc.ImageProcessingPipelines import SkyMapPolygons, OverlapFinder


description = 'Compute overlaps of sensor-visits with a skymap.'

parser = argparse.ArgumentParser(description=description)
parser.add_argument('--repo', type=str, default=None,
help='Data repository containing the skymap object')
parser.add_argument('--opsim_db_file', type=str, default=None,
help='OpSim db file')
parser.add_argument('--opsim_constraint', type=str, default=None,
help='Selection constraint on Summary table in opsim db')
parser.add_argument('--processes', type=int, default=1,
help='Number of processes to run concurrently')
parser.add_argument('--outfile', type=str, default='overlaps.pickle',
help='Name of output pickle file to contain the DataFrame')
args = parser.parse_args()

repo = args.repo if args.repo is not None else \
('/global/cfs/cdirs/lsst/production/DC2_ImSim/Run2.2i'
'/desc_dm_drp/v19.0.0-v1/rerun/run2.2i-coadd-wfd-dr6-v1')

opsim_db_file = args.opsim_db_file if args.opsim_db_file is not None else \
('/global/cfs/cdirs/descssim/DC2'
'/minion_1016_desc_dithered_v4_trimmed.db')

butler = dp.Butler(repo)
skymap = butler.get('deepCoadd_skyMap')
skymap_polygons = SkyMapPolygons(skymap)
overlap_finder = OverlapFinder(opsim_db_file, skymap_polygons)

if args.opsim_constraint is not None:
df = overlap_finder.opsim_db.query(args.opsim_constraint)
else:
df = overlap_finder.opsim_db

all_visits = list(df['obsHistID'])

if args.processes == 1:
df = overlap_finder.get_overlaps(all_visits)
else:
indexes = np.linspace(0, len(all_visits), args.processes + 1, dtype=int)
visit_lists = [all_visits[imin:imax] for imin, imax in
zip(indexes[:-1], indexes[1:])]

with multiprocessing.Pool(processes=args.processes) as pool:
workers = []
for visits in visit_lists:
workers.append(pool.apply_async(overlap_finder.get_overlaps,
(visits,)))
pool.close()
pool.join()
dfs = [_.get() for _ in workers]

df = pd.concat(dfs)

df.to_pickle(args.outfile)
47 changes: 47 additions & 0 deletions bin.src/drp_resource_estimator.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
#!/usr/bin/env python
"""
Script to estimate computing resources for DRP image processing.
"""
import os
import argparse
import sqlite3
import pandas as pd
from desc.ImageProcessingPipelines import extract_coadds, \
tabulate_pipe_task_resources, total_node_hours


parser = argparse.ArgumentParser(description=('Computing resource estimator '
'for DRP processing'))
parser.add_argument('overlaps_db_file', type=str,
help='sqlite3 file containing the overlaps table')
parser.add_argument('--coadd_df_file', type=str, default='coadd_df.pickle',
help='pickle file containing the coadd summary dataframe')
parser.add_argument('--knl_factor', type=float, default=8,
help='slow-down factor for running on KNL vs Haswell')
parser.add_argument('--verbose', default=True, action='store_false',
help='option to enable verbose output')

args = parser.parse_args()

if args.overlaps_db_file.endswith('.pickle'):
df = pd.read_pickle(args.overlaps_db_file)
else:
with sqlite3.connect(args.overlaps_db_file) as con:
df = pd.read_sql('select * from overlaps', con)

# coadds in each band with # visits per coadd:
if not os.path.isfile(args.coadd_df_file):
coadd_df = extract_coadds(df, verbose=args.verbose)
coadd_df.to_pickle(args.coadd_df_file)
else:
coadd_df = pd.read_pickle(args.coadd_df_file)

pt_df = tabulate_pipe_task_resources(df, coadd_df, verbose=args.verbose)

node_hours, node_hours_opt = total_node_hours(pt_df,
cpu_factor=args.knl_factor)

print()
print(pt_df)
print()
print(f'KNL node days: {node_hours/24.:.1f} ({node_hours_opt/24.:.1f})')
1 change: 0 additions & 1 deletion python/__init__.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,2 @@
from __future__ import absolute_import
import pkgutil
__path__ = pkgutil.extend_path(__path__, __name__)
6 changes: 6 additions & 0 deletions python/desc/ImageProcessingPipelines/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
try:
from version import *
except ImportError:
pass
from .tabulate_pipe_task_resources import *
from .visit_skymap_overlaps import *
110 changes: 110 additions & 0 deletions python/desc/ImageProcessingPipelines/pipe_task_resource_usage.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
"""
Resource usage estimates for DRP pipe_tasks as run on a
Cori-Haswell node as of June 2020. See
https://github.com/LSSTDESC/gen3_workflow/wiki/Resource-Usage-Testing-with-Gen2-pipe_tasks
"""
__all__ = ['pipe_tasks']


def nImages(num_visits):
"""
Convert the number of visits to the mean number of images
contributing to a given patch. This is the average conversion
ratio for a WFD survey candence.
"""
return 0.77*num_visits


def processCcd():
"""
Return the average cpu time in hours and required memory in GB for
a single instance of processCcd.
"""
return 2.5/60., 1.


def makeCoaddTempExp():
"""
Return the average cpu time in hours and required memory in GB for
a single instance of makeCoaddTempExp.
"""
return 2.2/60., 1.6


def assembleCoadd(num_visits):
"""
Return the average cpu time in hours and required memory in GB for
a single instance of assembleCoadd, scaling with the number of visits.
"""
return 0.4*num_visits/60., 1.5


def cpu_mem_visit_scaling(num_visits, cpu0, cpu_index, mem0, mem_index):
"""
Return tuple of cpu time (hours), required memory (GB) for pipe tasks
that operate on coadds. Power-law scalings were fit to cpu_mins and
mem_GB values found from DR6 studies on a Cori-Haswell node.
"""
n_images = nImages(num_visits)
cpu_mins = cpu0*n_images**cpu_index
mem_GB = mem0*n_images**mem_index
return cpu_mins/60., mem_GB


def detectCoaddSources(num_visits):
"""
Return a tuple of cpu hours and required memory (in GB) for
detectCoaddSources using the fit parameters found from DR6 studies.
"""
return cpu_mem_visit_scaling(num_visits, 0.23, 0.78, 0.68, 0.26)


def mergeCoaddDetections():
"""
Return a tuple of cpu hours and required memory (in GB) for
mergeCoaddDetections. These are the worst case values found
from studies of DR6 data run on a Cori-Haswell node.
"""
return 2.7/60., 0.7


def deblendCoaddSources(num_visits):
"""
Return a tuple of cpu hours and required memory (in GB) for
deblendCoaddSources using the fit parameters found from DR6 studies.
"""
return cpu_mem_visit_scaling(num_visits, 0.16, 1.40, 0.30, 0.41)


def measureCoaddSources(num_visits):
"""
Return a tuple of cpu hours and required memory (in GB) for
measureCoaddSources using the fit parameters found from DR6 studies.
"""
return cpu_mem_visit_scaling(num_visits, 1.80, 1.20, 0.43, 0.36)


def mergeCoaddMeasurements():
"""
Return a tuple of cpu hours and required memory (in GB) for
mergeCoaddMeasurements. These are the worst case values found
from studies of DR6 data run on a Cori-Haswell node.
"""
return 1/60., 2.8


def forcedPhotCoadd(num_visits):
"""
Return a tuple of cpu hours and required memory (in GB) for
forcedPhotCoadd using the fit parameters found from DR6 studies.
"""
return cpu_mem_visit_scaling(num_visits, 2.20, 1.20, 0.43, 0.36)


pipe_task_funcs = ['processCcd', 'makeCoaddTempExp', 'assembleCoadd',
'detectCoaddSources', 'mergeCoaddDetections',
'deblendCoaddSources', 'measureCoaddSources',
'mergeCoaddMeasurements', 'forcedPhotCoadd']


pipe_tasks = {_: eval(_) for _ in pipe_task_funcs}
Loading