Skip to content

[Core][MPI] Add non-blocking (asynchronous) operations to DataCommunicator - #14835

Open
loumalouomega wants to merge 7 commits into
KratosMultiphysics:masterfrom
loumalouomega:core/async-data-communicator
Open

loumalouomega wants to merge 7 commits into
KratosMultiphysics:masterfrom
loumalouomega:core/async-data-communicator

Conversation

@loumalouomega

@loumalouomega loumalouomega commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

📝 Description

DataCommunicator only supported blocking calls. This PR adds non-blocking point-to-point and collective operations. They return a new MPI-agnostic DataCommunicatorRequest, which is the core's equivalent of MPI_Request.

Key changes

  • DataCommunicatorRequest (kratos/includes/data_communicator_request.h): a move-only handle with Wait(), Test(), IsCompleted() and a static WaitAll(). A default-constructed request counts as already completed. Destroying a pending request waits for it to finish.
  • New DataCommunicator methods:
    • ISend / IRecv (MPI_Isend / MPI_Irecv)
    • IBarrier (MPI_Ibarrier)
    • IBroadcast (MPI_Ibcast)
    • ISumAll / IMinAll / IMaxAll (MPI_Iallreduce)
    • Types: the same as Send/Recv (char, int, unsigned int, long unsigned int, double, array_1d, Vector, Matrix, plus std::vector of each), and std::string for send/recv/broadcast.
  • MPIDataCommunicator: the request owns the MPI_Request and the message buffers. Received data that isn't contiguous (e.g. std::vector<Vector>) is copied back on completion.
  • Serial DataCommunicator: operations complete immediately; IRecv throws, like Recv.
  • Python: int/double variants plus ISendString, IRecvString and IBarrier. They return a DataCommunicatorRequest that keeps the buffers alive; read the result with GetResult().

Usage rules: output buffers must be sized in advance, including the inner shapes of Vector/Matrix entries, because no size information is sent. Pair ISend with IRecv.

std::vector<double> recv(n);
std::vector<DataCommunicatorRequest> requests;
requests.push_back(r_comm.IRecv(recv, prev_rank, tag));
requests.push_back(r_comm.ISend(send, next_rank, tag));
// ... overlap computation ...
DataCommunicatorRequest::WaitAll(requests);

double global_sum = 0.0;
auto request = r_comm.ISumAll(local_value, global_sum);
// ...
request.Wait();
req = comm.IRecvDoubles(2, recv_rank)
comm.ISendDoubles([1.0, 2.0], send_rank).Wait()
req.Wait()
values = req.GetResult()

total = comm.ISumAll(rank)
KM.DataCommunicatorRequest.WaitAll([total])
print(total.GetResult())

Validation

New tests:

  • kratos/tests/cpp_tests/includes/test_data_communicator_request.cpp (serial)
  • kratos/mpi/tests/cpp_tests/sources/test_mpi_data_communicator_async.cpp (MPI, ring exchange, all types, Test() polling, WaitAll, completion on destruction, collectives)
  • testAsync* cases in kratos/mpi/tests/test_mpi_data_communicator_python.py

Results:

  • Serial and MPI tests pass with mpirun -np 1/2/3.
  • The full KratosMPICoreTest suite has no failures.

🆕 Changelog

  • Added DataCommunicatorRequest, an MPI-agnostic handle for non-blocking communication
  • Added ISend, IRecv, IBarrier, IBroadcast, ISumAll, IMinAll and IMaxAll to DataCommunicator / MPIDataCommunicator
  • Exposed the non-blocking operations to Python

@loumalouomega loumalouomega added Enhancement Kratos Core Parallel-MPI Distributed memory parallelism for HPC / clusters labels Oct 1, 2026
@loumalouomega
loumalouomega marked this pull request as ready for review October 2, 2026 11:37
@loumalouomega
loumalouomega requested review from a team as code owners October 2, 2026 11:37

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Enhancement Kratos Core Parallel-MPI Distributed memory parallelism for HPC / clusters Python

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant