Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
63 commits
Select commit Hold shift + click to select a range
e706906
run in subcalcs
CB-quakemodel Jul 31, 2026
ea20261
fix sampling
CB-quakemodel Jul 31, 2026
a068bb5
new QA test instead
CB-quakemodel Jul 31, 2026
46c0506
changelog
CB-quakemodel Jul 31, 2026
0dde90c
clean
CB-quakemodel Jul 31, 2026
738b911
test quantiles too
CB-quakemodel Jul 31, 2026
8e911e5
atol
CB-quakemodel Jul 31, 2026
9263d11
error if not classical or disagg
CB-quakemodel Jul 31, 2026
5bd21f6
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Jul 31, 2026
6898bfb
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 1, 2026
9a46e0f
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 2, 2026
bbd9cdd
sequential approach
CB-quakemodel Aug 3, 2026
466ff6d
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 3, 2026
7b4c5f1
cleanup
CB-quakemodel Aug 3, 2026
33e30b4
remove sort
CB-quakemodel Aug 3, 2026
f613b26
revert
CB-quakemodel Aug 3, 2026
6b31a82
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 4, 2026
36c7b16
remove dup code
CB-quakemodel Aug 4, 2026
0822dea
increase src_grp limit (probably bad idea)
CB-quakemodel Aug 4, 2026
d8fcc6d
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 4, 2026
916d516
clean
CB-quakemodel Aug 5, 2026
eef2a45
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 5, 2026
8e15c13
more cleaning
CB-quakemodel Aug 5, 2026
d179092
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 5, 2026
2a1be44
clean
CB-quakemodel Aug 5, 2026
04d70e6
revert
CB-quakemodel Aug 5, 2026
606f910
revert
CB-quakemodel Aug 5, 2026
e41be4a
revert
CB-quakemodel Aug 5, 2026
8a6c61d
revert
CB-quakemodel Aug 5, 2026
d9deb52
try splitting preclassical too
CB-quakemodel Aug 5, 2026
5b5904c
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 5, 2026
cf65d13
unit test
CB-quakemodel Aug 5, 2026
78e6fae
helper func
CB-quakemodel Aug 5, 2026
82cb0f0
widen task_no dtype
CB-quakemodel Aug 6, 2026
b5281e0
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 6, 2026
e3f2aad
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 6, 2026
f90eced
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 9, 2026
3b4d31d
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 10, 2026
a2b863d
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 11, 2026
10fe788
work with extendModel
CB-quakemodel Aug 11, 2026
c29f2ca
readme
CB-quakemodel Aug 11, 2026
13ac4a2
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 11, 2026
5858de4
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 13, 2026
f65de2f
test - widen dtype of gid to u32
CB-quakemodel Aug 13, 2026
62b9d31
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 14, 2026
5bc5fb3
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 14, 2026
9e45266
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 15, 2026
5d1a759
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 15, 2026
7cad9ca
remove dtype changes
CB-quakemodel Aug 18, 2026
ed44fb5
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 18, 2026
f70b2b4
dtypes widened
CB-quakemodel Aug 18, 2026
d94c95e
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 20, 2026
d959050
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 21, 2026
56025cc
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 22, 2026
beef244
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 28, 2026
e46ba78
revert dtype widen
CB-quakemodel Aug 31, 2026
52fce15
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Aug 31, 2026
6c0b6cc
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Sep 1, 2026
eafbaf7
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Sep 1, 2026
6acada0
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Sep 1, 2026
8d6c5bd
fix merge errors
CB-quakemodel Sep 11, 2026
3485729
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Sep 12, 2026
1a7c20a
Merge branch 'master' of github.com:gem/oq-engine into subcalc
CB-quakemodel Sep 14, 2026
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
55 changes: 50 additions & 5 deletions openquake/calculators/classical.py
Original file line number Diff line number Diff line change
Expand Up @@ -645,14 +645,59 @@ def _execute(self, sgs, ds):
OQ_TASK_NO = os.environ.get('OQ_TASK_NO', '')
if OQ_TASK_NO:
allargs = [allargs[int(OQ_TASK_NO)]]
if self.few_sites or oq.disagg_by_src:
smap = parallel.Starmap(
classical_disagg, allargs, h5=self.datastore.hdf5)
task_func = (classical_disagg if (self.few_sites or oq.disagg_by_src)
else classical)
if oq.sequential_source_models and not OQ_TASK_NO:
acc = self._run_sequential_source_models(allargs, task_func)
else:
smap = parallel.Starmap(classical, allargs, h5=self.datastore.hdf5)
acc = smap.reduce(self.agg_dicts, AccumDict(accum=0.))
smap = parallel.Starmap(
task_func, allargs, h5=self.datastore.hdf5)
acc = smap.reduce(self.agg_dicts, AccumDict(accum=0.))
self._post_execute(acc)

def _run_sequential_source_models(self, allargs, task_func):
"""
Run one Starmap per sourceModel branch sequentially, so that
only a single source model's tasks are running at one time.

NOTE: Unless extendModel is being used, the sharing of src_groups
over base source models is not permitted. If extendModel is being
used, then the sharing of src_groups over base source models is
permitted and a "shared" batch is dispatched last. The memory
footprint of this "shared" batch could be similar (or equal) to
that observed in the "regular" (i.e., none-sequential) approach
if extendModel is heavily used in the logic tree (because many
of the src_grps would be piled into this final "shared" batch).
"""
# Map each src_group id to its sourceModel branch_id
smb_of_grp = {
gid: smb
for smb, gids in self.csm.grp_ids_by_source_model().items()
for gid in gids}

# Partition the task-arg tuples by sourceModel branch
partitions = AccumDict(accum=[])
for args in allargs:
# Strip any tile suffix ("5-2" -> 5) to recover grp_id
gid = int(args[0][0].split('-')[0])
partitions[smb_of_grp[gid]].append(args)

# Sort per-source model partitions for reproducible
# order with "shared" batch last
per_sm_keys = sorted(k for k in partitions if k is not None)
ordered_keys = per_sm_keys + (
[None] if None in partitions else [])

# Run one source model at a time
acc = AccumDict(accum=0.)
for smb in ordered_keys:
part = partitions[smb]
logging.info('Source model %r: %d tasks', smb, len(part))
smap = parallel.Starmap(task_func, part, h5=self.datastore.hdf5)
acc = smap.reduce(self.agg_dicts, acc)

return acc

def _post_execute(self, acc):
# save the rates and performs some checks
oq = self.oqparam
Expand Down
113 changes: 81 additions & 32 deletions openquake/calculators/preclassical.py
Original file line number Diff line number Diff line change
Expand Up @@ -298,22 +298,16 @@ def populate_csm(self):
self.store()
logging.info('Building cmakers')
trt_smrs = csm.get_trt_smrs()
self.cmakers = get_cmakers(trt_smrs, csm.full_lt, oq)
self.datastore.hdf5.save_vlen('trt_smrs', trt_smrs)
if oq.sequential_source_models:
grp_ids_by_batch = [
numpy.array(
[sg.sources[0].grp_id for sg in sgs_batch], U32)
for _, _, sgs_batch, _ in csm.iter_source_model_batches()]
self.datastore.hdf5.save_vlen('grp_ids_by_batch', grp_ids_by_batch)
sites = csm.sitecol if csm.sitecol else None
if sites is None:
logging.warning('No sites??')

L = oq.imtls.size
Gfull = self.full_lt.gfull([cm.trt_smrs for cm in self.cmakers])
Gt = sum(len(cm.gsims) for cm in self.cmakers)
extra = f'<{Gfull}' if Gt < Gfull else ''
if sites is not None:
nbytes = 4 * len(self.sitecol) * L * Gt
# Gt is known before starting the preclassical
logging.warning(f'Global RateMap of %s ({Gt=}%s)',
general.humansize(nbytes), extra)

if sites and not self.few_sites:
# in SAM from 539,831 -> 11,430 sites
lowres = sites.lower_res(res=4)[0] # res=4 ~39 km
Expand All @@ -324,43 +318,99 @@ def populate_csm(self):
sf = SourceFilter(sites, oq.maximum_distance)
else:
sf = SourceFilter(None)
atomic_sources = []
normal_sources = []
reqv = 'reqv' in oq.inputs
if reqv:
logging.warning(
'Using equivalent distance approximation and '
'collapsing hypocenters and nodal planes')
multifaults = []
multifaults = [src for sg in csm.src_groups for src in sg
if src.code == b'F']
if multifaults:
with hdf5.File(multifaults[0].hdf5path, 'r') as h5:
secparams = h5['secparams'][:]
logging.warning(
'There are %d multiFaultSources (secparams=%s)',
len(multifaults), general.humansize(secparams.nbytes))
else:
secparams = ()
if oq.sequential_source_models:
# Bound preclassical memory by iterating one source model
# at a time and rebuild cmakers at end
self._run_batched(sf, secparams, reqv)
self.cmakers = get_cmakers(trt_smrs, csm.full_lt, oq)
else:
self._run_regular(trt_smrs, sf, secparams, reqv)
L = oq.imtls.size
Gfull = self.full_lt.gfull([cm.trt_smrs for cm in self.cmakers])
Gt = sum(len(cm.gsims) for cm in self.cmakers)
extra = f'<{Gfull}' if Gt < Gfull else ''
if sites is not None:
nbytes = 4 * len(self.sitecol) * L * Gt
logging.warning(f'Global RateMap of %s ({Gt=}%s)',
general.humansize(nbytes), extra)
allsources = csm.get_sources()
self.store_source_info(source_data(allsources))

def _run_batched(self, sf, secparams, reqv):
"""
Run preclassical per source-model batch when the flag of
sequential_source_models is True.
"""
oq = self.oqparam
csm = self.csm
for batch_id, sm_id, sgs_batch, trt_smrs_batch in (
csm.iter_source_model_batches()):
logging.info(
'Preclassical batch %d (source model %r): %d src_groups',
batch_id, sm_id, len(sgs_batch))
cmakers_batch = get_cmakers(trt_smrs_batch, csm.full_lt, oq)
cmaker_by_grp = {
sg.sources[0].grp_id: cm
for sg, cm in zip(sgs_batch, cmakers_batch.to_array())}
atomic_batch = []
normal_batch = []
for sg in sgs_batch:
for src in sg:
if reqv and sg.trt in oq.inputs['reqv']:
if src.source_id not in oq.reqv_ignore_sources:
collapse_nphc(src)
grp_id = sg.sources[0].grp_id
if sg.atomic:
cmaker_by_grp[grp_id].set_weight(sg, sf)
atomic_batch.extend(sg)
else:
normal_batch.extend(sg)
self._process(atomic_batch, normal_batch, sf, secparams,
cmaker_by_grp=cmaker_by_grp)

def _run_regular(self, trt_smrs, sf, secparams, reqv):
"""
Run preclassical in a single pass over all src_groups when
sequential_source_models is False.
"""
oq = self.oqparam
csm = self.csm
self.cmakers = get_cmakers(trt_smrs, csm.full_lt, oq)
atomic_sources = []
normal_sources = []
cmakers = self.cmakers.to_array()
for sg in csm.src_groups:
for src in sg:
if src.code == b'F':
multifaults.append(src)
if reqv and sg.trt in oq.inputs['reqv']:
if src.source_id not in oq.reqv_ignore_sources:
collapse_nphc(src)
grp_id = sg.sources[0].grp_id
# do nothing for atomic sources except counting the ruptures
if sg.atomic:
# compute weight sequentially
cmakers[grp_id].set_weight(sg, sf)
atomic_sources.extend(sg)
else:
normal_sources.extend(sg)
if multifaults:
with hdf5.File(multifaults[0].hdf5path, 'r') as h5:
secparams = h5['secparams'][:]
logging.warning(
'There are %d multiFaultSources (secparams=%s)',
len(multifaults), general.humansize(secparams.nbytes))
else:
secparams = ()
self._process(atomic_sources, normal_sources, sf, secparams)
allsources = csm.get_sources()
self.store_source_info(source_data(allsources))

def _process(self, atomic_sources, normal_sources, sf, secparams):
def _process(self, atomic_sources, normal_sources, sf, secparams,
cmaker_by_grp=None):
if cmaker_by_grp is None:
cmaker_by_grp = dict(enumerate(self.cmakers.to_array()))
# run preclassical in parallel for non-atomic sources
if normal_sources:
sources_by_key = groupby(
Expand All @@ -371,10 +421,9 @@ def _process(self, atomic_sources, normal_sources, sf, secparams):
# avoid a segfault in macOS
self.datastore.swmr_on()
smap = parallel.Starmap(preclassical, h5=self.datastore.hdf5)
cmakers = self.cmakers.to_array()
num_tasks = len(sources_by_key)
for grp_id, srcs in sources_by_key.items():
cmaker = cmakers[grp_id]
cmaker = cmaker_by_grp[grp_id]
cmaker.gsims = list(cmaker.gsims) # reducing data transfer
pointlike = [src for src in srcs
if hasattr(src, 'nodal_plane_distribution')]
Expand Down
54 changes: 50 additions & 4 deletions openquake/calculators/tests/logictree_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,10 +32,10 @@
from openquake.qa_tests_data.logictree import (
case_01, case_02, case_03, case_04, case_05, case_06, case_07, case_08,
case_09, case_10, case_11, case_12, case_13, case_14, case_15, case_16,
case_17, case_18, case_19, case_20, case_21, case_22, case_23, case_25,
case_26, case_28, case_29, case_30, case_31, case_32, case_33, case_36,
case_39, case_45, case_46, case_52, case_56, case_58, case_59, case_67,
case_68, case_71, case_73, case_79, case_80, case_83, case_84)
case_17, case_18, case_19, case_20, case_21, case_22, case_23, case_24,
case_25, case_26, case_28, case_29, case_30, case_31, case_32, case_33,
case_36, case_39, case_45, case_46, case_52, case_56, case_58, case_59,
case_67, case_68, case_71, case_73, case_79, case_80, case_83, case_84)

ae = numpy.testing.assert_equal
aac = numpy.testing.assert_allclose
Expand Down Expand Up @@ -501,6 +501,31 @@ def test_case_23_bis(self):
ns = len(self.calc.datastore['source_info'])
assert ns == 26

def test_case_24(self):
# Parity check: with sequential_source_models=true the hazard
# statistics (mean and quantiles) must match the regular
# (all-in-one Starmap) approach for both full enumeration and
# sampling. A small tolerance is used because task reduction
# order can differ across runs

# Full enumeration
self.run_calc(case_24.__file__, 'job.ini',
sequential_source_models='true')
seq_full = self.calc.datastore['hcurves-stats'][:]
self.run_calc(case_24.__file__, 'job.ini')
reg_full = self.calc.datastore['hcurves-stats'][:]
aac(seq_full, reg_full, atol=1e-6, rtol=1e-6)

# Sampling
self.run_calc(case_24.__file__, 'job.ini',
sequential_source_models='true',
number_of_logic_tree_samples='10')
seq_sampled = self.calc.datastore['hcurves-stats'][:]
self.run_calc(case_24.__file__, 'job.ini',
number_of_logic_tree_samples='10')
reg_sampled = self.calc.datastore['hcurves-stats'][:]
aac(seq_sampled, reg_sampled, atol=1e-6, rtol=1e-6)

def test_case_25(self):
# BCHydro-style correlated uncertainties (alt1 + alt2 + alt3)
# sampled to keep the calc fast (highly simplified version)
Expand Down Expand Up @@ -804,6 +829,27 @@ def test_case_83(self):
self.run_calc(case_83.__file__, 'job_expanded_LT.ini')
[fname_ex] = export(('hcurves/mean', 'csv'), self.calc.datastore)
self.assertEqualFiles(fname_em, fname_ex)
reg_full = self.calc.datastore['hcurves-stats'][:]

# Check that sequential approach matches regular with
# full enumeration
self.run_calc(case_83.__file__, 'job_extendModel.ini',
sequential_source_models='true')
seq_full = self.calc.datastore['hcurves-stats'][:]
aac(seq_full, reg_full, atol=1e-6, rtol=1e-6)

# Run regular approach with sampling
self.run_calc(case_83.__file__, 'job_extendModel.ini',
number_of_logic_tree_samples='10')
reg_sampled = self.calc.datastore['hcurves-stats'][:]

# Check that sequential approach matches regular with
# sampling
self.run_calc(case_83.__file__, 'job_extendModel.ini',
sequential_source_models='true',
number_of_logic_tree_samples='10')
seq_sampled = self.calc.datastore['hcurves-stats'][:]
aac(seq_sampled, reg_sampled, atol=1e-6, rtol=1e-6)

def test_case_83_eb(self):
# event based sampling with double extendModel
Expand Down
18 changes: 18 additions & 0 deletions openquake/commonlib/oqvalidation.py
Original file line number Diff line number Diff line change
Expand Up @@ -793,6 +793,14 @@
Example: *ses_seed = 123*.
Default: 42

sequential_source_models:
Flag used in classical and disaggregation calculations to dispatch
tasks one top-level sourceModel branch at a time (one Starmap per
source model, run sequentially). Not compatible with source models
that share sources across top-level branches.
Example: *sequential_source_models = true*.
Default: false

shakemap_id:
Used in ShakeMap calculations to download a ShakeMap from the USGS site
Example: *shakemap_id = usp000fjta*.
Expand Down Expand Up @@ -1261,6 +1269,7 @@ class OqParam(valid.ParamSet):
ses_per_logic_tree_path = valid.Param(
valid.compose(valid.nonzero, valid.positiveint), 1)
ses_seed = valid.Param(valid.positiveint, 42)
sequential_source_models = valid.Param(valid.boolean, False)
shakemap_id = valid.Param(valid.nice_string, None)
# example: shakemap_uri = {'kind': 'usgs_id', 'id': 'XXX'}
shakemap_uri = valid.Param(valid.dictionary, {})
Expand Down Expand Up @@ -2315,6 +2324,15 @@ def is_valid_disagg_by_src(self):
return self.ps_grid_spacing == 0
return True

def is_valid_sequential_source_models(self):
"""
sequential_source_models is only useable in classical and
disaggregation calculations
"""
if self.sequential_source_models:
return self.calculation_mode in ('classical', 'disaggregation')
return True

def is_valid_concurrent_tasks(self):
"""
At most you can use 30_000 tasks
Expand Down
Loading
Loading