Skip to content
Open
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: 2 additions & 2 deletions astra.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ inputs:
outputs:
- id: final_cat
type: data
format: fits
format: hdf5
description: >-
Per-tile shear catalogue family, the terminal science product
(make_cat_runner).
Expand Down Expand Up @@ -2024,7 +2024,7 @@ analyses:
outputs:
- id: tile_final_cat
type: data
format: fits
format: hdf5
description: The assembled per-tile science catalogue (final_cat family).
inputs: [ngmix_chunks]
decisions:
Expand Down
72 changes: 43 additions & 29 deletions scripts/python/create_final_cat.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,9 @@

"""Script create_final_cat.py

Create and update hdf5 file of all final ShapePipe output FITS files, runs of
ShapePipe module ``make_catalogue_runner``. Supercedes `merge_final_cat.py`.
Create and update hdf5 file of all per-tile final catalogues, the output of
ShapePipe module ``make_cat_runner`` (hdf5, one dataset per column; FITS for
catalogues written before that format). Supercedes `merge_final_cat.py`.

Usage: in parent dir of patches:
create_final_cat.py -p ~/shapepipe/workflow/config/cfis/final_cat.param -i . -P 7 -v -m final_cat_P7.hdf5
Expand Down Expand Up @@ -374,9 +375,14 @@ def get_patch_group(hdf5_file, patch, verbose=False):
return patch_group


def read_data(fits_file, params):
def read_data(cat_file, params):
"""Read the parameter list's columns out of one catalogue.

The catalogue is hdf5, one dataset per column (a vector column is 2-D),
or, for one with a ``.fits`` suffix, a FITS table at HDU
``params["hdu_num"]``. Returns the requested columns and a structured
dtype describing them.

@sc [label:schema] read-data-raises-on-missing-column
A requested column the catalogue lacks raises `KeyError` naming it; it is
never skipped or filled. `copy_data` keeps only columns present in the
Expand All @@ -385,31 +391,29 @@ def read_data(fits_file, params):
silently narrower, with that slot's exposure identity gone. Enforced by
tests/unit/test_final_cat_merge_invariants.py.
"""
with fits.open(fits_file) as hdu_list:
try:
if cat_file.endswith(".fits"):
with fits.open(cat_file) as hdu_list:
data = hdu_list[params["hdu_num"]].data
except:
print(f"Error with ID {id}, file{fits_file}")
raise
dtype = data.dtype
else:
with h5py.File(cat_file, "r") as cat:
dtype = np.dtype([(col, ds.dtype, ds.shape[1:])
for col, ds in cat.items()])
wanted = params["param_list"] or dtype.names
data = {col: cat[col][()] for col in wanted if col in cat}

# If columns not given on input: read all column names
if params["param_list"] is None:
params["param_list"] = [col for col in data.keys()]

# RAISE, do not print and fall through. The bare `except:` this replaces
# left extracted_data and dtype unbound, so the caller's own error was an
# UnboundLocalError from the return statement below, naming neither the
# file nor the column that was actually missing.
present = set(data.dtype.names or ())
missing = [col for col in params["param_list"] if col not in present]
params["param_list"] = list(dtype.names)

missing = [col for col in params["param_list"] if col not in dtype.names]
if missing:
raise KeyError(
f"{fits_file}: missing {len(missing)} of the "
f"{cat_file}: missing {len(missing)} of the "
f"{len(params['param_list'])} requested column(s): "
f"{' '.join(missing)}"
)
extracted_data = {col: data[col] for col in params["param_list"]}
dtype = data.dtype

return extracted_data, dtype

Expand Down Expand Up @@ -496,26 +500,36 @@ def collect_tile_ids_image_sims(patch_path):
return result


# make_cat_runner writes hdf5; catalogues made before that are FITS.
FINAL_CAT_SUFFIXES = (".hdf5", ".fits")


def find_final_cat(id, id_path, run_prefix):
"""Path of this tile's final catalogue, or None.

Flat layout first: the unified workflow copies one catalogue per tile to
<products>/tiles/<shard>/<tile>/final_cat-<tile>.fits, in DOT form and with
no run sub-tree. Then the legacy layout, where the catalogue sits under the
tile's own shapepipe run dir in DASH form, newest run wins.
<products>/tiles/<shard>/<tile>/final_cat-<tile>.hdf5, in DOT form and
with no run sub-tree. Then the run-dir layout, where the catalogue sits
under the tile's own shapepipe run dir in DASH form, newest run wins.
Either layout may hold an hdf5 or a FITS catalogue; hdf5 is preferred.
"""
flat = os.path.join(id_path, f"final_cat-{id}.fits")
if os.path.exists(flat):
return flat
for suffix in FINAL_CAT_SUFFIXES:
flat = os.path.join(id_path, f"final_cat-{id}{suffix}")
if os.path.exists(flat):
return flat

base_pattern = os.path.join(id_path, "output", run_prefix)
all_matches = [d for d in glob.glob(base_pattern) if os.path.isdir(d)]
if not all_matches:
return None
newest_dir = max(all_matches, key=os.path.getmtime)
id_dash = re.sub(r"\.", "-", id)
legacy = f"{newest_dir}/make_cat_runner/output/final_cat-{id_dash}.fits"
return legacy if os.path.exists(legacy) else None
for suffix in FINAL_CAT_SUFFIXES:
in_run = (f"{newest_dir}/make_cat_runner/output/"
f"final_cat-{id_dash}{suffix}")
if os.path.exists(in_run):
return in_run
return None


def process(params):
Expand Down Expand Up @@ -577,13 +591,13 @@ def process(params):
print(f"Skipping {id} (already processed)")
continue

fits_file = find_final_cat(id, id_path, run_prefix)
if fits_file is None:
cat_file = find_final_cat(id, id_path, run_prefix)
if cat_file is None:
if params["verbose"]:
print(f"Final cat for {id} not found, continuing")
continue

extracted_data, dtype = read_data(fits_file, params)
extracted_data, dtype = read_data(cat_file, params)

structured_data = copy_data(params["param_list"], extracted_data, dtype)

Expand Down
Loading
Loading