Arcane  4.2.3.0
Developer documentation
Loading...
Searching...
No Matches
MpiParallelMng.cc
1// -*- tab-width: 2; indent-tabs-mode: nil; coding: utf-8-with-signature -*-
2//-----------------------------------------------------------------------------
3// Copyright 2000-2026 CEA (www.cea.fr) IFPEN (www.ifpenergiesnouvelles.com)
4// See the top-level COPYRIGHT file for details.
5// SPDX-License-Identifier: Apache-2.0
6//-----------------------------------------------------------------------------
7/*---------------------------------------------------------------------------*/
8/* MpiParallelMng.cc (C) 2000-2026 */
9/* */
10/* Parallelism manager using MPI. */
11/*---------------------------------------------------------------------------*/
12/*---------------------------------------------------------------------------*/
13
14#include "arcane/utils/Collection.h"
15#include "arcane/utils/Enumerator.h"
16#include "arcane/utils/ScopedPtr.h"
17#include "arcane/utils/PlatformUtils.h"
18#include "arcane/utils/TimeoutException.h"
19#include "arcane/utils/NotImplementedException.h"
20#include "arcane/utils/ArgumentException.h"
21#include "arcane/utils/ITraceMng.h"
22#include "arcane/utils/ValueConvert.h"
23#include "arcane/utils/Exception.h"
24#include "arcane/utils/HPReal.h"
25
26#include "arcane/core/IIOMng.h"
27#include "arcane/core/Timer.h"
28#include "arcane/core/IItemFamily.h"
29#include "arcane/core/IParallelTopology.h"
30#include "arcane/core/parallel/IStat.h"
31#include "arcane/core/internal/SerializeMessage.h"
32#include "arcane/core/internal/ParallelMngInternal.h"
33#include "arcane/core/internal/MachineShMemWinMemoryAllocator.h"
34
35#include "arcane/parallel/mpi/MpiParallelMng.h"
36#include "arcane/parallel/mpi/MpiParallelDispatch.h"
37#include "arcane/parallel/mpi/MpiTimerMng.h"
38#include "arcane/parallel/mpi/MpiSerializeMessage.h"
39#include "arcane/parallel/mpi/MpiParallelNonBlockingCollective.h"
40#include "arcane/parallel/mpi/MpiDatatype.h"
41#include "arcane/parallel/mpi/IVariableSynchronizerMpiCommunicator.h"
42
43#include "arcane/impl/ParallelReplication.h"
44#include "arcane/impl/SequentialParallelMng.h"
45#include "arcane/impl/internal/ParallelMngUtilsFactoryBase.h"
46#include "arcane/impl/internal/VariableSynchronizer.h"
47
48#include "arccore/message_passing_mpi/internal/MpiMessagePassingMng.h"
49#include "arccore/message_passing_mpi/internal/MpiSerializeDispatcher.h"
50#include "arccore/message_passing_mpi/internal/MpiRequestList.h"
51#include "arccore/message_passing_mpi/internal/MpiAdapter.h"
52#include "arccore/message_passing_mpi/internal/MpiLock.h"
53#include "arccore/message_passing_mpi/internal/MpiMachineShMemWinBaseInternalCreator.h"
54#include "arccore/message_passing_mpi/internal/MpiContigMachineShMemWinBaseInternal.h"
55#include "arccore/message_passing_mpi/internal/MpiMachineShMemWinBaseInternal.h"
56#include "arccore/message_passing/Dispatchers.h"
58#include "arccore/message_passing/internal/SerializeMessageList.h"
59
60#include "arcane_packages.h"
61
62//#define ARCANE_TRACE_MPI
63
64/*---------------------------------------------------------------------------*/
65/*---------------------------------------------------------------------------*/
66
67namespace Arcane
68{
69using namespace Arcane::MessagePassing;
70using namespace Arcane::MessagePassing::Mpi;
72
73/*---------------------------------------------------------------------------*/
74/*---------------------------------------------------------------------------*/
75
76extern "C++" IIOMng*
77arcaneCreateIOMng(IParallelMng* psm);
78
79#if defined(ARCANE_HAS_MPI_NEIGHBOR)
80// Defined in MpiNeighborVariableSynchronizeDispatcher
82arcaneCreateMpiNeighborVariableSynchronizerFactory(MpiParallelMng* mpi_pm,
83 Ref<IVariableSynchronizerMpiCommunicator> synchronizer_communicator);
84#endif
86arcaneCreateMpiBlockVariableSynchronizerFactory(MpiParallelMng* mpi_pm, Int32 block_size, Int32 nb_sequence);
88arcaneCreateMpiVariableSynchronizerFactory(MpiParallelMng* mpi_pm);
90arcaneCreateMpiDirectSendrecvVariableSynchronizerFactory(MpiParallelMng* mpi_pm);
92arcaneCreateMpiLegacyVariableSynchronizerFactory(MpiParallelMng* mpi_pm);
93#if defined(ARCANE_HAS_PACKAGE_NCCL)
95arcaneCreateNCCLVariableSynchronizerFactory(IParallelMng* mpi_pm);
96#endif
97
98/*---------------------------------------------------------------------------*/
99/*---------------------------------------------------------------------------*/
100
101MpiParallelMngBuildInfo::
102MpiParallelMngBuildInfo(MPI_Comm comm, MPI_Comm machine_comm)
103: is_parallel(false)
104, comm_rank(MessagePassing::A_NULL_RANK)
105, comm_nb_rank(0)
106, mpi_comm(comm)
107, mpi_machine_comm(machine_comm)
108, is_mpi_comm_owned(true)
109{
110 ::MPI_Comm_rank(comm, &comm_rank);
111 ::MPI_Comm_size(comm, &comm_nb_rank);
112}
113
114/*---------------------------------------------------------------------------*/
115/*---------------------------------------------------------------------------*/
116
117/*---------------------------------------------------------------------------*/
118/*---------------------------------------------------------------------------*/
122class VariableSynchronizerMpiCommunicator
124{
125 public:
126
127 explicit VariableSynchronizerMpiCommunicator(MpiParallelMng* pm)
128 : m_mpi_parallel_mng(pm)
129 {}
130 ~VariableSynchronizerMpiCommunicator() override
131 {
132 _checkFreeCommunicator();
133 }
134 MPI_Comm communicator() const override
135 {
136 return m_topology_communicator;
137 }
138 void compute(VariableSynchronizer* var_syncer) override
139 {
140 Int32ConstArrayView comm_ranks = var_syncer->communicatingRanks();
141 const Int32 nb_message = comm_ranks.size();
142
143 MpiParallelMng* pm = m_mpi_parallel_mng;
144
145 MPI_Comm old_comm = pm->communicator();
146
147 UniqueArray<int> destinations(nb_message);
148 for (Integer i = 0; i < nb_message; ++i) {
149 destinations[i] = comm_ranks[i];
150 }
151
152 _checkFreeCommunicator();
153
154 int r = MPI_Dist_graph_create_adjacent(old_comm, nb_message, destinations.data(), MPI_UNWEIGHTED,
155 nb_message, destinations.data(), MPI_UNWEIGHTED,
156 MPI_INFO_NULL, 0, &m_topology_communicator);
157
158 if (r != MPI_SUCCESS)
159 ARCANE_FATAL("Error '{0}' in MPI_Dist_graph_create", r);
160
161 // Checks that the rank order for the MPI implementation is the same as the one we have in
162 // the VariableSynchronizer.
163 {
164 int indegree = 0;
165 int outdegree = 0;
166 int weighted = 0;
167 MPI_Dist_graph_neighbors_count(m_topology_communicator, &indegree, &outdegree, &weighted);
168
169 if (indegree != nb_message)
170 ARCANE_FATAL("Bad value '{0}' for 'indegree' (expected={1})", indegree, nb_message);
171 if (outdegree != nb_message)
172 ARCANE_FATAL("Bad value '{0}' for 'outdegree' (expected={1})", outdegree, nb_message);
173
174 UniqueArray<int> srcs(indegree);
175 UniqueArray<int> dsts(outdegree);
176
177 MPI_Dist_graph_neighbors(m_topology_communicator, indegree, srcs.data(), MPI_UNWEIGHTED, outdegree, dsts.data(), MPI_UNWEIGHTED);
178
179 for (int k = 0; k < outdegree; ++k) {
180 int x = dsts[k];
181 if (x != comm_ranks[k])
182 ARCANE_FATAL("Invalid destination rank order k={0} v={1} expected={2}", k, x, comm_ranks[k]);
183 }
184
185 for (int k = 0; k < indegree; ++k) {
186 int x = srcs[k];
187 if (x != comm_ranks[k])
188 ARCANE_FATAL("Invalid source rank order k={0} v={1} expected={2}", k, x, comm_ranks[k]);
189 }
190 }
191 }
192
193 private:
194
195 MpiParallelMng* m_mpi_parallel_mng = nullptr;
196 MPI_Comm m_topology_communicator = MPI_COMM_NULL;
197
198 private:
199
200 void _checkFreeCommunicator()
201 {
202 if (m_topology_communicator != MPI_COMM_NULL)
203 MPI_Comm_free(&m_topology_communicator);
204 m_topology_communicator = MPI_COMM_NULL;
205 }
206};
207
208/*---------------------------------------------------------------------------*/
209/*---------------------------------------------------------------------------*/
216class MpiVariableSynchronizer
217: public VariableSynchronizer
218{
219 public:
220
221 MpiVariableSynchronizer(IParallelMng* pm, const ItemGroup& group,
222 Ref<IDataSynchronizeImplementationFactory> implementation_factory,
224 : VariableSynchronizer(pm, group, implementation_factory)
225 , m_topology_info(topology_info)
226 {
227 }
228
229 public:
230
231 void compute() override
232 {
234 // If not null, calculate the topology
235 if (m_topology_info.get())
236 m_topology_info->compute(this);
237 }
238
239 private:
240
242};
243
244/*---------------------------------------------------------------------------*/
245/*---------------------------------------------------------------------------*/
246
247class MpiParallelMngUtilsFactory
249{
250 public:
251
252 MpiParallelMngUtilsFactory()
253 : m_synchronizer_version(2)
254 {
255 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "1")
256 m_synchronizer_version = 1;
257 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "2")
258 m_synchronizer_version = 2;
259 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "3")
260 m_synchronizer_version = 3;
261 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "4") {
262 m_synchronizer_version = 4;
263 String v = platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_BLOCK_SIZE");
264 if (!v.null()) {
265 Int32 block_size = 0;
266 if (!builtInGetValue(block_size, v))
267 m_synchronize_block_size = block_size;
268 m_synchronize_block_size = std::clamp(m_synchronize_block_size, 0, 1000000000);
269 }
270 v = platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_NB_SEQUENCE");
271 if (!v.null()) {
272 Int32 nb_sequence = 0;
273 if (!builtInGetValue(nb_sequence, v))
274 m_synchronize_nb_sequence = nb_sequence;
275 m_synchronize_nb_sequence = std::clamp(m_synchronize_nb_sequence, 1, 1024 * 1024);
276 }
277 }
278 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "5")
279 m_synchronizer_version = 5;
280 if (platform::getEnvironmentVariable("ARCANE_SYNCHRONIZE_VERSION") == "6")
281 m_synchronizer_version = 6;
282 }
283
284 public:
285
287 {
288 return _createSynchronizer(pm, family->allItems());
289 }
290
292 {
293 return _createSynchronizer(pm, group);
294 }
295
296 private:
297
298 Ref<IVariableSynchronizer> _createSynchronizer(IParallelMng* pm, const ItemGroup& group)
299 {
301 MpiParallelMng* mpi_pm = ARCANE_CHECK_POINTER(dynamic_cast<MpiParallelMng*>(pm));
302 ITraceMng* tm = pm->traceMng();
304 // Only displays information for the group of all cells to avoid displaying
305 // the same message multiple times.
306 bool do_print = (group.isAllItems() && group.itemKind() == IK_Cell);
307 if (m_synchronizer_version == 2) {
308 if (do_print)
309 tm->info() << "Using MpiSynchronizer V2";
310 generic_factory = arcaneCreateMpiVariableSynchronizerFactory(mpi_pm);
311 }
312 else if (m_synchronizer_version == 3) {
313 if (do_print)
314 tm->info() << "Using MpiSynchronizer V3";
315 generic_factory = arcaneCreateMpiDirectSendrecvVariableSynchronizerFactory(mpi_pm);
316 }
317 else if (m_synchronizer_version == 4) {
318 if (do_print)
319 tm->info() << "Using MpiSynchronizer V4 block_size=" << m_synchronize_block_size
320 << " nb_sequence=" << m_synchronize_nb_sequence;
321 generic_factory = arcaneCreateMpiBlockVariableSynchronizerFactory(mpi_pm, m_synchronize_block_size, m_synchronize_nb_sequence);
322 }
323 else if (m_synchronizer_version == 5) {
324 if (do_print)
325 tm->info() << "Using MpiSynchronizer V5";
326 topology_info = createRef<VariableSynchronizerMpiCommunicator>(mpi_pm);
327#if defined(ARCANE_HAS_MPI_NEIGHBOR)
328 generic_factory = arcaneCreateMpiNeighborVariableSynchronizerFactory(mpi_pm, topology_info);
329#else
330 throw NotSupportedException(A_FUNCINFO, "Synchronize implementation V5 is not supported with this version of MPI");
331#endif
332 }
333#if defined(ARCANE_HAS_PACKAGE_NCCL)
334 else if (m_synchronizer_version == 6) {
335 if (do_print)
336 tm->info() << "Using NCCLSynchronizer";
337 generic_factory = arcaneCreateNCCLVariableSynchronizerFactory(mpi_pm);
338 }
339#endif
340 else {
341 if (do_print)
342 tm->info() << "Using MpiSynchronizer V1";
343 generic_factory = arcaneCreateMpiLegacyVariableSynchronizerFactory(mpi_pm);
344 }
345 if (!generic_factory.get())
346 ARCANE_FATAL("No factory created");
347 return createRef<MpiVariableSynchronizer>(pm, group, generic_factory, topology_info);
348 }
349
350 private:
351
352 Integer m_synchronizer_version = 1;
353 Int32 m_synchronize_block_size = 32000;
354 Int32 m_synchronize_nb_sequence = 1;
355};
356
357/*---------------------------------------------------------------------------*/
358/*---------------------------------------------------------------------------*/
359
360/*---------------------------------------------------------------------------*/
361/*---------------------------------------------------------------------------*/
362
364: public ParallelMngInternal
365{
366 public:
367
368 explicit Impl(MpiParallelMng* pm)
369 : ParallelMngInternal(pm)
370 , m_parallel_mng(pm)
371 , m_alloc(makeRef(new MachineShMemWinMemoryAllocator(pm)))
372 {}
373
374 ~Impl() override = default;
375
376 public:
377
378 Int32 masterParallelIORank() const override { return m_parallel_mng->commRank(); }
379 Int32 nbSendersToMasterParallelIO() const override { return 1; }
380
382 {
383 m_parallel_mng->traceMng()->info() << "initializeWindowCreator() MPI";
384 m_parallel_mng->adapter()->initializeWindowCreator(m_parallel_mng->machineCommunicator());
385 }
386
388 {
389 if (m_shmem_available == 1) {
390 return true;
391 }
392
393 if (m_shmem_available == 0) {
394 Ref<IParallelTopology> topo = m_parallel_mng->_internalUtilsFactory()->createTopology(m_parallel_mng);
395 if (topo->machineRanks().size() == m_parallel_mng->adapter()->windowCreator()->machineRanks().size()) {
396 m_shmem_available = 1;
397 return true;
398 }
399 // Problem with MPI. May occur if MPICH is compiled in ch3:sock mode.
400 m_shmem_available = 2;
401 return false;
402 }
403
404 return false;
405 }
406
408 {
409 return makeRef(m_parallel_mng->adapter()->windowCreator()->createWindow(sizeof_segment, sizeof_type));
410 }
411
413 {
414 return makeRef(m_parallel_mng->adapter()->windowCreator()->createDynamicWindow(sizeof_segment, sizeof_type));
415 }
416
418 {
419 return MemoryAllocationOptions{ m_alloc.get() };
420 }
421
423 {
424 return m_parallel_mng->adapter()->windowCreator()->machineRanks();
425 }
426
427 void machineBarrier() override
428 {
429 m_parallel_mng->adapter()->windowCreator()->machineBarrier();
430 }
431
432 private:
433
434 MpiParallelMng* m_parallel_mng;
436
437 // 0 = Attribute not initialized
438 // 1 = Shared memory available
439 // 2 = Shared memory not available
440 Int8 m_shmem_available = 0;
441};
442
443/*---------------------------------------------------------------------------*/
444/*---------------------------------------------------------------------------*/
445
446/*---------------------------------------------------------------------------*/
447/*---------------------------------------------------------------------------*/
448
449MpiParallelMng::
450MpiParallelMng(const MpiParallelMngBuildInfo& bi)
451: ParallelMngDispatcher(ParallelMngDispatcherBuildInfo(bi.commRank(),bi.commSize(),MP::Communicator(bi.mpiComm())))
452, m_trace(bi.trace_mng)
453, m_thread_mng(bi.thread_mng)
454, m_world_parallel_mng(bi.world_parallel_mng)
455, m_timer_mng(bi.timer_mng)
456, m_replication(new ParallelReplication())
457, m_is_parallel(bi.is_parallel)
458, m_comm_rank(bi.commRank())
459, m_comm_size(bi.commSize())
460, m_stat(bi.stat)
461, m_communicator(bi.mpiComm())
462, m_machine_communicator(bi.mpiMachineComm())
463, m_is_communicator_owned(bi.is_mpi_comm_owned)
464, m_mpi_lock(bi.mpi_lock)
465, m_non_blocking_collective(nullptr)
466, m_utils_factory(createRef<MpiParallelMngUtilsFactory>())
467, m_parallel_mng_internal(new Impl(this))
468{
469 if (!m_world_parallel_mng) {
470 m_trace->debug() << "[MpiParallelMng] No m_world_parallel_mng found, reverting to ourselves!";
471 m_world_parallel_mng = this;
472 }
473}
474
475/*---------------------------------------------------------------------------*/
476/*---------------------------------------------------------------------------*/
477
478MpiParallelMng::
479~MpiParallelMng()
480{
481 delete m_parallel_mng_internal;
482 delete m_non_blocking_collective;
483 m_sequential_parallel_mng.reset();
484 if (m_is_communicator_owned) {
485 MpiLock::Section ls(m_mpi_lock);
486 MPI_Comm_free(&m_communicator);
487 MPI_Comm_free(&m_machine_communicator);
488 }
489 delete m_replication;
490 delete m_io_mng;
491 if (m_is_timer_owned)
492 delete m_timer_mng;
493 arcaneCallFunctionAndTerminateIfThrow([&]() { m_adapter->destroy(); });
494 delete m_datatype_list;
495}
496
497/*---------------------------------------------------------------------------*/
498/*---------------------------------------------------------------------------*/
499
500namespace
501{
502
503 /*---------------------------------------------------------------------------*/
504 /*---------------------------------------------------------------------------*/
505 // Class to create the different dispatchers
506 class DispatchCreator
507 {
508 public:
509
510 DispatchCreator(ITraceMng* tm, IMessagePassingMng* mpm, MpiAdapter* adapter, MpiDatatypeList* datatype_list)
511 : m_tm(tm)
512 , m_mpm(mpm)
513 , m_adapter(adapter)
514 , m_datatype_list(datatype_list)
515 {}
516
517 public:
518
519 template <typename DataType> MpiParallelDispatchT<DataType>*
520 create()
521 {
522 MpiDatatype* dt = m_datatype_list->datatype(DataType());
523 return new MpiParallelDispatchT<DataType>(m_tm, m_mpm, m_adapter, dt);
524 }
525
526 ITraceMng* m_tm;
527 IMessagePassingMng* m_mpm;
528 MpiAdapter* m_adapter;
529 MpiDatatypeList* m_datatype_list;
530 };
531
532 /*---------------------------------------------------------------------------*/
533 /*---------------------------------------------------------------------------*/
534
535 class ControlDispatcherDecorator
536 : public ParallelMngDispatcher::DefaultControlDispatcher
537 {
538 public:
539
540 ControlDispatcherDecorator(IParallelMng* pm, MpiAdapter* adapter)
541 : ParallelMngDispatcher::DefaultControlDispatcher(pm)
542 , m_adapter(adapter)
543 {}
544
545 IMessagePassingMng* commSplit(bool keep) override
546 {
547 return m_adapter->commSplit(keep);
548 }
549 MP::IProfiler* profiler() const override { return m_adapter->profiler(); }
550 void setProfiler(MP::IProfiler* p) override { m_adapter->setProfiler(p); }
551
552 private:
553
554 MpiAdapter* m_adapter;
555 };
556} // namespace
557
558/*---------------------------------------------------------------------------*/
559/*---------------------------------------------------------------------------*/
560
562build()
563{
564 ITraceMng* tm = traceMng();
565 if (!m_timer_mng) {
566 m_timer_mng = new MpiTimerMng(tm);
567 m_is_timer_owned = true;
568 }
569
570 // Created the associated sequential manager.
571 {
573 bi.setTraceMng(traceMng());
574 bi.setCommunicator(communicator());
575 bi.setThreadMng(threadMng());
576 m_sequential_parallel_mng = arcaneCreateSequentialParallelMngRef(bi);
577 }
578
579 // Indicates that reduces must be performed in processor order
580 // in order to guarantee deterministic execution
581 bool is_ordered_reduce = false;
582 if (platform::getEnvironmentVariable("ARCANE_ORDERED_REDUCE") == "TRUE")
583 is_ordered_reduce = true;
584 m_datatype_list = new MpiDatatypeList(is_ordered_reduce);
585
586 ARCANE_CHECK_POINTER(m_stat);
587
588 MpiAdapter* adapter = new MpiAdapter(m_trace, m_stat->toArccoreStat(), m_communicator, m_mpi_lock);
589 m_adapter = adapter;
590 auto mpm = _messagePassingMng();
591
592 // NOTE: this instance will be destroyed by the ParallelMngDispatcher
593 auto* control_dispatcher = new ControlDispatcherDecorator(this, m_adapter);
594 _setControlDispatcher(control_dispatcher);
595
596 // NOTE: this instance will be destroyed by the ParallelMngDispatcher
597 auto* serialize_dispatcher = new MpiSerializeDispatcher(m_adapter, mpm);
598 m_mpi_serialize_dispatcher = serialize_dispatcher;
599 _setSerializeDispatcher(serialize_dispatcher);
600
601 DispatchCreator creator(m_trace, mpm, m_adapter, m_datatype_list);
602 this->createDispatchers(creator);
603
604 m_io_mng = arcaneCreateIOMng(this);
605
606 m_non_blocking_collective = new MpiParallelNonBlockingCollective(tm, this, adapter);
607 m_non_blocking_collective->build();
608 if (m_mpi_lock)
609 m_trace->info() << "Using mpi with locks.";
610
611 m_parallel_mng_internal->initializeWindowCreator();
612}
613
614/*---------------------------------------------------------------------------*/
615/*---------------------------------------------------------------------------*/
616
617/*----------------------------------------------------------------------------*/
618/*---------------------------------------------------------------------------*/
619
622{
623 Trace::Setter mci(m_trace, "Mpi");
624 if (m_is_initialized) {
625 m_trace->warning() << "MpiParallelMng already initialized";
626 return;
627 }
628
629 // Initialization of MpiParallelMng
630 m_trace->info() << "Initialisation de MpiParallelMng";
631 m_sequential_parallel_mng->initialize();
632
633 m_adapter->setTimeMetricCollector(timeMetricCollector());
634
635 m_is_initialized = true;
636}
637
638/*---------------------------------------------------------------------------*/
639/*---------------------------------------------------------------------------*/
640
641void MpiParallelMng::
642sendSerializer(ISerializer* s, Int32 rank)
643{
644 Trace::Setter mci(m_trace, "Mpi");
645 Timer::Phase tphase(timeStats(), TP_Communication);
647 m_mpi_serialize_dispatcher->legacySendSerializer(s, { MessageRank(rank), mpi_tag, Blocking });
648}
649
650/*---------------------------------------------------------------------------*/
651/*---------------------------------------------------------------------------*/
652
655{
656 return m_utils_factory->createSendSerializeMessage(this, rank)._release();
657}
658
659/*---------------------------------------------------------------------------*/
660/*---------------------------------------------------------------------------*/
661
662Request MpiParallelMng::
663sendSerializer(ISerializer* s, Int32 rank, [[maybe_unused]] ByteArray& bytes)
664{
665 Trace::Setter mci(m_trace, "Mpi");
666 Timer::Phase tphase(timeStats(), TP_Communication);
668 return m_mpi_serialize_dispatcher->legacySendSerializer(s, { MessageRank(rank), mpi_tag, NonBlocking });
669}
670
671/*---------------------------------------------------------------------------*/
672/*---------------------------------------------------------------------------*/
673
674void MpiParallelMng::
675broadcastSerializer(ISerializer* values, Int32 rank)
676{
677 Timer::Phase tphase(timeStats(), TP_Communication);
678 m_mpi_serialize_dispatcher->broadcastSerializer(values, MessageRank(rank));
679}
680
681/*---------------------------------------------------------------------------*/
682/*---------------------------------------------------------------------------*/
683
684void MpiParallelMng::
685recvSerializer(ISerializer* values, Int32 rank)
686{
687 Trace::Setter mci(m_trace, "Mpi");
688 Timer::Phase tphase(timeStats(), TP_Communication);
690 m_mpi_serialize_dispatcher->legacyReceiveSerializer(values, MessageRank(rank), mpi_tag);
691}
692
693/*---------------------------------------------------------------------------*/
694/*---------------------------------------------------------------------------*/
695
698{
699 return m_utils_factory->createReceiveSerializeMessage(this, rank)._release();
700}
701
702/*---------------------------------------------------------------------------*/
703/*---------------------------------------------------------------------------*/
704
706probe(const PointToPointMessageInfo& message) -> MessageId
707{
708 return m_adapter->probeMessage(message);
709}
710
711/*---------------------------------------------------------------------------*/
712/*---------------------------------------------------------------------------*/
713
716{
717 return m_adapter->legacyProbeMessage(message);
718}
719
720/*---------------------------------------------------------------------------*/
721/*---------------------------------------------------------------------------*/
722
723Request MpiParallelMng::
724sendSerializer(const ISerializer* s, const PointToPointMessageInfo& message)
725{
726 return m_mpi_serialize_dispatcher->sendSerializer(s, message);
727}
728
729/*---------------------------------------------------------------------------*/
730/*---------------------------------------------------------------------------*/
731
732Request MpiParallelMng::
733receiveSerializer(ISerializer* s, const PointToPointMessageInfo& message)
734{
735 return m_mpi_serialize_dispatcher->receiveSerializer(s, message);
736}
737
738/*---------------------------------------------------------------------------*/
739/*---------------------------------------------------------------------------*/
740
743{
744 for (Integer i = 0, is = requests.size(); i < is; ++i)
745 m_adapter->freeRequest(requests[i]);
746}
747
748/*---------------------------------------------------------------------------*/
749/*---------------------------------------------------------------------------*/
750
751void MpiParallelMng::
752_checkFinishedSubRequests()
753{
754 m_mpi_serialize_dispatcher->checkFinishedSubRequests();
755}
756
757/*---------------------------------------------------------------------------*/
758/*---------------------------------------------------------------------------*/
759
760Ref<IParallelMng> MpiParallelMng::
761sequentialParallelMngRef()
762{
763 return m_sequential_parallel_mng;
764}
765
768{
769 return m_sequential_parallel_mng.get();
770}
771
772/*---------------------------------------------------------------------------*/
773/*---------------------------------------------------------------------------*/
774
777{
778 if (m_stat)
779 m_stat->print(m_trace);
780}
781
782/*---------------------------------------------------------------------------*/
783/*---------------------------------------------------------------------------*/
784
786barrier()
787{
788 traceMng()->flush();
789 m_adapter->barrier();
790}
791
792/*---------------------------------------------------------------------------*/
793/*---------------------------------------------------------------------------*/
794
797{
798 m_adapter->waitAllRequests(requests);
799 _checkFinishedSubRequests();
800}
801
802/*---------------------------------------------------------------------------*/
803/*---------------------------------------------------------------------------*/
804
807{
808 return _waitSomeRequests(requests, false);
809}
810
811/*---------------------------------------------------------------------------*/
812/*---------------------------------------------------------------------------*/
813
816{
817 return _waitSomeRequests(requests, true);
818}
819
820/*---------------------------------------------------------------------------*/
821/*---------------------------------------------------------------------------*/
822
823UniqueArray<Integer> MpiParallelMng::
824_waitSomeRequests(ArrayView<Request> requests, bool is_non_blocking)
825{
826 UniqueArray<Integer> results;
827 UniqueArray<bool> done_indexes(requests.size());
828
829 m_adapter->waitSomeRequests(requests, done_indexes, is_non_blocking);
830 for (int i = 0; i < requests.size(); i++) {
831 if (done_indexes[i])
832 results.add(i);
833 }
834 return results;
835}
836
837/*---------------------------------------------------------------------------*/
838/*---------------------------------------------------------------------------*/
839
840ISerializeMessageList* MpiParallelMng::
841_createSerializeMessageList()
842{
843 return new MP::internal::SerializeMessageList(messagePassingMng());
844}
845
846/*---------------------------------------------------------------------------*/
847/*---------------------------------------------------------------------------*/
848
851{
852 return m_utils_factory->createGetVariablesValuesOperation(this)._release();
853}
854
855/*---------------------------------------------------------------------------*/
856/*---------------------------------------------------------------------------*/
857
860{
861 return m_utils_factory->createTransferValuesOperation(this)._release();
862}
863
864/*---------------------------------------------------------------------------*/
865/*---------------------------------------------------------------------------*/
866
869{
870 return m_utils_factory->createExchanger(this)._release();
871}
872
873/*---------------------------------------------------------------------------*/
874/*---------------------------------------------------------------------------*/
875
878{
879 return m_utils_factory->createSynchronizer(this, family)._release();
880}
881
882/*---------------------------------------------------------------------------*/
883/*---------------------------------------------------------------------------*/
884
886createSynchronizer(const ItemGroup& group)
887{
888 return m_utils_factory->createSynchronizer(this, group)._release();
889}
890
891/*---------------------------------------------------------------------------*/
892/*---------------------------------------------------------------------------*/
893
896{
897 return m_utils_factory->createTopology(this)._release();
898}
899
900/*---------------------------------------------------------------------------*/
901/*---------------------------------------------------------------------------*/
902
904replication() const
905{
906 return m_replication;
907}
908
909/*---------------------------------------------------------------------------*/
910/*---------------------------------------------------------------------------*/
911
914{
915 delete m_replication;
916 m_replication = v;
917}
918
919/*---------------------------------------------------------------------------*/
920/*---------------------------------------------------------------------------*/
921
922IParallelMng* MpiParallelMng::
923_createSubParallelMng(MPI_Comm sub_communicator)
924{
925 // If null, this rank is not part of the sub-communicator
926 if (sub_communicator == MPI_COMM_NULL)
927 return nullptr;
928
929 int sub_rank = -1;
930 MPI_Comm_rank(sub_communicator, &sub_rank);
931
932 MPI_Comm sub_machine_communicator = MPI_COMM_NULL;
933 MPI_Comm_split_type(sub_communicator, MPI_COMM_TYPE_SHARED, sub_rank, MPI_INFO_NULL, &sub_machine_communicator);
934
935 MpiParallelMngBuildInfo bi(sub_communicator, sub_machine_communicator);
936 bi.is_parallel = isParallel();
937 bi.stat = m_stat;
938 bi.timer_mng = m_timer_mng;
939 bi.thread_mng = m_thread_mng;
940 bi.trace_mng = m_trace;
941 bi.world_parallel_mng = m_world_parallel_mng;
942 bi.mpi_lock = m_mpi_lock;
943
944 IParallelMng* sub_pm = new MpiParallelMng(bi);
945 sub_pm->build();
946 return sub_pm;
947}
948
949/*---------------------------------------------------------------------------*/
950/*---------------------------------------------------------------------------*/
951
952Ref<IParallelMng> MpiParallelMng::
953_createSubParallelMngRef(Int32 color, Int32 key)
954{
955 if (color < 0)
956 color = MPI_UNDEFINED;
957 MPI_Comm sub_communicator = MPI_COMM_NULL;
958 MPI_Comm_split(m_communicator, color, key, &sub_communicator);
959 IParallelMng* sub_pm = _createSubParallelMng(sub_communicator);
960 return makeRef(sub_pm);
961}
962
963/*---------------------------------------------------------------------------*/
964/*---------------------------------------------------------------------------*/
965
966IParallelMng* MpiParallelMng::
967_createSubParallelMng(Int32ConstArrayView kept_ranks)
968{
969 MPI_Group mpi_group = MPI_GROUP_NULL;
970 MPI_Comm_group(m_communicator, &mpi_group);
971 Integer nb_sub_rank = kept_ranks.size();
972 UniqueArray<int> mpi_kept_ranks(nb_sub_rank);
973 for (Integer i = 0; i < nb_sub_rank; ++i)
974 mpi_kept_ranks[i] = (int)kept_ranks[i];
975
976 MPI_Group final_group = MPI_GROUP_NULL;
977 MPI_Group_incl(mpi_group, nb_sub_rank, mpi_kept_ranks.data(), &final_group);
978 MPI_Comm sub_communicator = MPI_COMM_NULL;
979
980 MPI_Comm_create(m_communicator, final_group, &sub_communicator);
981 MPI_Group_free(&final_group);
982 return _createSubParallelMng(sub_communicator);
983}
984
985/*---------------------------------------------------------------------------*/
986/*---------------------------------------------------------------------------*/
987
996: public MpiRequestList
997{
998 using Base = MpiRequestList;
999
1000 public:
1001
1002 explicit RequestList(MpiParallelMng* pm)
1003 : Base(pm->m_adapter)
1004 , m_parallel_mng(pm)
1005 {}
1006
1007 public:
1008
1009 void _wait(Parallel::eWaitType wait_type) override
1010 {
1011 Base::_wait(wait_type);
1012 m_parallel_mng->_checkFinishedSubRequests();
1013 };
1014
1015 private:
1016
1017 MpiParallelMng* m_parallel_mng;
1018};
1019
1020/*---------------------------------------------------------------------------*/
1021/*---------------------------------------------------------------------------*/
1022
1028
1029/*---------------------------------------------------------------------------*/
1030/*---------------------------------------------------------------------------*/
1031
1034{
1035 return m_utils_factory;
1036}
1037
1038/*---------------------------------------------------------------------------*/
1039/*---------------------------------------------------------------------------*/
1040
1041bool MpiParallelMng::
1042_isAcceleratorAware() const
1043{
1045}
1046
1047/*---------------------------------------------------------------------------*/
1048/*---------------------------------------------------------------------------*/
1049
1050} // End namespace Arcane
1051
1052/*---------------------------------------------------------------------------*/
1053/*---------------------------------------------------------------------------*/
#define ARCANE_CHECK_POINTER(ptr)
Macro returning the pointer ptr if it is not null or throwing an exception if it is null.
#define ARCANE_FATAL(...)
Macro throwing a FatalErrorException.
Brief list of message exchange functions.
Modifiable view of an array of type T.
constexpr Integer size() const noexcept
Returns the size of the array.
void add(ConstReferenceType val)
Adds element val to the end of the array.
const T * data() const
Access to the root of the array without any protection.
Constant view of an array of type T.
constexpr Integer size() const noexcept
Number of elements in the array.
Operations to access variable values from another subdomain.
Interface of the input/output manager.
Definition IIOMng.h:37
Interface of an entity family.
Definition IItemFamily.h:85
virtual ItemGroup allItems() const =0
Group of all entities.
Information exchange between processors.
Interface of the parallelism manager for a subdomain.
virtual ITraceMng * traceMng() const =0
Trace manager.
virtual void build()=0
Constructs the instance.
Brief information on parallel subdomain replication.
Information on the computing core allocation topology.
virtual void flush()=0
Flushes all streams.
virtual TraceMessage info()=0
Stream for an information message.
Sends values across different processors.
Interface of a specific MPI communicator for synchronizations.
Interface of a variable synchronization service.
Mesh entity group.
Definition ItemGroup.h:51
eItemKind itemKind() const
Group kind. This is the kind of its elements.
Definition ItemGroup.h:114
bool isAllItems() const
Indicates if the group is that of all entities.
Definition ItemGroup.cc:607
Communicator for message exchange.
Information about the source of a message.
void _wait(eWaitType wait_type) override
Performs the wait or test.
Request receiveSerializer(ISerializer *s, const PointToPointMessageInfo &message) override
Receiving message.
Request sendSerializer(const ISerializer *s, const PointToPointMessageInfo &message) override
Sending message.
Information for sending/receiving a point-to-point message.
Serializing message using a BasicSerializer.
static MessageTag defaultTag()
Default tag for serialization messages.
Manages the MPI_Datatypes associated with Arcane types.
Ref< IVariableSynchronizer > createSynchronizer(IParallelMng *pm, IItemFamily *family) override
Returns an interface to synchronize variables on the group of the family family.
Ref< IVariableSynchronizer > createSynchronizer(IParallelMng *pm, const ItemGroup &group) override
Returns an interface to synchronize variables on the group group.
void initializeWindowCreator() override
Method allowing the initialization of the windowCreator specific to the implementation.
Ref< IContigMachineShMemWinBaseInternal > createContigMachineShMemWinBase(Int64 sizeof_segment, Int32 sizeof_type) override
Method allowing the creation of a memory window on the node.
Ref< IMachineShMemWinBaseInternal > createMachineShMemWinBase(Int64 sizeof_segment, Int32 sizeof_type) override
Method allowing the creation of a dynamic memory window on the node.
bool isMachineShMemWinAvailable() override
Method allowing to know if shared memory mode is supported.
ConstArrayView< Int32 > machineRanks() override
Method allowing retrieval of the ranks of the sub-domains of the computing node.
void machineBarrier() override
Method allowing a barrier for the sub-domains of the computing node.
Int32 masterParallelIORank() const override
Int32 nbSendersToMasterParallelIO() const override
MemoryAllocationOptions machineShMemWinMemoryAllocator() override
Method allowing retrieval of a shared memory allocator.
Specialization of MpiRequestList for MpiParallelMng.
void _wait(Parallel::eWaitType wait_type) override
Performs the wait or test.
Parallelism manager using MPI.
IParallelMng * worldParallelMng() const override
Parallelism manager over all allocated resources.
MessageSourceInfo legacyProbe(const PointToPointMessageInfo &message) override
Probes if messages are available.
void barrier() override
Performs a barrier.
UniqueArray< Integer > waitSomeRequests(ArrayView< Request > requests) override
Blocks while waiting for one of the rvalues requests to complete.
void waitAllRequests(ArrayView< Request > requests) override
Blocks while waiting for the rvalues requests to complete.
MessageId probe(const PointToPointMessageInfo &message) override
Probes if messages are available.
void build() override
Constructs the instance.
void printStats() override
Prints statistics related to this parallelism manager.
bool m_is_initialized
true if already initialized
IParallelMng * sequentialParallelMng() override
Returns a sequential parallelism manager.
IThreadMng * threadMng() const override
Thread manager.
ITimerMng * timerMng() const override
Timer manager.
void initialize() override
Initializes the parallelism manager.
IVariableSynchronizer * createSynchronizer(IItemFamily *family) override
Returns an interface for synchronizing variables on the group of the family.
ISerializeMessage * createSendSerializer(Int32 rank) override
Creates a non-blocking message to send serialized data to rank rank.
bool isParallel() const override
Returns true if the execution is parallel.
Ref< IParallelMngUtilsFactory > _internalUtilsFactory() const override
Factory for utility functions.
ITraceMng * traceMng() const override
Trace manager.
IParallelTopology * createTopology() override
Creates an instance containing information about the rank topology of this manager.
Communicator communicator() const override
MPI communicator associated with this instance.
IParallelExchanger * createExchanger() override
Returns an interface for transferring messages between processors.
IParallelReplication * replication() const override
Replication information.
void setReplication(IParallelReplication *v) override
Sets the Replication Information.
Ref< Parallel::IRequestList > createRequestListRef() override
Creates a request list for this manager.
ITransferValuesParallelOperation * createTransferValuesOperation() override
Returns an operation to transfer values between subdomains.
UniqueArray< Integer > testSomeRequests(ArrayView< Request > requests) override
Tests if one of the rvalues requests is complete.
IGetVariablesValuesParallelOperation * createGetVariablesValuesOperation() override
Returns an operation to retrieve the values of a variable on the entities of another subdomain.
ISerializeMessage * createReceiveSerializer(Int32 rank) override
Creates a non-blocking message to receive serialized data from rank rank.
void freeRequests(ArrayView< Parallel::Request > requests) override
Frees the requests.
Timer manager using the MPI library.
Definition MpiTimerMng.h:41
void compute() override
Recalculates the synchronization information.
Redirects the message management of sub-domains according to the argument type.
IMessagePassingMng * messagePassingMng() const override
Associated Arccore message passing manager.
ITimeMetricCollector * timeMetricCollector() const override
Interface for collecting execution times (can be null).
ITimeStats * timeStats() const override
Associated statistics manager (can be null).
Base class of a factory for IParallelMng utility functions.
Brief information on parallel subdomain replication.
InstanceType * get() const
Associated instance or nullptr if none.
Reference to an instance.
bool null() const
Returns true if the string is null.
Definition String.cc:306
Positions the phase of the currently executing action.
Definition Timer.h:142
1D data vector with value semantics (STL style).
MPI_Comm communicator() const override
Retrieves the specific communicator from the topology.
void compute(VariableSynchronizer *var_syncer) override
Calculates the specific communicator.
Interface of a variable synchronization service.
void compute() override
Creation of the list of synchronization elements.
Int32ConstArrayView communicatingRanks() override
Ranks of subdomains with which communication occurs.
Declarations of types and methods used by message exchange mechanisms.
String getEnvironmentVariable(const String &name)
Environment variable named name.
-- tab-width: 2; indent-tabs-mode: nil; coding: utf-8-with-signature --
std::int8_t Int8
Signed integer type of 8 bits.
Ref< TrueType > createRef(Args &&... args)
Creates an instance of type TrueType with arguments Args and returns a reference to it.
ARCANE_MPI_EXPORT bool arcaneIsAcceleratorAwareMPI()
Indicates if the current MPI runtime supports accelerators.
Definition ArcaneMpi.cc:85
std::int64_t Int64
Signed integer type of 64 bits.
Int32 Integer
Type representing an integer.
Array< Byte > ByteArray
Dynamic one-dimensional array of characters.
Definition UtilsTypes.h:115
ConstArrayView< Int32 > Int32ConstArrayView
C equivalent of a 1D array of 32-bit integers.
Definition UtilsTypes.h:476
@ IK_Cell
Cell mesh entity.
void arcaneCallFunctionAndTerminateIfThrow(std::function< void()> function)
Calls the function function and calls std::terminate() if an exception occurs.
auto makeRef(InstanceType *t) -> Ref< InstanceType >
Creates a reference on a pointer.
std::int32_t Int32
Signed integer type of 32 bits.
Info to construct an MpiParallelMng.
Information to construct a SequentialParallelMng.