XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration Shunyu Yao
Dimitrios S. Nikolopoulos
Ali R. Butt
[email protected] Virginia Tech Blacksburg, VA, USA
[email protected] Virginia Tech Blacksburg, VA, USA
[email protected] Virginia Tech Blacksburg, VA, USA
arXiv:2609.16531v1 [cs.DC] 15 Sep 2026
Abstract Message Passing Interface (MPI) has been the dominant programming model for High Performance Computing (HPC) for three decades, and as HPC workloads increasingly migrate to cloud infrastructure for scalability and cost efficiency, MPI applications must contend with an execution environment fundamentally unlike traditional supercomputers: ephemeral resources, dynamic pricing and preemptable instances. In such a volatile setting, the ability to relocate running MPI processes between nodes without restarting the job is a necessity for cost-effective, resilient execution. Existing approaches either require restarting the entire job from a global checkpoint, or transparently intercepting the full MPI stack at prohibitive complexity. To address these challenges, we propose XMPIaaS, a cooperative migration system for MPI that enables selective process group migration on-the-fly. When a cloud instance is scheduled for preemption, only the affected ranks are relocated while the remaining processes briefly quiesce and resume in place, avoiding the cost of a full-job checkpoint. XMPIaaS tackles this through a cooperative protocol between the MPI process management runtime and rank processes. We expose an XMPI_quiesce interface built atop the MPI Sessions API that allows applications to mark safe migration points, and we extend the Hydra process manager to orchestrate the full migration lifecycle: rank quiescence, CRIU checkpoint/restore, proxy relaunch on the target node, and seamless rank reconnection. We evaluate and show that the cooperative quiesce phase accounts for less than 1.4% of total migration downtime, and that this downtime is governed by the migrating node’s rank count alone, independent of job size, and the instrumentation introduces no measurable overhead during normal execution.
CCS Concepts • Computer systems organization → Cloud computing; Reliability; • Computing methodologies → Parallel programming languages.
1
Introduction
The Message Passing Interface (MPI) has served as the de-facto standard [7, 8, 21, 22, 43] for parallel communication in High Performance Computing (HPC), forming the backbone of large-scale scientific simulations. Traditionally, MPI applications are deployed on dedicated supercomputers and HPC clusters under a rigid model with static resource allocation and exclusive node [58]. As cloud platforms evolve to offer increasingly powerful compute, HPC communities are increasingly migrating their workloads to cloud platforms being attracted by on-demand scalability, resource elasticity and flexible pricing [45, 54].
However, the conventional deployment model of MPI applications is fundamentally at odds with the cloud computing paradigm, where resources are ephemeral and preemption can happen at any time. Several cloud deployment models exist for HPC workloads, each presenting distinct challenges for MPI. The simplest way of renting dedicated virtual machines or bare metals in an Infrastructure-as-a-Service (IaaS) model offers little advantage over traditional clusters beyond hardware procurement convenience, resources remain statically allocated and exclusively held [45]. Preemptible instances (AWS Spot, Azure Spot) reduce costs for up to 90% [1, 41], but the provider may reclaim any node at any time, and MPI’s tightly coupled runtime and processes cannot survive partial node loss. Serverless and Function-as-a-Service platforms [33, 52] offer the most cloud-native execution model with automatic scaling, per-invocation billing, and zero infrastructure management, but it imposes strict resource limits per instance and assume stateless, short-lived workloads, fundamentally incompatible with MPI’s long-running, stateful processes [26]. Across all these models, the common barrier is the same: MPI lacks a mechanism for individual processes to survive the preemption and resource boundaries in a volatile cloud environment. The natural solution is to evacuate affected processes from terminating instances and migrate them onto available ones for continued execution. Enabling migration for MPI ranks — the individual MPI worker processes — is difficult for three reasons. First, MPI ranks maintain complex network states from TCP socket buffers to RDMA queue pairs, memory registrations, and completion queues. These states are unfit for checkpointing because they are opaque to userspace [16, 48]. General-purpose checkpoint tools cannot serialize these states, and relocating a rank to a new node renders them stale. Thus, network states must be torn down at checkpoint and reconstructed at restore. Second, MPI ranks are tightly coupled. Different MPI applications have their distinct rank communication patterns, and MPI ranks synchronize through collective operations. Therefore, migrating ranks breaks the collectively established communication topology, disrupts applications’ embedded synchronization pattern, and requires coordinated action across the entire job to participate in rebuilding new connections. Third, existing MPI process managements are designed around a static process-to-node mapping established at launch time [12, 44], they provide no mechanism for a rank to change its hosting node, update its network identity, or rejoin the MPI process group from a new location mid-execution. Prior works have sought to address resilience for MPI applications[3, 6, 11, 13, 18, 24, 37, 38, 42, 46, 50, 55–57] (See Section 7). Representatively, MANA [18] achieves transparent checkpoint/restore by virtualizing internal MPI states, using a split-process architecture
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt
to isolate the MPI library and discard network context at checkpoint time. The transparency costs tracking and replaying every MPI object’s lifecycle, and requires the entire job to be checkpointed no matter the granularity of disruption. Charm++ and AMPI [9, 10, 13, 23, 29, 34, 35] decompose a job into migratable virtual ranks that the runtime’s load balancer can relocate, evacuating a node without changing the rank count the application sees. However, AMPI requires privatizing global and static variables and an over-decomposed launch, and it relocates its own virtual ranks rather than a stock MPI process. User-Level Fault Mitigation (ULFM) [11] extends the MPI standard with fault-detection and communicator-repair semantics, allowing applications to survive process failures. But ULFM provides no mechanism to preserve or relocate the ranks’ progress, and it requires applications to explicitly implement its own recovery logic from scratch. The ability to evacuate only affected instances and minimize disruption to the rest of the job for modern MPI is left unaddressed. We propose XMPIaaS, a cooperative framework for MPI that enables the selective migration of the MPI process group at the execution time. At the time of preemption, XMPIaaS allows provider to issue notifications to MPI process management proxies and MPI ranks on only affected instances, thus avoiding the cost of a fulljob checkpoint. Then XMPIaaS orchestrates a migration protocol between MPI application processes and the MPI process manager. On the application side, XMPIaaS exposes an XMPI_quiesce interface built atop the MPI Sessions API [43], allowing developers to mark safe migration points at natural iteration boundaries. When triggered, all ranks collectively tear down their network resources and park at the quiesce point. After the affected ranks have been relocated, all ranks synchronize and jointly open a fresh session before resuming computation. Ranks that are not migrating never undergo checkpointing, they simply wait at the quiesce point and rejoin once the relocated ranks have been restored. On the process manager side, we extend MPI’s process management layer to coordinate the migration: spinning up a replacement proxy on the destination node, snapshotting only the affected ranks using CRIU [16], performing an atomic control-channel handover from the departing proxy to its replacement, and restoring each rank with substituted process management control sockets. We implement XMPIaaS on top of MPICH’s Hydra process manager [21, 44], though its design relies only on abstractions common to all major MPI process managers and is portable to other implementations such as Open MPI’s PRRTE [12]. Three key insights underpin XMPIaaS’s design. First, at a cooperative quiesce point, the process is reduced to its simplest possible form: application memory and a single process management socket with no network connections. This makes all the rank states local and trivially checkpointable by CRIU without any MPI-specific plugins or virtualization. Second, migration is per-instance, not per-job: only the ranks on the affected node are checkpointed and relocated, while non-migrating ranks experience a brief quiesce and resume in place, avoiding the cost of a full-fleet checkpoint. Third, we build XMPIaaS on top of MPI libraries, with design that is easily portable to other standard-conforming MPI library and process manager implementations. Together, these properties make XMPIaaS a natural fit for the cloud, where preemption strikes individual instances with little warning.
The main contributions of this paper are summarized as follows: • We design an application-side cooperative quiesce built on the MPI Sessions API that allows all ranks to tear down and reconstruct their MPI state within a single process lifetime. • We design and implement a runtime-side migration orchestration protocol within MPICH’s Hydra process manager that coordinates rank quiescence, selective checkpointing, proxy replacement, and seamless process group reconstruction. • We evaluate XMPIaaS on a multi-node cluster with four proxy applications and show that the cooperative quiesce phase accounts for less than 1.4% of total migration downtime, image transfer dominates at 60–92%, and migration downtime is governed by the migrating node’s rank count alone, independent of job size, the instrumentation introduces no measurable overhead during normal execution.
2 Background 2.1 Cloud Deployment Models for HPC Cloud platforms offer several deployment models for computeintensive workloads, each presenting different tradeoffs for MPI applications. The most straightforward is dedicated Infrastructureas-a-Service (IaaS), where users rent virtual machines or bare-metal instances for the duration of a job [45]. This model mirrors traditional cluster computing and is fully compatible with MPI, but resources remain statically allocated and exclusively held. Preemptible instances (Spot Instances on AWS, Spot VMs on Azure, and Spot VMs on Google Cloud) provide access to spare cloud capacity at substantially reduced cost [1, 20, 41]. Amazon EC2 Spot Instances offer up to 90% discount compared to on-demand prices [1], and Azure Spot VMs similarly offer discounts of up to 90% off standard pay-as-you-go prices [41]. The tradeoff is that when Azure needs capacity back, the infrastructure will evict Spot VMs with 30 seconds notice [41], while AWS provides a two-minute warning before reclaiming a Spot Instance [1, 4]. This preemption model is designed for stateless, fault-tolerant workloads, it is not practical to run parallel MPI jobs or stateful HPC services on these lowcost instances under current tooling [25, 54]. MPI’s tightly coupled processes cannot survive partial node loss. When one instance is reclaimed, all ranks holding pairwise connections to the evicted ranks are left with broken links, and the entire job typically must be restarted. Enabling MPI to tolerate preemption would unlock a significant portion of cloud capacity for scientific computing. Serverless and Function-as-a-Service (FaaS) platforms represent the most cloud-native execution model, offering automatic scaling, per-invocation billing, and zero infrastructure management [2, 19, 33, 40]. Today, serverless computing poses challenges for HPC workloads due to resource limits imposed by cloud providers, including maximum memory, CPU, and runtime restrictions [15, 52]. More fundamentally, FaaS platforms contracts stateless, short-lived function invocations, while MPI requires long-running processes with persistent state and direct inter-process communication [26]. A long-running MPI simulation cannot complete within a single function invocation’s lifetime. However, a migration-capable MPI runtime could bridge this gap and survive the resource boundaries of serverless instances. If ranks can be checkpointed before a
XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration
function’s time limit expires and restored into a fresh invocation, the serverless platform effectively becomes a source of renewable short-lived compute slots stitched together into a long-running job [39]. This reframing transforms serverless from an incompatible execution model into a viable, elastically scaled substrate for MPI, provided that migration overhead remains small relative to the function’s available runtime.
2.2
MPI Fault Tolerance
Prior work on MPI resilience imposes its cost in one of three places. Transparent checkpoint/restart places the cost on the MPI implementation, preserving a rank across a checkpoint requires interposing on the MPI interface, virtualizing and replaying internal MPI state, and supplying transport-specific machinery to discard and rebuild network contexts that cannot be serialized [3, 18, 29, 57]. This requires either MPI implementation dependence or explicit application rewrite against a dedicated customized execution model. Malleability and communicator-repair approaches instead place the cost on the application requiring user to modify compute logic to tolerate rank shrinking and expansion, redistributing data into a new decomposition, or supplying explicit recovery logic [11, 14, 30– 32, 38]. Whole-job checkpoint/restart has a cascading cost due to that a preemption confined to limited instance can trigger a checkpoint, a rollback, or a reconfiguration across every rank, at an expense that scales with the size of the job rather than the size of the loss [6, 18, 42, 50]. There still lacks an approach with balanced tradeoff among the above cost to accommodate MPI in the volatile environment on cloud.
2.3
MPI Sessions
the rank processes locally on its node. The internal mechanisms can be explained at Section 3. Ranks communicate with this process management hierarchy through the Process Management Interface (PMI) [5, 12], a wire protocol between MPI rank processes and their process manager. Through PMI, ranks perform key-value store operations, fence barriers, and process lifecycle commands such as spawn. Its primary role during initialization is coordinating the exchange of endpoint addresses needed for communicator construction. The per-node proxy acts as a relay, forwarding PMI requests upstream to the central launcher, which serves as the authoritative key-value store and barrier coordinator. PMI has evolved through three generations: PMI-1, PMI-2, and PMIx [5, 12]. They differ in framing, scalability optimizations, and feature set, but all three share the same core semantics.
2.5
CRIU
CRIU [16] is a Linux tool that can freeze a running process, serialize its states of memory pages, file descriptors, register contents, signal masks, and memory mappings to a set of image files on disk, and later restore the process from those images, potentially on a different node. CRIU operates mainly in userspace, leveraging kernel interfaces such as ptrace, /proc, and prctl to capture and reconstruct process state. CRIU can checkpoint standard POSIX resources such as regular files, pipes, Unix sockets, TCP connections. CRIU also provides a file-descriptor inheritance mechanism (criu_add_inherit_fd) that allows the restoring process to receive fresh file descriptors in place of the originals. However, CRIU cannot serialize state that lives in kernel subsystems or hardware devices without explicit plugin support. Notably, CRIU does not support InfiniBand [48]. Device file descriptors, kernel-side RDMA objects (queue pairs, completion queues, memory registrations), and device-backed memory mappings are all opaque to CRIU’s checkpoint machinery.
The MPI 4.0 standard [43] introduced the Sessions model as an alternative to the traditional MPI_Init/MPI_Finalize lifecycle. Under the legacy model, an application initializes MPI exactly once and finalizes it exactly once, there is no standard mechanism to tear down MPI state and reinitialize it within the same process for multiple 3 Process Management in Hydra times. With the Sessions API, MPI resource management is decoupled from process lifetime, and an application may repeat arbitrary This section describes the internal architecture of MPICH’s Hydra rounds of Session creation (MPI_Session_init), obtaining an MPI process manager that XMPIaaS repurposes for migration. We introprocess group from a named process set (MPI_Group_from_session_pset), duce the three-tier topology of process management launcher, proconstructing a communicator (MPI_Comm_create_from_group), cess management proxy, and application rank by walking through and finalizing the session (MPI_Session_finalize). Each session the process management lifecycle of startup (3.1 and 3.2) and tearcycle yields a fresh set of MPI resources with no residual state down (3.3). from the previous session, and is effectively equivalent to a complete MPI_Init/MPI_Finalize lifecycle. This ability to repeatedly tear 3.1 Architecture down and reconstruct MPI state within a single process lifetime Hydra’s central design choice is a strict separation between process opens door for XMPIaaS’s migration scheme that must preserve management and computation: all coordination state resides in a rank’s application state across relocation without restarting the a single launcher and the per-node proxies, while ranks hold no process. management state of their own. This separation arises from a three-
2.4
MPI Process Management
MPI implementations rely on a process management layer to launch, monitor, and coordinate rank processes throughout a job’s lifetime. In MPICH, this role is filled by Hydra [44], which follows a threetier architecture. A central launcher process (mpiexec) spawns one proxy (pmip) per node on each allocated node via SSH or a resource manager such as SLURM. Each proxy in turn forks and executes
tier hierarchy shown in Figure 1. launcher (mpiexec) sits at the top as the job launcher and central coordinator. For each node in the allocation, it SSH-launches a proxy (pmip), and each proxy in turn fork-execs the MPI ranks assigned to that node. The mpiexec–pmip control channel is always a TCP socket, regardless of the transport configured for MPI data traffic between ranks. Over this channel, the launcher and proxies communicate through Hydra commands. One example of such commands can be
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt
mpiexec(hydra) Hydra Control
PMI control
PMI control
rank process
rank process
rank process
Hydra Control
pmi proxy
pmi proxy rank process
KVS
Inter-Node Rank Connection
Intra-Node Rank Connection
Compute Node
rank process
rank process
rank process
rank process
Intra-Node Rank Connection
Compute Node
Figure 1: Hydra’s three-tier process management topology. The TCP control channel connects mpiexec to each pmip, while a Unix socketpair carries PMI messages between each pmip and its local ranks.
launch commands carrying the executable path, environment, and rank assignments. Hydra commands define types covering events such as process launch, PMI request forwarding, I/O relay, and exit status reporting. The pmip-rank control channel is a Unix socketpair created by pmip during startup. Through PMI calls, pmip and proxies communicate over this PMI socket for rank setup and health management such as rank stdout and stderr redirection or business card exchanges (See 3.2).
3.2
Address Exchange
After the above architecture spawns, as part of communicator creation within MPI_Init or MPI_Session_init calls, ranks discover one another through a key-value store coordinated by mpiexec to exchange business card, which is a string encoding a rank’s transport endpoint, including its hostname, port, and connection tag. Each rank publishes its business card by issuing a PMI_put on its socketpair. In Hydra’s implementation, the proxy buffers these puts locally and flushes them to mpiexec as a single batch when the rank enters a PMI barrier. On the mpiexec side, it serves as the authoritative KVS store, ingesting each batch and tracks an epoch counter per contributing process. When every rank in the job has reached the same epoch, mpiexec signals all proxies, which unblock their ranks. Each rank then retrieves its peers’ business cards via PMI_get and opens peer connections for MPI data traffic. From this point on, ranks communicate directly with their peers over these transport connections, bypassing pmip and mpiexec entirely. This put–barrier–get cycle is the only point at which all ranks must globally converge through the process management layer.
3.3
Teardown
Teardown reverses the startup sequence. When a rank calls MPI_Finalize or MPI_Session_finalize, the MPI library closes all transport connections and releases their associated resources. After finalize returns, the rank retains no MPI network state; its only remaining tie to the runtime is the PMI socketpair. When the rank subsequently exits,
pmip detects the termination and end-of-file on the socketpair and stdout/stderr pipes. Once all local ranks have terminated, pmip reports their exit statuses upstream via the command protocol and itself exits.
3.4
Observation
MPI_Session_init / MPI_Session_finalize allows multiple rounds of communicator states creation and teardown (See 2.3), this enables checkpointing opportunities for ranks between two MPI Sessions, in which they are ordinary user-space processes whose sole kernelvisible connection to the MPI runtime is a single Unix file descriptor. This is the window in which a process can be checkpointed without capturing any MPI transport state.
4 Design 4.1 Workflow Overview XMPIaaS enables selective migration of individual MPI ranks midexecution without checkpointing the entire job. When a cloud provider issues a preemption notice for a node, XMPIaaS migrates only the affected ranks to a replacement instance while the rest of the job pauses. The workflow has three stages: first, all ranks cooperatively quiesce, tearing down their MPI sessions and parking at a known-safe state (4.2); second, the process manager checkpoints the affected ranks, relocates them to a replacement proxy on the new node, and redirects the control plane; third, all ranks re-initialize their sessions and exchange updated endpoint addresses to establish fresh connections. Normal computation then resumes from where it left off (4.3). XMPIaaS matches the granularity of the cloud preemption model. Per-instance migration prevents a preemption on a single instance from cascading into checkpointing and restarting every rank across every node. Only affected instances are checkpointed/restored, the remaining ranks experience only a brief quiesce stall rather than a full checkpoint and restart cycle, minimizing the downtime disruption.
4.2
Cooperative Quiesce
To react to a preemption notice, XMPIaaS exposes a CLI control utility, hydra_ctl, for cloud admin to trigger a migration notification. At the time of being notified, all ranks must enter a quiesce state where they pause any computations and communications with no incomplete communications. This is due to MPI connections are inherently pairwise: if a migrating rank tears down its transport endpoint, the non-migrating peer holding the other end of that connection is left with a broken link. Rather than attempting to handle these half-open connections ( which would require error recovery logic in the MPI library), we require all ranks to tear down their sessions cooperatively, so that every connection is closed cleanly from both ends. We design an XMPI_quiesce (Algorithm 1) interface for applications to insert calls at periodic safe points, marking where migration and its resulting quiesce may safely occur. Such safe points typically coincide with iteration boundaries in a simulation’s main loop (Algorithm 2). The application developer, who understands the program’s communication structure, identifies the safe quiesce points. When no migration notice is sent, XMPI_quiesce
XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration
Algorithm 1 The XMPI_quiesce procedure. During normal execution, the flag check at line 1 causes an immediate return. When migration is requested, the rank tears down its MPI session (lines 4–7), signals the proxy and blocks until relocation completes (lines 8–9), then reconstructs a fresh session (lines 11–13).
4.3
Migration Protocol
We implement a new protocol on top of Hydra to conduct the selective migration procedure, shown in Figure 2. pmip(M) and rank(M) denote the proxy and ranks on the node being evacuated; pmip(N) and rank(N) denote proxies and ranks on nodes not involved in the 1: if migration_requested = 0 then migration; pmip(M’) denotes the replacement proxy spawned on 2: return {no-op during normal execution} the target node; and rank(M’) denotes the restored ranks after they 3: end if have been relocated to the target node under pmip(M’). 4: MPI_Barrier(comm) Migration begins when mpiexec receives a hydra_ctl migra5: MPI_Comm_free(comm) tion command that contains the migration target and destinations 6: MPI_Group_free(group) 1 . It first spawns a replacement proxy, pmip(M’), on the target 7: MPI_Session_finalize(session) node. once pmip(M’) connects 2 , mpiexec initiates the quiesce 8: write(pmi_fd, "xmpi_quiesce\n") sequence by sending out quiesce Hydra command 3 . At receiv9: read(pmi_fd) {block until proxy responds} ing the command, migrating proxy pmip(M) and non-migrating 10: migration_requested ← 0 proxy pmip(N) mark their ranks as migrating and non-migrating 11: MPI_Session_init(session) separately 4 . Each proxy relays the notification to its ranks as a 12: MPI_Group_from_session_pset(session, group) POSIX SIGUSR2 5 , notifying both rank(M) and rank(N) that a node 13: MPI_Comm_create_from_group(group, comm) preemption notification has received. The rank’s handler sets a "migration_requested" flag and returns immediately to continue executing application code with no checkpoint, no exit at this point. At the next call site to XMPI_quiesce 6 (Section 4.2), Both Algorithm 2 LULESH integration with XMPIaaS. A single rank(M) and rank(N) first execute an MPI barrier to drain all inXMPI_quiesce call at the iteration boundary (line 6) is the only flight messages, ensuring no peer’s pending send is discarded by application modification required. an early teardown. They then tear down their network states gracefully through MPI’s own session API 7 . Finally, all ranks write 1: MPI_init() "xmpi_quiesce" on their PMI socket to notify the proxies that they 2: domain ← initialize simulation state have reached a checkpoint-safe state 8 . After this point, rank(M) 3: while time < stoptime do and rank(N) behave differently. rank(M) block on a read from 4: TimeIncrement(domain) the PMI socket, waiting for a "xmpi_continue" response of their 5: LagrangeLeapFrog(domain) {physics + MPI comms} pmip(M’) after relocation completes. rank(M) remain parked this 6: XMPI_quiesce {safe point: no in-flight messages} way while pmip(M) proceeds to checkpoint them. For rank(N), the 7: end while proxy responds immediately with "xmpi_continue" after rank(N) 8: MPI_finalize() sends "xmpi_quiesce", unblocking rank(N) to begin session re-initialization, they will wait at the PMI barrier embedded in the re-initialization calls until migrated ranks rank(M’) arrive 9 12 . After receiving "xmpi_quiesce" from rank(M), pmip(M) invokes functions as a no-op. All ranks will enter quiesce procedure when CRIU to checkpoint the processes 10 . Because each rank has already calling XMPI_quiesce if there is notification received. Within the finalized its MPI session before reaching this phase, the rank holds XMPI_quiesce, ranks will call an MPI barrier to drain any potenno MPI transport state at this point, only application-level memtial in-flight messages or partially completed collectives. Migrating ory and the PMI socket. CRIU writes rank images to disk together ranks utilizes MPI sessions API to gracefully cleanup old network with the fd number of the rank’s PMI socket (for later recovery context and start new sessions after migration, so that there is purposes). CRIU then terminates the ranks and images are transalways no network states left at the time of checkpoint. ferred by pmip(M) directly to the target node hosting pmip(M’) 11 . MPI Sessions are the enabling mechanism for XMPIaaS’s quiIf multiple nodes are migrating, their pmip(M) initiate the transfers esce. Unlike MPI_Finalize that permanently terminates a process’s asynchronously and independently. MPI participation, multiple session cycles can exist with a rank’s When pmip(M) has checkpointed all its local ranks, it sends a lifetime (See 2.3). This makes the quiesce-migrate-rejoin cycle or Hydra command upstream to mpiexec, distinguishing an inteneven multiple migration rounds possible. Because teardown and tional migration termination from crashes or exits 13 . mpiexec reconstruction go through the MPI implementation’s own session then replaces its proxy table slot for pmip(M) to pmip(M’). From path and XMPIaaS never touches transport internals, XMPIaaS is this point forward, every Hydra command that mpiexec sends to genuinely interconnect-agnostic. this proxy slot goes to pmip(M’) on the target node. We chose to require application participation by calling XMPI_quiesce The replacement proxy pmip(M’), already running on the tarrather than interpose on the entire MPI API surface like MANA, beget node, receives a Hydra command from mpiexec and restores cause cooperative quiesce together with MPI’s own session utilities each rank from its CRIU image rather than fork-execing fresh promakes clean communication state teardown and rebuild, eliminatcesses 14 . Since the restored rank’s original PMI socket pointed ing the need for MPI state virtualization, replay logic, and transportspecific plugins, at the cost of a single application-inserted call.
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt Migration
❶ Notice
mpiexec
2->3
pmip
pmip
mpiexec
❷
Spawn
Connect
pmip(M’)
Quiesce
rank
rank
rank
rank
rank
rank
rank
Node 1
Node 2
Node 3
❺
XMPI_Queisce()
XMPI_Queisce()
❻
rank(N) rank(N)
zzz
rank(N)
zzz
rank(N)
Node 1
zzz
zzz
rank(M) rank(M)
❼
zzz
pmip(M’)
❺
rank(N)
❹rank(M)
rank(M)
rank(N)
rank(N)
rank(M)
rank(M)
Node 1
pmip(M’)
❽
quiesced
zzz
SIGUSR2
Node 2
Node 3
mpiexec
pmip(M)
❽
quiesced
❸
❹rank(N)
mpiexec
pmip(N)
Quiesce
pmip(M)
pmip(N) SIGUSR2
rank
❸
rank(M) rank(M)
XMPI_Queisce()
zzz
XMPI_Queisce()
Node 2 ❼
pmip(M)
pmip(N) continue
zzz
Node 3
❾
pmip(M’)
❿
ckpt
⓫
rank(N)
rank(N)
rank(M)
rank(M)
rank(M’) zz
rank(M’) zz
rank(N)
rank(N)
rank(M)
rank(M)
rank(M’) zz
rank(M’) zz
Node 1
Node 2
z
z
z
z
Node 3
mpiexec
mpiexec departing ⓭
pmip(M)
pmip(N)
rstr
continue
⓮
rank(N)
rank(N)
rank(M’)
rank(M’)
rank(N)
rank(N)
rank(N)
rank(N)
rank(M’)
rank(M’)
rank(N)
rank(N)
Waiting …
Waiting … ⓬
Node 1
pmip(M’)
pmip(N)
pmip(M’)
⓯
rank(M’)
rank(M’)
rank(M’)
rank(M’)
⓰
Node 2
Node 3
Node 1
Node 2
Node 3
Figure 2: Migration protocol timeline. Table 1: Process management abstractions across MPI implementations.
MPICH Open MPI Intel MPI Slurm Cray PALS
Proxy pmip prted pmip slurmstepd palsd
PMI PMI-1/2 PMIx PMI-1/2 PMI-1/2/x PMIx
Topology Flat Tree (64) Flat Tree Tree (32)
to pmip(M), which no longer exists, pmip(M’) creates fresh replacements and instructs CRIU to map them onto the fd numbers recorded during checkpoint. The restored rank wakes holding the same fd numbers it had before, but the other end of each now belongs to pmip(M’) on the new node. pmip(M’) writes "xmpi_continue" on the new PMI socket, unblocking the rank, which clears its migration flag and proceeds to MPI_Session_init 15 . All ranks then re-execute the KVS cycle from Section 3.2. All ranks republish business cards with rank(M’) soliciting their new addresses on the target node. When the barrier releases, all ranks retrieve updated addresses and open fresh MPI connections 16 . Execution resumes from the point where XMPI_quiesce was called.
4.4
Design Portability
XMPIaaS’s migration protocol depends on three abstractions seen in Hydra: a central coordinator that holds process management
state, per-node proxies that relay traffic between the coordinator and ranks, and a PMI channel connecting each rank to its local proxy. These abstractions are not Hydra-specific, they are commonly present in major MPI process managers (Table 1). Open MPI’s PRRTE uses prted as its per-node proxy with PMIx over Unix sockets, Slurm’s slurmstepd serves the same role with configurable PMI-1, PMI-2, or PMIx support, and Intel MPI uses an unmodified fork of Hydra. They share similar structures, and thus XMPIaaS’s design is highly portable to other process management implementation. The rank-side code uses only standard MPI Sessions API calls and requires no modification inside implementations, thus XMPIaaS can work on top of any standard-conformant MPI. The one architectural difference worth addressing is proxy topology. Hydra connects all proxies directly to mpiexec in a flat tree, while PRRTE organizes proxies into a radix tree with a default fanout of 64, and Cray PALS uses a tree with fanout 32. Under a hierarchical topology, the coordinator-to-proxy command path would traverse intermediate proxies, requiring them to forward migration commands rather than receive them directly. However, this distinction is largely academic for real workloads: facility data from NERSC and NREL [47, 49] shows that 95–98% of HPC jobs by count run on fewer than 64 nodes, with medians of 1–4 nodes. At these scales, PRRTE’s tree collapses to a single level, which is functionally identical to Hydra’s flat topology. Extending XMPI to multi-level trees is straightforward but orthogonal to our core contribution, and we defer it to future work (Section 6).
XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration
Evaluation
This section seeks to answer the following questions: • Q1: What overhead does XMPIaaS introduce during normal execution when no migration occurs? • Q2: How long are migrating ranks unavailable during a migration event, and what dominates the downtime? • Q3: How does the stall imposed on non-migrating ranks scale with per-node rank count, job size, and evacuation width? • Q4: Does XMPIaaS remain stable across repeated migrations, with no state leaks or degradation?
5.1
Experimental Setup
Hardware. We use a 9-node cluster of identical machines, each equipped with two Intel Xeon Gold 6242 CPUs (16 cores per socket, 2.80 GHz base, 3.90 GHz turbo), 192 GB DDR4 RAM, and 10 Gbps Ethernet. Nodes run Ubuntu 22.04 (kernel 5.15.0). Software. MPICH 4.2.2 with the CH4:OFI communication device (FI_PROVIDER=tcp). CRIU version 4.2. All benchmarks are compiled with -O3. Image transfer is handled by scp; reducing transfer latency is orthogonal to XMPIaaS’s design and left to the underlying transport. Topology. All experiments run on a nine-node cluster: one node is dedicated to the launcher and runs mpiexec only, hosting no compute ranks, while the remaining eight nodes serve as compute nodes and migration targets. We denote by 𝐶 the number of compute nodes, each running a single Hydra proxy (pmip) that hosts a fixed number of ranks per node (rpn), and by 𝑀 the number of proxies evacuated in a migration event; the two are bounded by 𝐶 + 𝑀 ≤ 8 and 𝑀 ≤ 𝐶. Because Hydra runs one proxy per node, each evacuated proxy is relocated to a distinct idle target, so 𝑀 concurrent evacuations consume 𝑀 free target nodes. Source and target nodes are kernel-matched so that CRIU can restore each checkpoint image on its target. During a migration, each source proxy transfers its rank images directly to its paired target over scp, independently of mpiexec and of every other migrating proxy; the 𝑀 transfers therefore proceed concurrently rather than being relayed serially through the launcher host. Benchmarks. We evaluate four proxy applications from the Mantevo [28] and ECP suites, selected to span distinct computational and communication patterns (Table 2): • LULESH 2.0 [36]: a shock hydrodynamics code with structured halo exchange between neighboring subdomains. • CoMD 1.1 [17]: a classical molecular dynamics code with irregular, data-dependent neighbor communication as atoms migrate between spatial domains. • HPCCG 1.0 [27]: a conjugate gradient solver dominated by sparse matrix-vector products and global allreduce operations. • miniAMR 1.0 [51]: an adaptive mesh refinement code with dynamic neighbor communication and communicator splitting as blocks are refined and coarsened.
Table 2: Benchmark applications and their communication characteristics. Benchmark LULESH 2.0 CoMD 1.1 HPCCG 1.0 miniAMR 1.0
Domain Shock hydro. Molecular dyn. CG solver Adaptive mesh
Comm. Pattern Structured halo exchange Irregular neighbor exch. SpMV + allreduce Dynamic comm + split
Baseline
Normalized Runtime
5
1.25 1.00
-3.8%
+9.8%
XMPIaaS +1.0%
+14.2%
0.75 0.50 0.25 0.00
LULESH CoMD
HPCCG miniAMR
Figure 3: Normalized runtime of XMPIaaS versus baseline (no migration triggered). Differences range from −3.8% to +14.2% across benchmarks, consistent with instruction cache alignment noise rather than systematic overhead. transferring tools since XMPIaaS do not focus on optmizing image transferring latency.
5.2
Normal-Run Overhead (Q1)
To measure the overhead XMPIaaS introduces during normal execution, we compare the XMPIaaS-instrumented binary (MPI Sessions, signal handler, and XMPI_quiesce called once per timestep) against an unmodified baseline using MPI_Init/MPI_COMM_WORLD. Both binaries are linked against the same MPICH build. We run each benchmark on a single node with 8 ranks. Figure 3 shows the normalized runtime for each benchmark. The direction of the difference flips between benchmarks: LULESH runs 3.8% faster under XMPIaaS, CoMD is 9.8% slower, HPCCG is within 1.0%, and miniAMR is 14.2% slower. This pattern where the sign of the difference varies across binaries despite identical hardware is characteristic of instruction cache alignment noise rather than real overhead. The per-iteration cost of XMPI_quiesce when no migration is pending consists of a single volatile read of the migration_requested flag followed by an untaken branch, which is unmeasurable at the granularity of our benchmarks. We conclude that XMPIaaS introduces no meaningful overhead during normal execution.
5.3
Migration Downtime Breakdown (Q2)
To understand what dominates migration downtime, we instruAll four were ported to XMPIaaS by replacing MPI_Init/MPI_COMM_WORLD ment the migration path with millisecond-resolution timestamps with MPI Sessions and inserting a single XMPI_quiesce call at each (clock_gettime) and decompose the total downtime of migratapplication’s main iteration boundary. We use scp as the image ing ranks into four phases: quiesce (barrier + session teardown),
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt
Downtime (ms)
20000
Quiesce CRIU Dump
Image Transfer Restore + Reconnect
15000
15.5s
10000 5000 0
6.1s
4.7s
4.1s
LULESH
CoMD
Set A B C D
Nodes (𝐶 ) 2–7 4 4 4
Ranks/Node 16 2–32 16 4–32
𝑀 1 1 1–4 1–4
Total Ranks 32–112 8–128 64 16–128
Table 4: LULESH scaling configurations. The perfect-cube rank constraint replaces the fixed-rpn node sweep with strong- and weak-scaling variants; swept axis in bold.
HPCCG miniAMR
Figure 4: Migration downtime breakdown by phase. Image transfer via scp dominates in all cases (60–92%). Quiesce is negligible (<55 ms). CRIU dump time scales with per-rank memory footprint.
CRIU dump, image transfer (scp), and restore + reconnect (CRIU restore, xmpi_continue, session re-initialization, and KVS address exchange). We migrate 4 ranks from one node to another across all four benchmarks. Migration downtime is the interval from the first CRIU dump to the last restore-and-reconnect across the migrating ranks. Figure 4 shows the stacked breakdown. Three findings stand out. First, image transfer dominates total downtime across all benchmarks, accounting for 91% in LULESH, 92% in CoMD, 81% in HPCCG, and 60% in miniAMR. This is an artifact of our testbed’s 10 Gbps Ethernet and the use of scp; on systems with a shared parallel filesystem, this phase is eliminated entirely since both the source and target nodes can access the checkpoint images directly. Second, the quiesce phase is negligible in all cases: 17 ms for LULESH, 55 ms for CoMD, 9 ms for HPCCG, and 7 ms for miniAMR. This confirms that cooperative session teardown via the MPI Sessions API adds minimal latency. Third, CRIU dump time scales with per-rank memory footprint: miniAMR, which maintains the largest working set due to its refined mesh blocks, incurs 5.4 s of dump time versus 260 ms for LULESH. Excluding image transfer, the core migration machinery (quiesce + dump + restore) completes in 424 ms for LULESH, 316 ms for CoMD, 1.2 s for HPCCG, and 6.2 s for miniAMR. These times represent the irreducible cost of the migration protocol itself and would be the total downtime on a system with shared storage.
5.4
Table 3: Flagship scaling configurations (CoMD, HPCCG). Each set varies one axis (bold) through the shared center 𝐶=4, rpn=16, 𝑀=1. Sets B and C are the 𝑀=1 row and rpn=16 column of the Set D grid.
Non-Migrating Rank Stall (Q3)
Section 5.3 established that migration downtime is dominated by image transfer and that CRIU dump scales with per-rank memory footprint. This subsection characterizes how the resulting downtime scales along three axes: the number of ranks hosted on a migrating node (rpn), the number of compute nodes in the job (𝐶), and the number of proxies evacuated concurrently (𝑀). The central finding is that downtime is governed almost entirely by rpn (the
Set A′ (strong) B′ (weak) C (width)
Nodes (𝐶 ) 1–4 2–4 3–4
Ranks/Node 2–64 4–16 2–16
𝑀 1 1 1–4
Total Ranks 8, 27, 64 8–64 8, 27, 64
rank count on the migrating node) and is insensitive to both total job size and evacuation width. We sweep three axes — per-node rank count (rpn), computenode count (𝐶), and evacuation width (𝑀) — through a shared center (𝐶=4, rpn=16, 𝑀=1); Table 3 lists the configurations for the flexible-rank benchmarks. LULESH’s perfect-cube rank constraint (total ∈ {8, 27, 64}) precludes a fixed-rpn node sweep, so we scale its job by strong and weak scaling instead (Table 4). Figure 5 plots migration downtime against per-node rank count. In every benchmark downtime is linear in rpn and passes through the origin: the fitted 𝑘 · rpn line (black caps) matches the measured bars, and the intercept is within noise of zero (e.g. +0.004 s for CoMD), indicating that the one term expected to grow with total job size, the global KVS address re-exchange, is negligible at these scales. The slope 𝑘 is benchmark-specific, ranging from 1.06 s/rank for CoMD to 1.37 s/rank for HPCCG, and tracks per-rank image size: HPCCG’s sparse 643 blocks yield the largest checkpoints and the steepest slope, CoMD’s compact atom lists the smallest and shallowest. This is the direct consequence of the transfer-bound behavior established in Section 5.3 — at a fixed per-rank footprint, the bytes evacuated from a node scale with its rank count, and transfer time with them. LULESH (top), whose perfect-cube constraint ties rank count to job size, makes the same point across three job sizes at once: bars spanning totals of 8, 27, and 64 ranks all fall on a single 1.15 s/rank line, so a given per-node density incurs the same downtime regardless of the total job it belongs to. Downtime is thus governed by the rank count on the migrating node alone, not the size of the job around it. Figure 6 varies the evacuation width 𝑀 — the number of proxies relocated concurrently — from one to four, holding per-node density fixed. At every rpn level and in all three benchmarks, downtime stays essentially flat as 𝑀 grows: evacuating four nodes at once costs almost the same as evacuating one, even though four times as many ranks are in flight. At rpn=16, going from 𝑀=1 to 𝑀=4 quadruples the ranks relocated (16 → 64) while downtime rises only +2.7% for LULESH, +1.2% for CoMD, and +4.5% for HPCCG. This flatness follows directly from the migration protocol: each source proxy transfers its checkpoint images directly to its paired
XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration
migration downtime (s)
80
1.15 s/rank fit
73.7
64
48 36.6 30.9
32
18.2
16 2.3
0
9.1
10.3
8
9
4.6
2
4
27
16
32
64
total = 8 total = 27 total = 64 ranks per node evacuated (rpn) 48
k ⋅ rpn fit
44.0
migration downtime (s)
40 33.9
32 24
21.9 17.0
16 11.0 8.5
8
5.5
4.3
2.8
2.1
0 2
4
8
16
32
2
4
8
16
CoMD
HPCCG
k = 1.06\,s/rank
k = 1.37\,s/rank
32
ranks per node evacuated (rpn)
Figure 5: Migration downtime versus per-node rank count (rpn): LULESH (top); CoMD and HPCCG (bottom).
target, independently of mpiexec and of every other migrating proxy (Section 4.3), so the 𝑀 transfers proceed in parallel and total downtime is the maximum over them rather than their sum. The small residual increase is ordered by per-rank image size — largest for HPCCG, smallest for CoMD — the signature of mild contention on the shared network fabric as concurrent image volume rises; it remains under 5% even in the most demanding case (HPCCG at rpn=32, 𝑀=4: 128 ranks and ∼5.4 GB pushed across four concurrent links). XMPIaaS thus lets an operator evacuate an arbitrary number of simultaneously preempted nodes for roughly the cost of evacuating one. Figure 7 isolates the effect of job size by holding per-node density fixed (rpn=16, 𝑀=1) and scaling the job from two to seven compute nodes (32 to 112 total ranks). Downtime stays flat throughout, at 16.95 s (±0.4%) for CoMD and 22.04 s (±0.6%) for HPCCG. Since only one node is evacuated, the cost is set by the 16 ranks it holds and is indifferent to how many other nodes participate. A preemption on one instance therefore never scales into a job-wide cost, which is the essence of selective migration.
5.5
Repeated Migration Stability (Q4)
To verify that XMPIaaS introduces no state leaks or performance degradation over multiple migration cycles, we run LULESH (8 ranks, 2 nodes, nx=45) and trigger 8 consecutive migrations at
30-second intervals, alternating the migrating proxy between two target nodes. We record the downtime of each migration event and the resident set size (VmRSS) of the mpiexec coordinator process. Figure 8 shows that migration downtime is effectively constant at 1012 ms with sub-millisecond variance (𝜎 < 1 ms) across all 8 events. The mpiexec VmRSS increases by only 8 kB over the entire sequence, confirming that the proxy replacement and epoch reclamation logic in mpiexec does not leak state. The fd-swap handover replaces the old proxy’s control socket with the new one in place, and the KVS epoch entries are reclaimed rather than accumulated. These results demonstrate that XMPIaaS can sustain repeated migrations indefinitely without degradation, as would be required in a longrunning cloud deployment subject to periodic preemption events.
6
Future Work
XMPIaaS currently requires the application developer to insert XMPI_quiesce calls at safe points, which is a reasonably trivial burden for iterative scientific codes, but a barrier for complex applications with irregular communication patterns or deeply nested call hierarchies where safe points are not immediately obvious. An automatic approach that identifies quiescent program points through static compiler analysis such as detecting loop boundaries where no MPI request is outstanding, or through lightweight runtime profiling that tracks in-flight message counts, would broaden applicability to unmodified legacy codes and remove the programming model requirement entirely. Our implementation targets Hydra’s flat proxy topology, in which all proxies connect directly to mpiexec. Extending the migration protocol to hierarchical proxy trees ( as used by PRRTE (fanout 64) and Cray PALS (fanout 32) at scale ) would require intermediate proxies to forward migration commands rather than receive them directly, and the fd-swap handover would need to propagate through the tree rather than occur at a single level. This extension would enable XMPIaaS to operate on jobs spanning hundreds of nodes, though as discussed in Section 4.4, fewer than 5% of HPC jobs by count exceed the threshold where hierarchical routing engages. The serverless deployment model described in Section 2.1 is a natural extension of our checkpoint/restore machinery. Each MPI rank would run as an independent function invocation, with XMPI_quiesce triggering a checkpoint to shared storage as the function’s time limit approaches. Fresh invocations would restore from the images and resume computation, making execution appear continuous while conforming to the serverless platform’s stateless, time-bounded execution model. The key challenges are adapting the process manager to operate across ephemeral function instances and keeping checkpoint/restore overhead small relative to the function’s available runtime.
7
Related Work
System-level checkpoint/restart. BLCR [24] provided kernellevel checkpoint/restart for Linux clusters and was integrated with several MPI implementations including LAM/MPI [53] and Open MPI, but required kernel module maintenance and was eventually abandoned as kernel APIs evolved. CRIU [16] succeeded BLCR by operating primarily in userspace, but neither tool can handle
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt
20
rpn = 16 (16 → 64 ranks)
rpn = 9
10
5
rpn = 4
rpn = 2
1
2
3
40
4 × ranks, +1.2%
rpn = 16 (16 → 64 ranks)
16
rpn = 8
8
4
1
2
nodes evacuated concurrently (M)
3
4 × ranks, +4.5%
rpn = 16 (16 → 64 ranks)
20
rpn = 8
10
rpn = 4
rpn = 4
4
2
rpn = 32
rpn = 32
32
migration downtime (s)
rpn = 32 4 × ranks, +2.7%
migration downtime (s)
migration downtime (s)
40
5
4
1
2
nodes evacuated concurrently (M)
3
4
nodes evacuated concurrently (M)
Figure 6: Migration downtime versus evacuation width 𝑀: LULESH (left), CoMD (center), HPCCG (right). 32
48
total ranks in job 64
80
96
112
32
48
total ranks in job 64
80
96
112
28
migration downtime (s)
migration downtime (s)
20 16 12
flat: 16.95\,s (±0.4\%)
8 4 0
24 20 16
flat: 22.04\,s (±0.6\%)
12 8 4 0
2
3
4
5
6
7
compute nodes C (rpn = 16, M = 1)
2
3
4
5
6
7
compute nodes C (rpn = 16, M = 1)
1025 1020 1015 1010 1005 1000
3140 3120 3100 3080 3060
Downtime mpiexec VmRSS
1
2
3
4
mpiexec VmRSS (kB)
Downtime (ms)
Figure 7: Migration downtime versus compute-node count 𝐶 at fixed rpn=16, 𝑀=1: CoMD (left), HPCCG (right).
5
6
Migration Number
7
8
Figure 8: Downtime and mpiexec memory footprint over 8 consecutive migrations. Downtime is flat at 1012 ms (𝜎 < 1 ms). Memory grows by only 8 kB total, confirming no state leak.
MPI’s network states, InfiniBand queue pairs, RDMA memory registrations, and device-backed memory mappings are all opaque to their checkpoint machinery. DMTCP [3] extended transparent checkpointing to distributed and multithreaded applications with a coordinator-based architecture, and serves as the foundation for MANA’s split-process approach. Whole-job checkpoint/restart. MANA [18] transparently checkpoints/restores MPI applications by interposing on the entire MPI API interface, maintaining a virtualized bookkeeping of internal MPI states. It employs a split-process architecture that separates the application from the MPI library through a proxy layer, discarding the network contexts that are difficult to preserve transparently at
checkpointing time. MANA 2.0 [57] improved scalability by introducing collective vector clocks to determine safe synchronization points algorithmically. MANA’s intrusive approach comes at considerable deployment complexity that it must track and replay every MPI object’s lifecycle, requires a dedicated DMTCP coordinator daemon [3] alongside the job, and each new transport (TCP, InfiniBand, Slingshot) demands a separate network-virtualization plugin. Application-level checkpointing libraries such as SCR [42] and FTI [6] enable efficient state serialization to local or parallel storage, but require the programmer to manually identify and save all relevant state. Works such as MANA, SCR and FTI all treat recovery as a whole-job event and necessitate a full job restart from the last checkpoint. A preemption on a single instance cascades into a full teardown. Malleability and resource reconfiguration. Malleable MPI frameworks allow a running job to change its process count in response to a resource manager’s decisions. Elastic MPI [14] extends MPICH and Slurm so that applications periodically poll for reconfiguration events, admitting or releasing processes as the scheduler dictates. DMR [31, 32] builds on MPI_Comm_spawn to expand or shrink a job at application-declared reconfiguration points, with a companion library redistributing application data into the new decomposition. Dynamic PSets [30] uses a dynamic resource manager to grant or reclaim an allocation and directs the job’s runtime daemon to spawn processes into it. All three treat the departing rank’s state as discarded, and the application must both run correctly at the new process count and redistribute its data accordingly. Charm++ [34][35] and its MPI interface AMPI [29] overdecomposes applications into virtual ranks, each rank is a serializable object
XMPIaaS: Towards Cloud Native MPI via Cooperative Process Migration
that the load balancer can relocate between nodes [13]. Recent work has extended Charm++ to handle spot instance preemption on cloud platforms [10] and elastic job scheduling [9]. However, AMPI still requires globals privatization and an over-decomposed launch, migrating its own virtual ranks rather than an unmodified MPI rank. Contained Recovery. User-Level Fault Mitigation (ULFM) [11, 38] extends the MPI standard with primitives for detecting process failures and repairing communicators, allowing applications to continue after a rank crashes. However, ULFM provides only the detection and communicator-repair mechanism, it preserves no process state and requires applications to implement their own recovery logic. SPBC [50] combines coordinated checkpointing with message logging, so that a failure rolls back only the affected process group, which replays its logged messages while the rest of the job waits in place rather than restarting. The logging overhead is paid on every message of every run, whether or not a failure ever occurs, and the recovered ranks must still re-execute the work performed since their last checkpoint. Process-level migration. Wang et al. [55] migrate the ranks of a single deteriorating node to a spare while the remaining ranks quiesce and resume in place, using BLCR-based [24] memory precopy and an MPI-level drain of in-flight messages. However, their design presumes a long-deprecated MPI runtime architecture, requesting the destination to already host a daemon in the job’s control plane. Wang et al. [56] migrate a containerized MPI rank group with CRIU. However, their work makes no mention to how MPI runtime handles the migration, and their evaluation only migrates single container MPI job.
8
Conclusion
We present XMPIaaS, a novel cooperative migration framework that enables selective relocation of MPI ranks mid-execution without checkpointing the entire job. By combining application-identified quiesce points with MPI’s own Sessions API for transport teardown, XMPIaaS reduces a migrating rank to an ordinary userspace process, trivially checkpointable by CRIU without MPI-specific plugins or state virtualization. On the process manager side, a lightweight protocol of three new Hydra commands and two out-of-band sentinel strings coordinates the full migration lifecycle: quiescence, selective checkpointing, proxy replacement via a single fd swap, and seamless reconnection through the existing KVS address exchange. Our evaluation across four proxy applications demonstrates that XMPIaaS’s cooperative design achieves its goals. The quiesce phase adds less than 55 ms of latency. The core migration machinery completes in 0.3–6.2 s depending on per-rank memory footprint. Migration downtime depends only on the migrating node’s rank count, and non-migrating ranks are never checkpointed. The system sustains repeated migrations with constant 1012 ms downtime and no state leaks over eight consecutive cycles. Normal-execution overhead is indistinguishable from noise. Two properties make XMPIaaS practical for cloud deployment. First, its per-instance granularity matches the cloud preemption model: only the affected node’s ranks are checkpointed, while the rest of the job continues with minimal disruption. Second, its reliance on standard MPI Sessions API calls on the rank side and
common process management abstractions on the runtime side makes the design portable across MPI implementations. As cloud platforms increasingly host HPC workloads on preemptible infrastructure, XMPIaaS provides a path toward resilient, cost-effective MPI execution without sacrificing the programming model that three decades of scientific software depend on.
References [1] Amazon Web Services. 2025. Amazon EC2 Spot Instances. https://aws.amazon. com/ec2/spot/. Accessed: 2025-05. [2] Amazon Web Services. 2025. AWS Lambda. https://aws.amazon.com/lambda/. Accessed: 2025-05. [3] Jason Ansel, Kapil Arya, and Gene Cooperman. 2009. DMTCP: Transparent Checkpointing for Cluster Computations and the Desktop. In 2009 IEEE International Symposium on Parallel & Distributed Processing (IPDPS). IEEE, 1–12. doi:10.1109/IPDPS.2009.5161063 [4] AWS HPC Blog. 2023. Checkpointing HPC applications using the Spot Instance two-minute notification from Amazon EC2. https: //aws.amazon.com/blogs/hpc/checkpointing-hpc-applications-using-thespot-instance-two-minute-notification-from-amazon-ec2/. Accessed: 2025-05. [5] Pavan Balaji, Darius Buntinas, David Goodell, William Gropp, Jayesh Krishna, Ewing Lusk, and Rajeev Thakur. 2010. PMI: A Scalable Parallel ProcessManagement Interface for Extreme-Scale Systems. In Recent Advances in the Message Passing Interface (EuroMPI 2010) (Lecture Notes in Computer Science, Vol. 6305). Springer, 31–41. doi:10.1007/978-3-642-15646-5_4 [6] Leonardo Bautista-Gomez, Seiji Tsuboi, Dimitri Komatitsch, Franck Cappello, Naoya Maruyama, and Satoshi Matsuoka. 2011. FTI: High Performance Fault Tolerance Interface for Hybrid Systems. In Proceedings of 2011 International Conference for High Performance Computing, Networking, Storage and Analysis (SC ’11). ACM, Article 32, 32 pages. doi:10.1145/2063384.2063427 [7] David E. Bernholdt, Swen Boehm, George Bosilca, Manjunath Gorentla Venkata, Ryan E. Grant, Thomas Naughton, Howard P. Pritchard, Martin Schulz, and Geoffroy R. Vallee. 2020. A Survey of MPI Usage in the US Exascale Computing Project. Concurrency and Computation: Practice and Experience 32, 3 (2020), e4851. doi:10.1002/cpe.4851 [8] David E. Bernholdt, George Bosilca, Aurelien Bouteiller, Ron Brightwell, Jan Ciesko, Matthew G. F. Dosanjh, Giorgis Georgakoudis, Ignacio Laguna, Scott Levy, Thomas Naughton, Stephen L. Olivier, Howard P. Pritchard, Whit Schonbein, Joseph Schuchart, and Amir Shehata. 2024. Taking the MPI Standard and the Open MPI Library to Exascale. The International Journal of High Performance Computing Applications (2024). Online first. doi:10.1177/10943420241265936 [9] Aditya Bhosale, Kavitha Chandrasekar, Laxmikant V. Kale, and Sara KokkilaSchumacher. 2025. An Elastic Job Scheduler for HPC Applications on the Cloud. In SC Workshops (SC-W ’25): Workshops of the International Conference for High Performance Computing, Networking, Storage and Analysis. ACM. arXiv:2510.15147. [10] Aditya Bhosale, Laxmikant V. Kale, and Sara Kokkila-Schumacher. 2025. Efficient and Cost-Effective HPC on the Cloud. In Proceedings of the 34th International Symposium on High-Performance Parallel and Distributed Computing (HPDC ’25). ACM. doi:10.1145/3731545.3744667 [11] Wesley Bland, Aurelien Bouteiller, Thomas Herault, George Bosilca, and Jack Dongarra. 2013. Post-Failure Recovery of MPI Communication Capability: Design and Rationale. The International Journal of High Performance Computing Applications 27, 3 (2013), 244–254. doi:10.1177/1094342013488238 [12] Ralph H. Castain, Joshua Hursey, Aurelien Bouteiller, and Daniel Solt. 2018. PMIx: Process Management for Exascale Environments. Parallel Comput. 79 (2018), 9–29. doi:10.1016/j.parco.2018.08.002 [13] Sayantan Chakravorty, Celso L. Mendes, and Laxmikant V. Kalé. 2006. Proactive Fault Tolerance in MPI Applications via Task Migration. In High Performance Computing (HiPC 2006) (Lecture Notes in Computer Science, Vol. 4297). Springer, 485–496. doi:10.1007/11945918_47 [14] Isaías Comprés, Ao Mo-Hellenbrand, Michael Gerndt, and Hans-Joachim Bungartz. 2016. Infrastructure and API Extensions for Elastic Execution of MPI Applications. In Proceedings of the 23rd European MPI Users’ Group Meeting (EuroMPI 2016). ACM, 82–97. doi:10.1145/2966884.2966917 [15] Marcin Copik, Grzegorz Kwasniewski, Maciej Besta, Michal Müller, and Torsten Hoefler. 2023. rFaaS: RDMA-enabled FaaS Platform for Serverless Scientific Computing. In Proc. IEEE Int. Parallel and Distributed Processing Symp. (IPDPS). 897–907. doi:10.1109/IPDPS54959.2023.00073 [16] CRIU Project. [n. d.]. CRIU – Checkpoint/Restore In Userspace. https://criu.org. Accessed: 2026-07-14. [17] ExMatEx. 2017. CoMD: A Classical Molecular Dynamics Proxy Application. https://github.com/ECP-copa/CoMD. Exascale Co-Design Center for Materials in Extreme Environments. [18] Rohan Garg, Gregory Price, and Gene Cooperman. 2019. MANA for MPI: MPIAgnostic Network-Agnostic Transparent Checkpointing. In Proceedings of the
Shunyu Yao, Dimitrios S. Nikolopoulos, and Ali R. Butt
28th International Symposium on High-Performance Parallel and Distributed Computing (HPDC ’19). ACM, 49–60. doi:10.1145/3307681.3325962 [19] Google Cloud. 2025. Cloud Run Functions. https://cloud.google.com/functions. Accessed: 2025-05. [20] Google Cloud. 2025. Spot VMs. https://cloud.google.com/compute/docs/ instances/spot. Accessed: 2025-05. [21] William Gropp, Ewing Lusk, Nathan Doss, and Anthony Skjellum. 1996. A HighPerformance, Portable Implementation of the MPI Message Passing Interface Standard. Parallel Comput. 22, 6 (1996), 789–828. doi:10.1016/0167-8191(96)000245 [22] William Gropp, Ewing Lusk, and Anthony Skjellum. 1994. Using MPI: Portable Parallel Programming with the Message-Passing Interface. MIT Press. [23] Nishant Gupta, Aditya Bhosale, Jae-Seung Yeom, and Laxmikant V. Kale. 2025. Towards an Adaptive Runtime System for Cloud-Native HPC. In arXiv preprint arXiv:2603.14630. https://arxiv.org/abs/2603.14630 [24] Paul H. Hargrove and Jason C. Duell. 2006. Berkeley Lab Checkpoint/Restart (BLCR) for Linux Clusters. Journal of Physics: Conference Series 46, 1 (2006), 494–499. Proceedings of SciDAC 2006. doi:10.1088/1742-6596/46/1/067 [25] Qiming He, Shunhui Zhou, Ben Kober, Dongping Qiu, Simon Impey, et al. 2012. SpotMPI: A Framework for Auction-Based HPC Computing Using Amazon Spot Instances. In Algorithms and Architectures for Parallel Processing (ICA3PP). Springer, 109–120. doi:10.1007/978-3-642-24669-2_11 [26] Joseph M. Hellerstein, Jose M. Faleiro, Joseph E. Gonzalez, Johann SchleierSmith, Vikram Sreekanti, Alexey Tumanov, and Chenggang Wu. 2019. Serverless Computing: One Step Forward, Two Steps Back. In Proc. 9th Biennial Conf. on Innovative Data Systems Research (CIDR). https://arxiv.org/abs/1812.03651 [27] Michael A. Heroux. 2007. HPCCG: A Simple Conjugate Gradient Benchmark Code. https://github.com/Mantevo/HPCCG. Sandia National Laboratories. [28] Michael A. Heroux, Douglas W. Doerfler, Paul S. Crozier, James M. Willenbring, H. Carter Edwards, Alan Williams, Mahesh Rajan, Eric R. Keiter, Heidi K. Thornquist, and Robert W. Numrich. 2009. Improving Performance via Miniapplications. In Sandia National Laboratories Technical Report. [29] Chao Huang, Orion Lawlor, and Laxmikant V. Kalé. 2004. Adaptive MPI. In Languages and Compilers for Parallel Computing (LCPC 2003) (Lecture Notes in Computer Science, Vol. 2958). Springer, 306–322. doi:10.1007/978-3-540-246442_20 [30] Dominik Huber, Martin Schreiber, Martin Schulz, Howard Pritchard, and Daniel J. Holmes. 2024. Design Principles of Dynamic Resource Management for HighPerformance Parallel Programming Models. arXiv:2403.17107 [cs.DC] doi:10. 48550/arXiv.2403.17107 [31] Sergio Iserte, Rafael Mayo, Enrique S. Quintana-Ortí, Vicenç Beltran, and Antonio J. Peña. 2018. DMR API: Improving Cluster Productivity by Turning Applications into Malleable. Parallel Comput. 78 (2018), 54–66. doi:10.1016/j. parco.2018.07.006 [32] Sergio Iserte, Rafael Mayo, Enrique S. Quintana-Ortí, and Antonio J. Peña. 2021. DMRlib: Easy-Coding and Efficient Resource Management for Job Malleability. IEEE Trans. Comput. 70, 9 (2021), 1443–1457. doi:10.1109/TC.2020.3022933 [33] Eric Jonas, Johann Schleier-Smith, Vikram Sreekanti, Chia-Che Tsai, Anurag Khandelwal, Qifan Pu, Vaishaal Shankar, Joao Carreira, Karl Krauth, Neeraja Yadwadkar, Joseph E. Gonzalez, Raluca Ada Popa, Ion Stoica, and David A. Patterson. 2019. Cloud Programming Simplified: A Berkeley View on Serverless Computing. Technical Report UCB/EECS-2019-3. EECS Department, University of California, Berkeley. https://arxiv.org/abs/1902.03383 [34] Laxmikant V. Kalé and Sanjeev Krishnan. 1993. CHARM++: A Portable Concurrent Object Oriented System Based on C++. In Proceedings of the Eighth Annual Conference on Object-Oriented Programming Systems, Languages, and Applications (OOPSLA ’93). ACM, 91–108. doi:10.1145/165854.165874 [35] Laxmikant V. Kalé and Sanjeev Krishnan. 1996. Charm++: Parallel Programming with Message-Driven Objects. In Parallel Programming using C++, Gregory V. Wilson and Paul Lu (Eds.). MIT Press, 175–213. [36] Ian Karlin, Jeff Keasler, and Rob Neely. 2013. LULESH 2.0 Updates and Changes. Technical Report LLNL-TR-641973. Lawrence Livermore National Laboratory. [37] Ignacio Laguna, David F. Richards, Todd Gamblin, Martin Schulz, and Bronis R. de Supinski. 2016. Evaluating and Extending User-Level Fault Tolerance in MPI Applications. In Int. Journal of High Performance Computing Applications, Vol. 30. 305–319. doi:10.1177/1094342015623623 [38] Nuria Losada, Patricia González, María J. Martín, George Bosilca, Aurélien Bouteiller, and Keita Teranishi. 2020. Fault Tolerance of MPI Applications in Exascale Systems: The ULFM Solution. Future Generation Computer Systems 106 (2020), 467–481. doi:10.1016/j.future.2020.01.026 [39] Xupeng Miao, Chunan Shi, Jiangfei Duan, Xiaoli Xi, Dahua Lin, Bin Cui, and Zhihao Jia. 2024. SpotServe: Serving Generative Large Language Models on Preemptible Instances. In Proc. 29th ACM Int. Conf. on Architectural Support for Programming Languages and Operating Systems (ASPLOS). doi:10.1145/3620666. 3651383 [40] Microsoft Azure. 2025. Azure Functions. https://azure.microsoft.com/en-us/ products/functions. Accessed: 2025-05.
[41] Microsoft Azure. 2025. Azure Spot Virtual Machines. https://azure.microsoft. com/en-us/products/virtual-machines/spot. Accessed: 2025-05. [42] Adam Moody, Greg Bronevetsky, Kathryn Mohror, and Bronis R. de Supinski. 2010. Design, Modeling, and Evaluation of a Scalable Multi-level Checkpointing System. In SC ’10: Proceedings of the 2010 ACM/IEEE International Conference for High Performance Computing, Networking, Storage and Analysis. IEEE, 1–11. doi:10.1109/SC.2010.18 [43] MPI Forum. 2021. MPI: A Message-Passing Interface Standard, Version 4.0. Technical Report. Message Passing Interface Forum. https://www.mpi-forum.org/ docs/mpi-4.0/mpi40-report.pdf [44] MPICH Project. 2025. Hydra Process Management Framework. https:// wiki.mpich.org/mpich/index.php/Hydra_Process_Management_Framework. Accessed: 2025-05. [45] Marco A. S. Netto, Rodrigo N. Calheiros, Eduardo R. Rodrigues, Renato L. F. Cunha, and Rajkumar Buyya. 2018. HPC Cloud for Scientific and Business Applications: Taxonomy, Vision, and Research Challenges. Comput. Surveys 51, 1 (2018), 1–29. doi:10.1145/3150224 [46] Bogdan Nicolae, Adam Moody, Elsa Gonsiorowski, Kathryn Mohror, and Franck Cappello. 2019. VeloC: Towards High Performance Adaptive Asynchronous Checkpointing at Large Scale. In Proc. IEEE Int. Parallel and Distributed Processing Symp. (IPDPS). 911–920. doi:10.1109/IPDPS.2019.00046 [47] Tapasya Patki et al. 2025. Workload Characterization of NERSC’s Cori Supercomputer. arXiv preprint arXiv:2501.12464 (2025). https://arxiv.org/abs/2501.12464 [48] Maksym Planeta, Jan Bierbaum, Leo Sahaya Daphne Antony, Torsten Hoefler, and Hermann Härtig. 2021. MigrOS: Transparent Operating Systems Live Migration Support for Containerised RDMA-Applications. In Proceedings of the 2021 USENIX Annual Technical Conference (USENIX ATC ’21). USENIX Association, 47–63. [49] Gonzalo P. Rodrigo, Erik Elmroth, Per-Olov Ostberg, and Lavanya Ramakrishnan. 2018. Towards Understanding HPC Users and Systems: A NERSC Case Study. J. Parallel and Distrib. Comput. 111 (2018), 206–217. doi:10.1016/j.jpdc.2017.09.002 [50] Thomas Ropars, Tatiana V. Martsinkevich, Amina Guermouche, André Schiper, and Franck Cappello. 2013. SPBC: Leveraging the Characteristics of MPI HPC Applications for Scalable Checkpointing. In Proceedings of the International Conference on High Performance Computing, Networking, Storage and Analysis (SC ’13). ACM, Article 8, 12 pages. doi:10.1145/2503210.2503271 [51] Aparna Sasidharan and Marc Snir. 2016. MiniAMR – A Miniapp for Adaptive Mesh Refinement. [52] Josef Spillner, Yaroslav Dorodko, and Jia Li. 2018. FaaSter, Better, Cheaper: The Prospect of Serverless Scientific Computing and HPC. In Proc. IEEE/ACM Int. Conf. on Utility and Cloud Computing Companion (UCC Companion). 154–160. doi:10.1109/UCC-Companion.2018.00076 [53] Jeffrey M. Squyres and Andrew Lumsdaine. 2003. A Component Architecture for LAM/MPI. In Recent Advances in Parallel Virtual Machine and Message Passing Interface (EuroPVM/MPI 2003) (Lecture Notes in Computer Science, Vol. 2840). Springer, 379–387. doi:10.1007/978-3-540-39924-7_52 [54] William Voorsluys and Rajkumar Buyya. 2012. Reliable Provisioning of Spot Instances for Compute-Intensive Applications. In Proc. IEEE Int. Conf. on Advanced Information Networking and Applications (AINA). 542–549. doi:10.1109/AINA. 2012.106 [55] Chao Wang, Frank Mueller, Christian Engelmann, and Stephen L. Scott. 2008. Proactive Process-Level Live Migration in HPC Environments. In SC ’08: Proceedings of the 2008 ACM/IEEE Conference on Supercomputing. IEEE, 1–12. doi:10.1145/1413370.1413414 [56] Xiaoning Wang, Yining Zhao, Shasha Lu, and Haili Xiao. 2025. Practice and Observation: Live Migration for MPI Workload. CCF Transactions on High Performance Computing 7, 5 (2025), 447–464. doi:10.1007/s42514-025-00227-0 [57] Yao Xu, Zhengji Zhao, Rohan Garg, Harsh Khetawat, Rebecca Hartman-Baker, and Gene Cooperman. 2021. MANA-2.0: A Future-Proof Design for Transparent Checkpointing of MPI at Scale. In 2021 SC Workshops Supplementary Proceedings (SCWS). IEEE, 68–78. [58] Andy B. Yoo, Morris A. Jette, and Mark Grondona. 2003. SLURM: Simple Linux Utility for Resource Management. In Job Scheduling Strategies for Parallel Processing (JSSPP) (Lecture Notes in Computer Science, Vol. 2862). Springer, 44–60. doi:10.1007/10968987_3