[FIX] Fixed RMA win creation in order to get it scaling

This commit is contained in:
Stack-1
2026-08-10 13:12:57 +02:00
parent 84d22c8429
commit 2a767a4097
2 changed files with 231 additions and 36 deletions
+156 -36
View File
@@ -1092,10 +1092,11 @@ contains
call mpi_win_unlock_all(rma_handle%win, iret)
rma_handle%window_open = .false.
end if
if (rma_handle%window_ready) then
call mpi_win_free(rma_handle%win, iret)
rma_handle%win = mpi_win_null
rma_handle%window_ready = .false.
! Dynamic window: detach the old buffer and keep the window. Detach
! is a local call, so growing the buffer no longer costs a collective.
if (rma_handle%buf_attached) then
call mpi_win_detach(rma_handle%win, y%combuf, iret)
rma_handle%buf_attached = .false.
end if
call y%new_buffer(ione*size(comm_indexes%v), info)
if (info /= psb_success_) then
@@ -1103,11 +1104,13 @@ contains
end if
end if
element_bytes = storage_size(y%combuf(1))/8
if (.not. rma_handle%window_ready) then
element_bytes = storage_size(y%combuf(1))/8
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_create(y%combuf, exposed_bytes, element_bytes, &
& mpi_info_null, ctxt%get_mpic(), rma_handle%win, iret)
! Created once and kept for the life of the handle. On a dynamic
! window buffers come and go through local attach/detach, so this
! collective is paid a single time instead of on every buffer change.
call mpi_win_create_dynamic(mpi_info_null, ctxt%get_mpic(), &
& rma_handle%win, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
@@ -1115,6 +1118,29 @@ contains
end if
rma_handle%window_ready = .true.
end if
if (.not. rma_handle%buf_attached) then
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_attach(rma_handle%win, y%combuf, exposed_bytes, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
call mpi_get_address(y%combuf, rma_handle%my_buf_addr, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
! Targets are addressed absolutely on a dynamic window, so neighbours
! must be told the new base before any Get/Put can reference it.
call rma_handle%publish_buf_addr(ctxt%get_mpic(), info)
if (info /= psb_success_) then
call psb_errpush(info,name); goto 9999
end if
rma_handle%buf_attached = .true.
end if
end if
if (buffer_size > 0) then
@@ -1158,7 +1184,12 @@ contains
end if
if (recv_count > 0) then
remote_disp = int(remote_base - 1, kind=MPI_ADDRESS_KIND)
! Dynamic window: disp_unit is 1 and the displacement is an
! absolute address in the target's space, so the offset has to be
! expressed in bytes rather than in window elements.
remote_disp = rma_handle%peer_buf_addr(neighbor_idx) &
& + int(remote_base - 1, kind=MPI_ADDRESS_KIND) &
& * int(element_bytes, kind=MPI_ADDRESS_KIND)
call mpi_get(y%combuf(recv_pos), recv_count, psb_mpi_r_dpk_, prc_rank, remote_disp, recv_count, psb_mpi_r_dpk_, &
& rma_handle%win, iret)
if (iret /= mpi_success) then
@@ -1314,10 +1345,11 @@ contains
call mpi_win_unlock_all(rma_handle%win, iret)
rma_handle%window_open = .false.
end if
if (rma_handle%window_ready) then
call mpi_win_free(rma_handle%win, iret)
rma_handle%win = mpi_win_null
rma_handle%window_ready = .false.
! Dynamic window: detach the old buffer and keep the window. Detach
! is a local call, so growing the buffer no longer costs a collective.
if (rma_handle%buf_attached) then
call mpi_win_detach(rma_handle%win, y%combuf, iret)
rma_handle%buf_attached = .false.
end if
call y%new_buffer(ione*size(comm_indexes%v), info)
if (info /= psb_success_) then
@@ -1325,11 +1357,13 @@ contains
end if
end if
element_bytes = storage_size(y%combuf(1))/8
if (.not. rma_handle%window_ready) then
element_bytes = storage_size(y%combuf(1))/8
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_create(y%combuf, exposed_bytes, element_bytes, &
& mpi_info_null, ctxt%get_mpic(), rma_handle%win, iret)
! Created once and kept for the life of the handle. On a dynamic
! window buffers come and go through local attach/detach, so this
! collective is paid a single time instead of on every buffer change.
call mpi_win_create_dynamic(mpi_info_null, ctxt%get_mpic(), &
& rma_handle%win, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
@@ -1337,6 +1371,29 @@ contains
end if
rma_handle%window_ready = .true.
end if
if (.not. rma_handle%buf_attached) then
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_attach(rma_handle%win, y%combuf, exposed_bytes, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
call mpi_get_address(y%combuf, rma_handle%my_buf_addr, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
! Targets are addressed absolutely on a dynamic window, so neighbours
! must be told the new base before any Get/Put can reference it.
call rma_handle%publish_buf_addr(ctxt%get_mpic(), info)
if (info /= psb_success_) then
call psb_errpush(info,name); goto 9999
end if
rma_handle%buf_attached = .true.
end if
end if
if (buffer_size > 0) then
@@ -1381,7 +1438,12 @@ contains
end if
if (send_count > 0) then
remote_disp = int(remote_base - 1, kind=MPI_ADDRESS_KIND)
! Dynamic window: disp_unit is 1 and the displacement is an
! absolute address in the target's space, so the offset has to be
! expressed in bytes rather than in window elements.
remote_disp = rma_handle%peer_buf_addr(neighbor_idx) &
& + int(remote_base - 1, kind=MPI_ADDRESS_KIND) &
& * int(element_bytes, kind=MPI_ADDRESS_KIND)
call mpi_put(y%combuf(send_pos), send_count, psb_mpi_r_dpk_, prc_rank, remote_disp, send_count, psb_mpi_r_dpk_, &
& rma_handle%win, iret)
if (iret /= mpi_success) then
@@ -2364,10 +2426,11 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
call mpi_win_unlock_all(rma_handle%win, iret)
rma_handle%window_open = .false.
end if
if (rma_handle%window_ready) then
call mpi_win_free(rma_handle%win, iret)
rma_handle%win = mpi_win_null
rma_handle%window_ready = .false.
! Dynamic window: detach the old buffer and keep the window. Detach
! is a local call, so growing the buffer no longer costs a collective.
if (rma_handle%buf_attached) then
call mpi_win_detach(rma_handle%win, y%combuf, iret)
rma_handle%buf_attached = .false.
end if
call y%new_buffer(buffer_size, info)
if (info /= psb_success_) then
@@ -2375,11 +2438,13 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
end if
end if
element_bytes = storage_size(y%combuf(1))/8
if (.not. rma_handle%window_ready) then
element_bytes = storage_size(y%combuf(1))/8
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_create(y%combuf, exposed_bytes, element_bytes, &
& mpi_info_null, ctxt%get_mpic(), rma_handle%win, iret)
! Created once and kept for the life of the handle. On a dynamic
! window buffers come and go through local attach/detach, so this
! collective is paid a single time instead of on every buffer change.
call mpi_win_create_dynamic(mpi_info_null, ctxt%get_mpic(), &
& rma_handle%win, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
@@ -2387,6 +2452,29 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
end if
rma_handle%window_ready = .true.
end if
if (.not. rma_handle%buf_attached) then
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_attach(rma_handle%win, y%combuf, exposed_bytes, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
call mpi_get_address(y%combuf, rma_handle%my_buf_addr, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
! Targets are addressed absolutely on a dynamic window, so neighbours
! must be told the new base before any Get/Put can reference it.
call rma_handle%publish_buf_addr(ctxt%get_mpic(), info)
if (info /= psb_success_) then
call psb_errpush(info,name); goto 9999
end if
rma_handle%buf_attached = .true.
end if
end if
if (buffer_size > 0) then
@@ -2424,7 +2512,10 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
goto 9999
end if
if (recv_count > 0) then
remote_disp = int((remote_base - 1) * n, kind=MPI_ADDRESS_KIND)
! Dynamic window: absolute target address, offset in bytes.
remote_disp = rma_handle%peer_buf_addr(neighbor_idx) &
& + int((remote_base - 1) * n, kind=MPI_ADDRESS_KIND) &
& * int(element_bytes, kind=MPI_ADDRESS_KIND)
call mpi_get(y%combuf(recv_pos), recv_count*n, psb_mpi_r_dpk_, prc_rank, remote_disp, recv_count*n, psb_mpi_r_dpk_, &
& rma_handle%win, iret)
if (iret /= mpi_success) then
@@ -2578,10 +2669,11 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
call mpi_win_unlock_all(rma_handle%win, iret)
rma_handle%window_open = .false.
end if
if (rma_handle%window_ready) then
call mpi_win_free(rma_handle%win, iret)
rma_handle%win = mpi_win_null
rma_handle%window_ready = .false.
! Dynamic window: detach the old buffer and keep the window. Detach
! is a local call, so growing the buffer no longer costs a collective.
if (rma_handle%buf_attached) then
call mpi_win_detach(rma_handle%win, y%combuf, iret)
rma_handle%buf_attached = .false.
end if
call y%new_buffer(buffer_size, info)
if (info /= psb_success_) then
@@ -2589,11 +2681,13 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
end if
end if
element_bytes = storage_size(y%combuf(1))/8
if (.not. rma_handle%window_ready) then
element_bytes = storage_size(y%combuf(1))/8
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_create(y%combuf, exposed_bytes, element_bytes, &
& mpi_info_null, ctxt%get_mpic(), rma_handle%win, iret)
! Created once and kept for the life of the handle. On a dynamic
! window buffers come and go through local attach/detach, so this
! collective is paid a single time instead of on every buffer change.
call mpi_win_create_dynamic(mpi_info_null, ctxt%get_mpic(), &
& rma_handle%win, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
@@ -2601,6 +2695,29 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
end if
rma_handle%window_ready = .true.
end if
if (.not. rma_handle%buf_attached) then
exposed_bytes = int(size(y%combuf),kind=MPI_ADDRESS_KIND) * int(element_bytes,kind=MPI_ADDRESS_KIND)
call mpi_win_attach(rma_handle%win, y%combuf, exposed_bytes, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
call mpi_get_address(y%combuf, rma_handle%my_buf_addr, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
call psb_errpush(info,name,m_err=(/iret/))
goto 9999
end if
! Targets are addressed absolutely on a dynamic window, so neighbours
! must be told the new base before any Get/Put can reference it.
call rma_handle%publish_buf_addr(ctxt%get_mpic(), info)
if (info /= psb_success_) then
call psb_errpush(info,name); goto 9999
end if
rma_handle%buf_attached = .true.
end if
end if
if (buffer_size > 0) then
@@ -2639,7 +2756,10 @@ end subroutine psi_dswap_neighbor_topology_multivect_persistent
goto 9999
end if
if (send_count > 0) then
remote_disp = int((remote_base - 1) * n, kind=MPI_ADDRESS_KIND)
! Dynamic window: absolute target address, offset in bytes.
remote_disp = rma_handle%peer_buf_addr(neighbor_idx) &
& + int((remote_base - 1) * n, kind=MPI_ADDRESS_KIND) &
& * int(element_bytes, kind=MPI_ADDRESS_KIND)
call mpi_put(y%combuf(send_pos), send_count*n, psb_mpi_r_dpk_, prc_rank, remote_disp, send_count*n, psb_mpi_r_dpk_, &
& rma_handle%win, iret)
if (iret /= mpi_success) then
@@ -39,6 +39,23 @@ module psb_comm_rma_mod
integer(psb_mpk_), allocatable :: notify_buf(:)
integer(psb_mpk_), allocatable :: notify_recv_reqs(:)
integer(psb_mpk_), allocatable :: notify_send_reqs(:)
!
! Dynamic-window support.
!
! MPI_Win_create is collective over the whole communicator and has to be
! repeated whenever the exposed buffer changes. Profiling on 448 ranks put
! it at ~0.9 s per call, against milliseconds for the MPI_Get that actually
! moves the data. With a dynamic window that collective is paid once and
! buffers are attached and detached locally.
!
! The price is that on a dynamic window the target displacement is an
! absolute address in the target's address space rather than an offset in
! window units, so every rank must publish the address of its own buffer to
! its neighbours each time it (re)attaches.
!
logical :: buf_attached = .false.
integer(kind=MPI_ADDRESS_KIND) :: my_buf_addr = 0
integer(kind=MPI_ADDRESS_KIND), allocatable :: peer_buf_addr(:)
contains
procedure, pass :: init => psb_comm_rma_init
procedure, pass :: free => psb_comm_rma_free
@@ -47,6 +64,7 @@ module psb_comm_rma_mod
procedure, pass :: init_memory_buffer_layout_tran => psb_comm_rma_ini_memory_buffer_layout_tran
procedure, pass :: set_swap_status => psb_comm_rma_set_swap_status
procedure, pass :: get_swap_status => psb_comm_rma_get_swap_status
procedure, pass :: publish_buf_addr => psb_comm_rma_publish_buf_addr
end type psb_comm_rma_handle
contains
@@ -63,6 +81,9 @@ contains
this%window_open = .false.
this%win_base = c_null_ptr
this%win_nelem = 0
this%buf_attached = .false.
this%my_buf_addr = 0
if (allocated(this%peer_buf_addr)) deallocate(this%peer_buf_addr)
this%layout_ready = .false.
this%layout_nnbr = -1
this%layout_send = -1
@@ -278,6 +299,11 @@ contains
end if
this%win_base = c_null_ptr
this%win_nelem = 0
! Freeing the window drops whatever is attached to it, so the attachment
! bookkeeping goes with it.
this%buf_attached = .false.
this%my_buf_addr = 0
if (allocated(this%peer_buf_addr)) deallocate(this%peer_buf_addr)
this%window_ready = .false.
this%layout_ready = .false.
this%layout_nnbr = -1
@@ -286,6 +312,55 @@ contains
call this%clear_memory_buffer_layout(info)
end subroutine psb_comm_rma_free
!
! Publish the address of the locally attached buffer to the neighbours, and
! collect theirs. On a dynamic window MPI_Get/MPI_Put take an absolute address
! in the target's address space, so this has to be redone whenever the buffer
! is re-attached. Point-to-point over the neighbour list, reusing the metadata
! tag already used for the displacement exchange: no collective.
!
subroutine psb_comm_rma_publish_buf_addr(this, icomm, info)
#ifdef PSB_MPI_MOD
use mpi
#endif
implicit none
#ifdef PSB_MPI_H
include 'mpif.h'
#endif
class(psb_comm_rma_handle), intent(inout) :: this
integer(psb_mpk_), intent(in) :: icomm
integer(psb_ipk_), intent(out) :: info
integer(psb_ipk_) :: k, nnbr
integer(psb_mpk_) :: prc_rank, iret
integer(psb_mpk_) :: p2pstat(mpi_status_size)
info = psb_success_
if (.not.allocated(this%peer_mpi_rank)) return
nnbr = size(this%peer_mpi_rank)
if (allocated(this%peer_buf_addr)) then
if (size(this%peer_buf_addr) /= nnbr) deallocate(this%peer_buf_addr)
end if
if (.not.allocated(this%peer_buf_addr)) then
allocate(this%peer_buf_addr(nnbr), stat=info)
if (info /= 0) then
info = psb_err_alloc_dealloc_
return
end if
end if
do k = 1, nnbr
prc_rank = int(this%peer_mpi_rank(k), psb_mpk_)
call mpi_sendrecv(this%my_buf_addr, 1, mpi_aint, prc_rank, psb_rma_meta_tag, &
& this%peer_buf_addr(k), 1, mpi_aint, prc_rank, psb_rma_meta_tag, &
& icomm, p2pstat, iret)
if (iret /= mpi_success) then
info = psb_err_mpi_error_
return
end if
end do
end subroutine psb_comm_rma_publish_buf_addr
! Transpose variant: peer_send_* is filled from comm_list RECV area,
! peer_recv_* from comm_list SEND area. Metadata exchange tells peers our
! recv displacement so they can GET/PUT to the correct location in swaptran.