Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 60 additions & 0 deletions applications/anatomy.rst
Original file line number Diff line number Diff line change
Expand Up @@ -368,3 +368,63 @@ Today, the Internet and the cloud have a symbiotic relationship. The
Internet provides the communication substrate that the cloud runs on,
while the cloud provides the computing substrate that enables ever more
powerful applications to be distributed across the Internet.

2.1.4 Other Programming Models
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~

So far we have looked at how scale a network applications, all the
while remaining consistent with the client-server paradigm and using
the Socket API. There have been extensions to the API, for example to
avoid copying messages into and out of user space, but the programming
model has remained remarkably stable: one side does a passive open, the
other does an active open, both sides send and receive messages, and
the receive operation blocks waiting for a message to arrive.

This programming model matches the needs of so-called "network
applications" (which we now sometimes call "cloud apps" or simply
"apps"), but that doesn't preclude other types of computations that
need to communicate over a network. The most notable example are
parallel programs that once ran on purpose-built supercomputers, but
today run in datacenters. And of these parallel programs, AI training
programs have become the dominant use case. These programs do not use
the Socket API and they are not best described as client-server.

The full complexity of AI workloads is beyond the scope of this book,
but to appreciate the networking implications it's enough to
understand two key points about the high-level communication pattern.
First, AI training proceeds in a sequence of iterations. For each
iteration, data is first "scattered" across multiple nodes, the nodes
then compute on the subset of data sent to each of them, and finally
the results are "gathered" back in a central node. To support this
(and similar patterns), the API supports *Scatter* and *Gather*
operations among a collective of nodes, the first implies a
one-to-many communication, and the second implies a many-to-one
communication.

Importantly, there is often a synchronization barrier between each
iteration, such that one iteration has to complete before the next
iteration can begin. This means that the last transfer to complete
during each transfer limits how fast the overall computation runs. In
other words, it's not how fast the fastest communication can be
implemented; it's how fast the slowest communication completes. Being
able to achieve low latency in the face of traffic bursts—a natural
consequence of one-to-many and many-to-one exchanges—is the central
performance challenge for this communication pattern.

The second point is that parallel programs are typically structured to
maximize parallelism (the number of threads running concurrently) and
minimize blocking (waiting for a message to arrive or another thread
to complete a task). To this end, communication is typically organized
around *work queues* and *completion queues*. The idea is that threads
asynchronously insert tasks (messages) into a work queue and are later
asynchronously informed that the work is complete. Threads that must
know that work has completed before they can safely proceed typically
*poll* the completion queue rather than invoke a blocking operation.

All of this is to say that networks enable a wide range of
applications, not just those that immediately come to mind when you
think of email, web surfing, or video. The API is the demarcation
point between these applications and the network, and while
understanding the Socket API takes us a long way, there are other
models. We return the API that AI training programs use in Chapter
|Message|.
42 changes: 8 additions & 34 deletions message/rdma.rst
Original file line number Diff line number Diff line change
Expand Up @@ -28,9 +28,10 @@ network-connected machines, as shown in :numref:`Figure %s
application processes to access (read and write) blocks of data from
each others' memory. Functionally, RDMA can be viewed as a special
case of an RPC, where only two "remote procedures" are supported:
*Read( )* and *Write( )*. That perspective glosses over a lot of
details, but it is helpful to see the similarities between the two
abstractions.
*Read( )* and *Write( )*. But that observation doesn't do justice to
how RDMA fundamentally changes the way applications interact with the
network. Section |Apps|.1.4 introduced this alternative programming
model; this section takes a closer look at the details.

.. _fig-rdma:
.. figure:: message/figures/rdma.png
Expand Down Expand Up @@ -115,33 +116,6 @@ running on GPUs. NCCL was created by NVIDIA (it is an acronym for
software package running on top of the Verbs API is available as open
source.

.. sidebar:: Understanding AI Workloads

*Our focus is on the network's role supporting AI workloads, and
not on the broader programming models, nor the role of GPUs in
executing those programs. In this context, the high-level
communication pattern for AI training models is to proceed in a
sequence of iterations. For each iteration, data is first
"scattered" across multiple nodes, the nodes then compute on the
subset of data sent to each of them, and finally the results are
"gathered" back in a central node. To support this (and similar
patterns), the NCCL API supports* **Scatter** and **Gather**
*operations among a collective of nodes, the first implies a
one-to-many communication, and the second implies a many-to-one
communication.*

*Importantly, there is often a synchronization barrier between each
iteration, such that one iteration has to complete before the next
iteration can begin. This means that the last flow to complete
during each transfer limits how fast the overall computation
runs. In other words, it's not how fast the fastest communication
can be implemented; it's how fast the slowest communication
completes. Being able to achieve low latency in the face of traffic
bursts—a natural consequence of one-to-many and many-to-one
exchanges—is the central performance challenge for the network. We
return to this topic in the next section when we look at possible
optimizations.*

Second, Ethernet continues to evolve, and in this particular
circumstance, offers an alternative to InfiniBand's "native" switches.
This effort is known as *Converged Ethernet (CE)*, and it makes it
Expand Down Expand Up @@ -285,9 +259,9 @@ message is ready to be sent.
might be scattered across multiple non-continuous memory
buffers. (Our particular "Hello World" message is located in a
single continuous buffer.) This is similar to the scatter/gather
operations mentioned in the earlier sidebar about AI workloads, but
in that case work is "scattered" across multiple nodes, and in this
case a message is "scattered" across multiple buffers on a single node.
operations mentioned in Section |Apps|.1.4, but in that
case work is "scattered" across multiple nodes, and in this case a
message is "scattered" across multiple buffers on a single node.

What might be equally surprising about this example is what little we
know when the ``ibv_wr_complete`` returns, which is only that the
Expand All @@ -296,7 +270,7 @@ the message has arrived at the remote server, or that it's been
successfully written into the remote buffer. Moreover, since this is a
one-way operation, the application on the remote server is not
explicitly notified when the write does in fact complete. The
application would need to either iteratively poll the target memory
application needs to either iteratively poll the target memory
location to see if it changes, or fall back to some other out-of-band
signal (of which RDMA provides several alternatives). The point is
that "message transfer" and "process synchronization" are separable
Expand Down