Deterministically forwarding results of endpoint processing unit computations through a network
Some embodiments provide a method of executing a distributed application with multiple EPUs that perform computations for the application. The EPUs are connected through a network having network elements. The method iteratively provides instructions to the EPUs to perform computations associated with the application. The method configures each EPU to store results of the computations performed by the EPU in a memory connected to the network. Computation results of at least two different EPUs are stored in at least two different memories connected to the network. The method configures, for each computation result in memory, at least one network element associated with the EPU to determine, based on a set of metrics that are defined to avoid creating congestion in the network, whether the result should be forwarded through the network and to forward the result through the network only after determining that the result should be forwarded.
1 . A method of executing a distributed application with a plurality of endpoint processing units (EPUs) that perform computations for the distributed application, the EPUs connected through a network comprising a plurality of network elements, each EPU having an associated network interface that connects the EPU to the network, the method comprising:
iteratively providing instructions to the EPUs to perform computations associated with the distributed application;
configuring each EPU to store results of the computations performed by the EPU in a memory connected to the network, wherein computation results of at least two different EPUs are stored in at least two different memories connected to the network; and
configuring, for each computation result stored by each EPU in memory, the network interface associated with the EPU to determine whether the result should be forwarded through the network and to forward the result through the network only after determining that the result should be forwarded through the network, wherein for a particular result computed by a particular EPU, said determination is performed based on a set of scheduling parameters that were proactively provided to the network interface for the particular result to preschedule the forwarding of the particular result before the particular EPU computed the result, said determinations performed for the EPUs to avoid creation of congestion in the network.
2 . The method of claim 1 , wherein the EPUs are graphics processing units (GPUs).
3 . The method of claim 1 , wherein the EPUs comprise at least one of graphics processing units (GPUs), tensor processing units (TPUs) and central processing units (CPUs).
4 . A method of executing a distributed application with a plurality of endpoint processing units (EPUs) that perform computations for the distributed application, the EPUs connected through a network comprising a plurality of network elements, the method comprising:
iteratively providing instructions to the EPUs to perform computations associated with the distributed application;
configuring each EPU to store results of the computations performed by the EPU in a memory connected to the network, wherein computation results of at least two different EPUs are stored in at least two different memories connected to the network; and
configuring, for each computation result stored in memory, at least one network element associated with the EPU to determine whether the result should be forwarded through the network and to forward the result through the network only after determining that the result should be forwarded through the network, said determination performed based on a set of metrics that are defined to avoid creating congestion in the network due to concurrent forwarding of more than a threshold level of a plurality of EPU computations results through overlapping paths in the network.
5 . The method of claim 1 , wherein the memories are EPU memories and the determination avoids creating congestion in the network by buffering computation results at source EPUs that perform the computations until there is enough capacity in the network to send the results to destination EPUs.
6 . The method of claim 1 , wherein the memories are EPU memories and the determination avoids creating congestion in the network by buffering computation results at source EPUs that perform the computations.
7 . The method of claim 1 , wherein the memories are memories of network interfaces of EPUs that connect the EPUs to the network and the determination avoids creating congestion in the network by buffering computation results at the network interface of source EPUs that perform the computations until there is enough capacity in the network to send the results to destination EPUs.
8 . The method of claim 1 , wherein the memories are memories of forwarding elements that form the network connecting the EPUs and the determination avoids creating congestion in the network by buffering computation results at one forwarding element until there is enough capacity in the rest of the network to send the results to destination EPUs.
9 . The method of claim 1 , wherein the set of scheduling parameters is provided to the particular EPU's network interface by a first-hop forwarding element that directly connects through a physical link to the particular EPU's network interface.
10 . A method of executing a distributed application with a plurality of endpoint processing units (EPUs) that perform computations for the distributed application, the EPUs connected through a network comprising a plurality of network elements, the method comprising:
iteratively providing instructions to the EPUs to perform computations associated with the distributed application;
configuring each EPU to store results of the computations performed by the EPU in a memory connected to the network, wherein computation results of at least two different EPUs are stored in at least two different memories connected to the network; and
configuring, for each computation result stored in memory, at least one network element associated with the EPU to determine whether the result should be forwarded through the network and to forward the result through the network only after determining that the result should be forwarded through the network,
wherein said determining, for a particular result computed by a particular EPU, comprises sending an in-band data message to a destination of the particular result to request a set of scheduling parameters to use to perform the determination and after receiving the set of scheduling parameters in response to the request, using the received set to perform the determination, said determination performed to avoid creation of congestion in the network.
11 . The method of claim 10 , wherein the set of scheduling parameters is provided by a last-hop forwarding element in a path from the particular EPU to the destination.
12 . The method of claim 11 , wherein a control plane circuit of the last-hop forwarding element provides the set of scheduling parameters.
13 . The method of claim 11 , wherein a control plane proxy server contacted by the last-hop forwarding element provides the set of scheduling parameters.
14 . A non-transitory machine readable medium storing program that when executed by at least one processor configures a network to forward communication between endpoint processing units (EPUs) that collectively execute a distributed application by performing computations associated with operations of the distributed application, the network comprising a plurality of network elements, the program comprising sets of instructions for:
for instructions that are provided to the EPUs to perform computations associated with the distributed application:
configuring each EPU to store results of the computations performed by the EPU in a memory connected to the network, wherein computation results of at least two different EPUs are stored in at least two different memories connected to the network; and
configuring, for each computation result stored in memory, at least one network element associated with the EPU to determine whether the result should be forwarded through the network and to forward the result through the network only after determining that the result should be forwarded through the network, wherein for a particular result computed by a particular EPU, said determination is performed based on a set of scheduling parameters that were proactively provided to the network element for the particular result to preschedule the forwarding of the particular result before the particular EPU computed the result, said determinations performed for the EPUs to avoid creation of congestion in the network.
15 . The non-transitory machine readable medium of claim 14 , wherein the EPUs are graphics processing units (GPUs).
16 . The non-transitory machine readable medium of claim 14 , wherein the EPUs comprise at least one of graphics processing units (GPUs), tensor processing units (TPUs) and central processing units (CPUs).
17 . The non-transitory machine readable medium of claim 14 , wherein the memories are EPU memories, and the determination avoids creating congestion in the network by buffering computation results at source EPUs that perform the computations until there is enough capacity in the network to send the results to destination EPUs.
18 . The method of claim 10 , wherein each EPU has an associated network interface that connects the EPU to the network, and the determination for a result computed by the EPU is performed by the network interface of the EPU.