From 96951b4688aa4bc5f2673ab8450708b7816fdaa0 Mon Sep 17 00:00:00 2001 From: LudwigBoess Date: Mon, 21 Sep 2026 18:29:55 +0000 Subject: [PATCH 1/3] read subdomain checkpoint metadata in one selection per variable instead of one per domain --- src/framework/domain/checkpoint/resume.cpp | 55 ++++++++++++---------- src/output/utils/readers.cpp | 22 +++++++++ src/output/utils/readers.h | 15 ++++++ 3 files changed, 66 insertions(+), 26 deletions(-) diff --git a/src/framework/domain/checkpoint/resume.cpp b/src/framework/domain/checkpoint/resume.cpp index d0818c935..13e2b6510 100644 --- a/src/framework/domain/checkpoint/resume.cpp +++ b/src/framework/domain/checkpoint/resume.cpp @@ -140,33 +140,36 @@ namespace ntt { std::numeric_limits::lowest()); } + // Each subdomain_* variable is read whole, once. Per-domain reads would + // cost one filesystem round-trip per domain per rank, since consecutive + // entries live in different writers' blocks. bool needs_reconstruction = false; - for (unsigned int dom_idx { 0 }; dom_idx < g_ndomains; ++dom_idx) { - for (auto d { 0u }; d < M::Dim; ++d) { - real_t x_min, x_max; - out::ReadVariable(io, - reader, - fmt::format("subdomain_x%d_min", d + 1), - x_min, - dom_idx); - out::ReadVariable(io, - reader, - fmt::format("subdomain_x%d_max", d + 1), - x_max, - dom_idx); - saved_extents[dom_idx].emplace_back(x_min, x_max); - global_extent[d].first = std::min(global_extent[d].first, x_min); - global_extent[d].second = std::max(global_extent[d].second, x_max); - - ncells_t nx; - out::ReadVariable(io, - reader, - fmt::format("subdomain_nx%d", d + 1), - nx, - dom_idx); - saved_ncells[dom_idx][d] = nx; - - if (nx != subdomain_ptr(dom_idx)->mesh.n_active()[d]) { + for (auto d { 0u }; d < M::Dim; ++d) { + std::vector x_min, x_max; + std::vector nx; + out::ReadVariableAll(io, + reader, + fmt::format("subdomain_x%d_min", d + 1), + x_min, + g_ndomains); + out::ReadVariableAll(io, + reader, + fmt::format("subdomain_x%d_max", d + 1), + x_max, + g_ndomains); + out::ReadVariableAll(io, + reader, + fmt::format("subdomain_nx%d", d + 1), + nx, + g_ndomains); + + for (unsigned int dom_idx { 0 }; dom_idx < g_ndomains; ++dom_idx) { + saved_extents[dom_idx].emplace_back(x_min[dom_idx], x_max[dom_idx]); + global_extent[d].first = std::min(global_extent[d].first, x_min[dom_idx]); + global_extent[d].second = std::max(global_extent[d].second, x_max[dom_idx]); + + saved_ncells[dom_idx][d] = nx[dom_idx]; + if (nx[dom_idx] != subdomain_ptr(dom_idx)->mesh.n_active()[d]) { needs_reconstruction = true; } } diff --git a/src/output/utils/readers.cpp b/src/output/utils/readers.cpp index 07e666caf..ea38b7d80 100644 --- a/src/output/utils/readers.cpp +++ b/src/output/utils/readers.cpp @@ -11,6 +11,7 @@ #include #include #include +#include namespace out { @@ -31,6 +32,22 @@ namespace out { } } + template + void ReadVariableAll(adios2::IO& io, + adios2::Engine& reader, + const std::string& quantity, + std::vector& data, + std::size_t count) { + auto var = io.InquireVariable(quantity); + if (var) { + data.resize(count); + var.SetSelection(adios2::Box({ 0 }, { count })); + reader.Get(var, data.data(), adios2::Mode::Sync); + } else { + raise::Error(fmt::format("Variable: %s not found", quantity.c_str()), HERE); + } + } + template void Read1DArray(adios2::IO& io, adios2::Engine& reader, @@ -110,6 +127,11 @@ namespace out { const std::string&, \ T&, \ std::size_t); \ + template void ReadVariableAll(adios2::IO&, \ + adios2::Engine&, \ + const std::string&, \ + std::vector&, \ + std::size_t); \ template void Read1DArray(adios2::IO&, \ adios2::Engine&, \ const std::string&, \ diff --git a/src/output/utils/readers.h b/src/output/utils/readers.h index c62489282..ea0570a56 100644 --- a/src/output/utils/readers.h +++ b/src/output/utils/readers.h @@ -4,6 +4,7 @@ * Defines generic reader functions. * @implements * - out::ReadVariable<> -> void + * - out::ReadVariableAll<> -> void * - out::Read1DArray<> -> void * - out::Read2DArray<> -> void * - out::ReadNDField<> -> void @@ -21,12 +22,26 @@ #include #include +#include namespace out { template void ReadVariable(adios2::IO&, adios2::Engine&, const std::string&, T&, std::size_t); + /** + * @brief Read the first `count` elements of a per-domain variable at once. + * @note One selection instead of `count` single-element reads: each read + * lands in a different writer's block, so per-element reads cost a + * separate filesystem round-trip apiece. + */ + template + void ReadVariableAll(adios2::IO&, + adios2::Engine&, + const std::string&, + std::vector&, + std::size_t); + template void Read1DArray(adios2::IO&, adios2::Engine&, From 7cb443606623c2339689c9fdd3d96d6d5a0b3d13 Mon Sep 17 00:00:00 2001 From: LudwigBoess Date: Mon, 21 Sep 2026 18:31:12 +0000 Subject: [PATCH 2/3] open the checkpoint on the ADIOS communicator instead of MPI_COMM_SELF so metadata is read once and broadcast --- src/framework/domain/checkpoint/resume.cpp | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/src/framework/domain/checkpoint/resume.cpp b/src/framework/domain/checkpoint/resume.cpp index 13e2b6510..c3c2b27c8 100644 --- a/src/framework/domain/checkpoint/resume.cpp +++ b/src/framework/domain/checkpoint/resume.cpp @@ -122,11 +122,8 @@ namespace ntt { adios2::IO io = ptr_adios->DeclareIO("Entity::CheckpointRead"); io.SetEngine("BPFile"); -#if !defined(MPI_ENABLED) + adios2::Engine reader = io.Open(fname, adios2::Mode::Read); -#else - adios2::Engine reader = io.Open(fname, adios2::Mode::Read, MPI_COMM_SELF); -#endif reader.BeginStep(); From 01bce5e37176edb44f235b8cef0485817168d548 Mon Sep 17 00:00:00 2001 From: LudwigBoess Date: Mon, 21 Sep 2026 18:38:45 +0000 Subject: [PATCH 3/3] add read-side ADIOS2 tuning for checkpoint resume with an optional reader thread-pool size --- input.example.toml | 16 ++++++++++++++++ src/framework/domain/checkpoint/resume.cpp | 11 +++++++++++ src/framework/parameters/parameters.cpp | 15 +++++++++++++++ src/global/defaults.h | 6 ++++++ src/output/utils/tuning.cpp | 17 +++++++++++++++++ src/output/utils/tuning.h | 19 +++++++++++++++++++ 6 files changed, 84 insertions(+) diff --git a/input.example.toml b/input.example.toml index b85c5f020..b45933696 100644 --- a/input.example.toml +++ b/input.example.toml @@ -750,6 +750,22 @@ # @default: 16777216 # @note: Scales with per-rank output volume; matches ADIOS2's default buffer_chunk_size = "" + # Size of the BP5 reader thread pool used when resuming (BP5 Threads) + # @type: int + # @default: 0 + # @note: 0 lets ADIOS2 size the pool itself. The best value depends on how + # many ranks share a node, so tune it per machine if restarts are slow + read_threads = "" + # How long to wait for the checkpoint file to appear, in seconds + # @type: int + # @default: 600 + # @note: Guards against a checkpoint directory that is slow to become + # visible on a shared filesystem + read_open_timeout_secs = "" + # Polling interval while waiting for the checkpoint file, in seconds + # @type: int + # @default: 1 + read_poll_secs = "" [diagnostics] # Number of timesteps between diagnostic logs diff --git a/src/framework/domain/checkpoint/resume.cpp b/src/framework/domain/checkpoint/resume.cpp index c3c2b27c8..822628e16 100644 --- a/src/framework/domain/checkpoint/resume.cpp +++ b/src/framework/domain/checkpoint/resume.cpp @@ -1,3 +1,4 @@ +#include "defaults.h" #include "enums.h" #include "global.h" @@ -10,6 +11,7 @@ #include "framework/parameters/parameters.h" #include "framework/specialization_registry.h" #include "output/utils/readers.h" +#include "output/utils/tuning.h" #if defined(MPI_ENABLED) #include @@ -122,6 +124,15 @@ namespace ntt { adios2::IO io = ptr_adios->DeclareIO("Entity::CheckpointRead"); io.SetEngine("BPFile"); + out::ApplyBp5ReadTuning( + io, + "BPFile", + { params.template get("adios2.read_threads", + defaults::adios2::read_threads), + params.template get("adios2.read_open_timeout_secs", + defaults::adios2::read_open_timeout_secs), + params.template get("adios2.read_poll_secs", + defaults::adios2::read_poll_secs) }); adios2::Engine reader = io.Open(fname, adios2::Mode::Read); diff --git a/src/framework/parameters/parameters.cpp b/src/framework/parameters/parameters.cpp index f3d8d507e..14c71c110 100644 --- a/src/framework/parameters/parameters.cpp +++ b/src/framework/parameters/parameters.cpp @@ -210,6 +210,21 @@ namespace ntt { "adios2", "buffer_chunk_size", defaults::adios2::buffer_chunk_size)); + set("adios2.read_threads", + toml::find_or(toml_data, + "adios2", + "read_threads", + defaults::adios2::read_threads)); + set("adios2.read_open_timeout_secs", + toml::find_or(toml_data, + "adios2", + "read_open_timeout_secs", + defaults::adios2::read_open_timeout_secs)); + set("adios2.read_poll_secs", + toml::find_or(toml_data, + "adios2", + "read_poll_secs", + defaults::adios2::read_poll_secs)); /* [diagnostics] -------------------------------------------------------- */ set("diagnostics.interval", diff --git a/src/global/defaults.h b/src/global/defaults.h index 82b4364cc..dd4a275ee 100644 --- a/src/global/defaults.h +++ b/src/global/defaults.h @@ -102,6 +102,12 @@ namespace ntt::defaults { const int aggregators_per_node = 0; const size_t max_shm_size = 4294967296ull; // 4 GiB const size_t buffer_chunk_size = 16777216ull; // 16 MiB + // Checkpoint-read knobs. read_threads == 0 leaves ADIOS2 auto-sizing its + // reader thread pool; the best value depends on the ranks-per-node layout, + // so it is left to the input file rather than fixed here. + const int read_threads = 0; + const int read_open_timeout_secs = 600; + const int read_poll_secs = 1; } // namespace adios2 namespace gca { diff --git a/src/output/utils/tuning.cpp b/src/output/utils/tuning.cpp index bad6fdcc4..2b54631dd 100644 --- a/src/output/utils/tuning.cpp +++ b/src/output/utils/tuning.cpp @@ -51,4 +51,21 @@ namespace out { io.SetParameter("OpenTimeoutSecs", "600"); } + void ApplyBp5ReadTuning(adios2::IO& io, + const std::string& engine, + const Bp5ReadTuning& bp5) { + const auto eng = fmt::toLower(engine); + if (eng != "bpfile" && eng != "bp5") { + return; + } + // Tolerate a checkpoint directory that is slow to become visible on a + // shared filesystem instead of failing the restart outright. + io.SetParameter("OpenTimeoutSecs", std::to_string(bp5.open_timeout_secs)); + io.SetParameter("BeginStepPollingFrequencySecs", + std::to_string(bp5.poll_secs)); + if (bp5.threads > 0) { + io.SetParameter("Threads", std::to_string(bp5.threads)); + } + } + } // namespace out diff --git a/src/output/utils/tuning.h b/src/output/utils/tuning.h index fbbe8c53d..b59cb1320 100644 --- a/src/output/utils/tuning.h +++ b/src/output/utils/tuning.h @@ -4,8 +4,10 @@ * for large-scale parallel filesystems. * @implements * - out::Bp5Tuning + * - out::Bp5ReadTuning * - out::TotalAggregators -> int * - out::ApplyBp5Tuning -> void + * - out::ApplyBp5ReadTuning -> void * @cpp: * - tuning.cpp * @namespaces: @@ -37,10 +39,27 @@ namespace out { // aggregators_per_node <= 0, which leaves ADIOS2 on its built-in default. auto TotalAggregators(int aggregators_per_node) -> int; + // Checkpoint-read knobs from the [adios2] toml section. The aggregation and + // buffering parameters above are write-side only and are deliberately not + // reused here: ADIOS2 accepts unknown parameters silently, so setting them + // on a reader would look like tuning while doing nothing. + struct Bp5ReadTuning { + const int threads { ntt::defaults::adios2::read_threads }; + const int open_timeout_secs { ntt::defaults::adios2::read_open_timeout_secs }; + const int poll_secs { ntt::defaults::adios2::read_poll_secs }; + }; + // Apply the [adios2] BP5 tuning to a freshly declared IO whose engine is // BPFile/BP5. A no-op for other engines. void ApplyBp5Tuning(adios2::IO&, const std::string& engine, const Bp5Tuning&); + // Apply the checkpoint-read knobs to a freshly declared reader IO whose + // engine is BPFile/BP5. A no-op for other engines. `threads <= 0` leaves + // ADIOS2 to size its own reader thread pool. + void ApplyBp5ReadTuning(adios2::IO&, + const std::string& engine, + const Bp5ReadTuning&); + } // namespace out #endif // OUTPUT_UTILS_TUNING_H