Enumerations | |
| enum class | FanoutPolicy : std::uint8_t { BOUNDED , UNBOUNDED } |
| Fanout policy controlling how messages are propagated. More... | |
Functions | |
| Actor | allgather (std::shared_ptr< Context > ctx, std::shared_ptr< Communicator > comm, std::shared_ptr< Channel > ch_in, std::shared_ptr< Channel > ch_out, OpID op_id, AllGather::Ordered ordered=AllGather::Ordered::YES) |
| Create an allgather actor for a single allgather operation. More... | |
| Actor | shuffler (std::shared_ptr< Context > ctx, std::shared_ptr< Communicator > comm, std::shared_ptr< Channel > ch_in, std::shared_ptr< Channel > ch_out, OpID op_id, shuffler::PartID total_num_partitions, shuffler::Shuffler::PartitionOwner partition_owner=shuffler::Shuffler::round_robin) |
| Launches a shuffler actor for a single shuffle operation. More... | |
| Actor | fanout (std::shared_ptr< Context > ctx, std::shared_ptr< Channel > ch_in, std::vector< std::shared_ptr< Channel >> chs_out, FanoutPolicy policy) |
| Broadcast messages from one input channel to multiple output channels. More... | |
| Actor | push_to_channel (std::shared_ptr< Context > ctx, std::shared_ptr< Channel > ch_out, std::vector< Message > messages) |
| Asynchronously pushes all messages from a vector into an output channel. More... | |
| Actor | pull_from_channel (std::shared_ptr< Context > ctx, std::shared_ptr< Channel > ch_in, std::vector< Message > &out_messages) |
| Asynchronously pulls all messages from an input channel into a vector. More... | |
SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. SPDX-License-Identifier: Apache-2.0
|
strong |
Fanout policy controlling how messages are propagated.
Definition at line 17 of file fanout.hpp.
| Actor rapidsmpf::streaming::actor::allgather | ( | std::shared_ptr< Context > | ctx, |
| std::shared_ptr< Communicator > | comm, | ||
| std::shared_ptr< Channel > | ch_in, | ||
| std::shared_ptr< Channel > | ch_out, | ||
| OpID | op_id, | ||
| AllGather::Ordered | ordered = AllGather::Ordered::YES |
||
| ) |
Create an allgather actor for a single allgather operation.
This is a streaming version of rapidsmpf::coll::AllGather that operates on packed data received through Channels.
| ctx | The streaming context to use. |
| comm | Communicator for the collective operation. |
| ch_in | Input channel providing PackedDatas to be gathered. |
| ch_out | Output channel where the gathered PackedDatas are sent. |
| op_id | Unique identifier for the operation. |
| ordered | If the extracted data should be sent to the output channel with sequence numbers corresponding to the global total order of input chunks. If yes, then the sequence numbers of the extracted data will be ordered first by rank and then by input sequence number. If no, the sequence number of the extracted chunks will have no relation to any input sequence order. |
| Actor rapidsmpf::streaming::actor::fanout | ( | std::shared_ptr< Context > | ctx, |
| std::shared_ptr< Channel > | ch_in, | ||
| std::vector< std::shared_ptr< Channel >> | chs_out, | ||
| FanoutPolicy | policy | ||
| ) |
Broadcast messages from one input channel to multiple output channels.
The actor continuously receives messages from the input channel and forwards them to all output channels according to the selected fanout policy, see FanoutPolicy.
Each output channel receives a deep copy of the same message.
| ctx | The actor context to use. |
| ch_in | Input channel from which messages are received. |
| chs_out | Output channels to which messages are broadcast. Must be at least 2. |
| policy | The fanout strategy to use (see FanoutPolicy). |
| std::invalid_argument | If an unknown fanout policy is specified or if the number of output channels is less than 2. |
| Actor rapidsmpf::streaming::actor::pull_from_channel | ( | std::shared_ptr< Context > | ctx, |
| std::shared_ptr< Channel > | ch_in, | ||
| std::vector< Message > & | out_messages | ||
| ) |
Asynchronously pulls all messages from an input channel into a vector.
Receives messages from the channel until it is closed and appends them to the provided output vector.
| ctx | The actor context to use. |
| ch_in | Input channel providing messages. |
| out_messages | Output vector to store the received messages. |
| Actor rapidsmpf::streaming::actor::push_to_channel | ( | std::shared_ptr< Context > | ctx, |
| std::shared_ptr< Channel > | ch_out, | ||
| std::vector< Message > | messages | ||
| ) |
Asynchronously pushes all messages from a vector into an output channel.
Sends each message of the input vector into the channel in order, marking the end of the stream once done.
| ctx | The actor context to use. |
| ch_out | Output channel to which messages will be sent. |
| messages | Input vector containing the messages to send. |
| std::invalid_argument | if any of the elements in messages is empty. |
| Actor rapidsmpf::streaming::actor::shuffler | ( | std::shared_ptr< Context > | ctx, |
| std::shared_ptr< Communicator > | comm, | ||
| std::shared_ptr< Channel > | ch_in, | ||
| std::shared_ptr< Channel > | ch_out, | ||
| OpID | op_id, | ||
| shuffler::PartID | total_num_partitions, | ||
| shuffler::Shuffler::PartitionOwner | partition_owner = shuffler::Shuffler::round_robin |
||
| ) |
Launches a shuffler actor for a single shuffle operation.
This is a streaming version of rapidsmpf::shuffler::Shuffler that operates on packed partition chunks using channels.
It consumes partitioned input data from the input channel and produces output chunks grouped by partition_owner.
| ctx | The context to use. |
| comm | Communicator for the collective operation. |
| ch_in | Input channel providing PartitionMapChunk to be shuffled. |
| ch_out | Output channel where the resulting PartitionVectorChunks are sent. |
| op_id | Unique operation ID for this shuffle. Must not be reused until all actors have called Shuffler::shutdown(). |
| total_num_partitions | Total number of partitions to shuffle the data into. |
| partition_owner | Function that maps a partition ID to its owning rank/node. |