I remember listening to Jane Street’s Ron Minsky on their podcast talking about this a few months ago, how TCP becomes the bottleneck in AI clusters. As an electrical engineer, I remember that circuit switching gave away to packet switching due to very sparse usage of the network when there are many actors going through it. It is not very efficient, but that wide variety of traffic makes it hard to optimize it since the flow patterns are too dynamic. A good analogy with car traffic is that downtown there are so many cars going to a large variety of places, that traffic lights – as inefficient as they are – are a solution that at least works good enough.
On the other hand, if the traffic follows a very predictable pattern, a custom implementation can be much more efficient. Specially nowadays machine learning can find much better solutions through reinforcement learning. And AI cluster data flow is much more predictable than what goes over the internet as a whole.
Google TPUs use optically routable fibers like circuit switching to be able to change the whole topology of the setup depending on training load/matrix dimensions.
I remember reading about their setup and it blew my mind. They have MEMS mirrors that can rotate to redirect one of thousands of fiber lines into any other fiber line by using other optics. That way you can take in like 1000 lines and optically connect it to any of the other 1000 lines.
Basically they use lens array (one lens per input fiber) to focus each specific beam on its assigned mirror and mirror directs the light to the chosen output that also each has a lens in front of the fiber going out.
I wonder if, as LLMs get accustomed to using homa or some other optimized protocol for, as the article lists, "chores such as weight gradients, model weights, KV cache entries, and checkpoints", whether we'll start to see TCP as a bottleneck for their post-trained interactions as well, for many of the reasons.
Homa has been around for a while. Here's the 2018 paper.[1]
The core idea: When a message arrives at the sender’s transport module, Homa
divides the message into two parts: an initial unscheduled portion (the first RTTbytes bytes), followed by a scheduled portion.
The sender transmits the unscheduled bytes immediately, using
one or more DATA packets. The scheduled bytes are not transmitted until requested explicitly by the receiver using GRANT packets.
So it sends blind for short requests, then needs a go-ahead from the receiver.
That's reasonable when the main application is a remote procedure call. It's reminiscent of QNX's networking protocol, which is also single packet message request/response but can also handle arbitrarily long messages.
What makes this work today is that per-packet processing overhead in hardware switches is low vs. per-byte overhead. In early software driven switches, per-packet overhead tended to dominate, and sending small packets was very inefficient. In modern hardware switches, where FPGAs are doing the processing, the per-packet overhead is low enough that small packets are not inefficient.
It's amusing that web stuff is so bloated today that any transaction under 1MB is considered "small".
So this is not a suitable protocol for open web use.
With regards to (~6:01) "congestion control is the responsibility of the sender" and "somehow we have to get the sending nodes to stop sending so fast". Does that not exist in Ethernet/RoCE?
> Link Level Flow Control: InfiniBand uses a credit-based algorithm to guarantee lossless HCA-to-HCA communication. RoCE runs on top of Ethernet. Implementations may require lossless Ethernet network for reaching to performance characteristics similar to InfiniBand. Lossless Ethernet is typically configured via Ethernet flow control or priority flow control (PFC). Configuring a Data center bridging (DCB) Ethernet network can be more complex than configuring an InfiniBand network.[19]
> A sending station (computer or network switch) may be transmitting data faster than the other end of the link can accept it. Using flow control, the receiving station can signal the sender requesting suspension of transmissions until the receiver catches up. Flow control on Ethernet can be implemented at the data link layer.
Flow control is better than nothing but it can cause congestion spreading and bufferbloat. QCN, Falcon, and Ultra Ethernet provide much better congestion control for RoCE but they also require newer hardware compared to Homa.
Ethernet flow control doesn't really fix congestion except in very special cases.
Consider a very simple topology:
A C
\ /
S1===S2
/ \
B D
Say hosts A and B are both sending data to C, as fast as they can, via switches S1 and S2 (which are connected via a high-speed link). And say the sum of these two flows is more than the capacity of the link to C.
S2 is receiving packets destined for C faster than it can forward them, but sending an Ethernet pause frame from S2 to S1 is not a very productive way to alleviate the situation, because it also disrupts any traffic that would be bound for D. It just moves the bottleneck elsewhere and causes collateral damage.
Maybe, but doesn’t RoCE / RDMA handle this at a higher level as well? It’s why PFC and ECN are required for these deployments, which effectively ensure everyone behaves nicely.
I agree that Ethernet flow control is insufficient, but given that NVidia in an act of brilliant foresight acquired Mellanox a decade ago, I’m fairly certain that this is how all these AI clusters are actually deployed, not using TCP, and maybe not even using Ethernet but infiniband instead.
if you spend the zillion dollars on switches and racks and power and cables and redundant power in your datacenter for infiniband(the only vendor: Nvidia), you get IB.
... which brings the question, "how is X solved?".
For IB the answer is generally a centralized controller, a lossless flow control algorithm and the credit system with subnet coordination. Ie. the packets that would be dropped due to buffer problems never get sent because the originating node knows it's out of credits. This means packet drops should be nonexistent and so the protocols make the assumption only one packet in a zillion ever gets dropped. On the down side: the tiniest of problems on IB networks have incredible impact (but it's realistic to just not have those)
IB doesn't interface with others. Yes, there's an ethernet bridge but it's not going to work nearly as well once you insert that bit of kit.
In ethernet packets are dropped when congestion is experienced, and you configure through priority queues whose packets get dropped first (you configure priority), and "stuff will sort itself out" like on the internet. RoCE tries to construct IB semantics from this, but there's limits. On the plus side: this is how the internet works and so the whole actually works somewhat reasonably over internet or other long-distance links (and tunnels into Neoclouds and/or Clouds). On the downside it's never going to match IB in speed (it's still going to get to 90% or so link speed, 98% if you just naively check the link utilization and ignore the effect retransmissions will have)
> if you spend the zillion dollars on switches and racks and power and cables and redundant power in your datacenter for infiniband(the only vendor: Nvidia), you get IB.
Worth noting that Nvidia/Mellanox is currently the only vendor for IB, but that has not always been the case. I had Intel IB switches for the backend network of my Isilon at my last job (we starting using them when they were 'just' Isilon, then EMC Isilon, and now Dell-EMC Isilon); only post-Dell did they start going Ethernet-only, and now with AI/ML they're re-offering IB as an option (LOL).
So Nvidia is kind of being 'rewarded' for sticking with a technology that everyone gave up on (I guess HPC wasn't big enough to bother with).
Besides IB, the other options being used in Top 500 are Slingshot (Eth-base?) and Omni-path (Category: Interconnect Family):
1. No way to detect whole RPC loss. Since there is no outer connection state, if every packet in the send-side of a RPC is lost then there is no way for a server to detect that it should issue a resend. The RPC is just lost to the ether. This affects small messages, like messages that fit in a single packet, more since there are fewer packets in the send-side.
2. Related to the above, there is no builtin encryption support. So, if you want encryption then you need to layer it either above or below.
3. Benchmarked performance is awful. The 60 kB average message case in [2] Table 4 takes 5(!) hyperthreads to average 20 Gbit/s. That is just 4 Gbit/s per hyperthread. Even a totally naive one-packet per system call network protocol design and implementation should get to ~8 Gbit/s per hyperthread. 30 Gbit/s per hyperthread is easy with just a little focus on performance.
4. Despite all the performance design problems in QUIC (though still faster than Homa) it already solves basically every problem Homa is trying to solve in a much cleaner way. Stream IDs correspond to RPC IDs. Stream Max corresponds to Grants. Multiple streams under single Client allows prioritization.
Except you do not randomly lose entire messages. You can compact small messages into packets. You get more precise RTT time allowing more accurate pacing/congestion calculations. You get builtin encryption. It survives ossified middleboxs. It has multiple ack frames/packets reducing ack overhead.
The only real difference is that Homa uses explicit receiver Resend instead of implicit sender Resend. Except that actually consumes significantly more receiver resources in non-trivial loss scenarios especially due to Resend packets only supporting a single Resend span. It also incurs higher latency and has higher requirements on the entire lossy network path due to all the extra transiting data that you do not want to lose.
And those are just serious problems off the top of my head after reloading the RFC into my head. I can come up with some more if needed.
Compared to TCP or Homa at a protocol design level? Not really.
However, software QUIC implementations are generally much slower than software TCP implementations. This is not due to any fundamental protocol design limitations, it is just that most/every QUIC implementation is poorly implemented for maximizing performance.
You can, of course, do multiple times more throughput than QUIC or TCP with a protocol better designed for performance, but that is likely orthogonal to your question.
Third party benchmarking [1] shows MsQuic solidly outperforming quiche by Cloudflare.
The official MsQuic benchmarks [2] tuned for the standard QUIC benchmark [3] with massively over-provisioned hardware does not even reach 8 Gbit/s on large sends with negligible RTT and loss.
If the highest performing public implementation can not even reach a paltry 8 Gbit/s it is fair to say they are all poorly performing.
As for what I would consider “okay” that would be around 30 Gbit/s per core, bottlenecking on decryption (assuming you do not use the allowed null encryption for benchmarking), or bottlenecking on ~1/3 of the bulk 1 Kbyte memcpy() throughput on your hardware, whatever comes first. “Good” is on the order of 100 Gbit/s per core or the same external limiting factors.
1. I do not see where they discuss that in the RFC. While that would work, it sounds awful.
That is a sender-driven resend and as such relies on the sender identifying the “send lost” condition. However, the entire protocol is designed around not doing that and thus you can only feasibly rely on “should have received a reply by now” as your timeout.
But Homa is intended to be a RPC protocol. So you send to server, wait for server to process command, then wait for server to send reply. Your timeout depends on the variable and heterogeneous command processing time.
Even if you were able to give a separate correctly tuned timeout for every possible RPC that is still awful. Any RPC with long processing time should not trigger the timeout until the expected reply time, but if that is far larger than the RTT then you are waiting a tremendous amount of time.
For instance, a select query that is only a few bytes long (and thus fits in one packet) could take seconds on a large database even though the database server is physically nearby and only microseconds away. In that case you would have to wait seconds before timing out instead of just microseconds like a ack-based design could achieve. The worst-case transport latency becomes application-specific instead of related to the application-independent transport parameters like RTT. That is troubling.
2. TCP can be forgiven for that given it predates asymmetric cryptography and even DES. Not considering it relevant in a new protocol this side of the millennium is much less reasonable.
There is nothing in the design that supports low latency relative to a modern protocol like QUIC. Basically the entire “latency” advantage it has over TCP is just that it does not head-of-line block which QUIC also does not do.
And then in basically every other respect it is worse including compatibility.
So that any sensitive data stupidly uploaded doesn't get leaked, which includes "zero-retention" corporate stuff. Anything that could get sniffed by employees and other agents. It really should be standard practice for any data passing between nodes, even if it's all within a site.
Homa was proposed as a datacenter-centric protocol for intra-DC traffic. In that context, with heterogeneous compute loads with various levels of criticality and confidentiality requirements, not considering how to layer encryption better is a design flaw.
A compelling architectural look at moving beyond traditional TCP to optimize latency and throughput for modern distributed AI workloads. Low-level networking is becoming the ultimate bottleneck.
One underlying claim here is that latency within the data center will start to matter more and more for future AI compute.
That is relevant for musks in-space computation. There, latency between compute nodes, due to being hundreds of kilometres apart, cannot be low.
This pretty much limits the size of AI models which can be trained on a space datacenter to a single satellite to avoid incurring huge latency overhead both in training and inference.
With the smartest models getting bigger and bigger, that really doesn't look good for space based compute.
MRC is an extension of RoCE that does out-of-order packet delivery, packet spraying across all available paths using SRv6 routing, combination of ECN and packet trimming for congestion control (so sender-based CC like TCP and unlike Homa).
One of the Packet Pushers podcasts just had an interview on this:
> Today's show puts the “heavy” in Heavy Networking with a discussion about Multipath Reliable Connection (MRC). MRC is about getting an even distribution of RDMA traffic across equal cost multipath links in massive data centers, and doing so over plain old lossy Ethernet. Ethan, Drew, and guests discuss how package spraying, resilience to link and fabric failures, and congestion control allow MRC to be a specialized transport for AI workloads.
> Our guests today are from Broadcom: Rip Sohan, Distinguished Engineer; Eric Davis, Master Engineer; and Eric Spada, Technical Director and Distinguished Engineer.
almost_usual | a day ago
https://www.usenix.org/system/files/atc21-ousterhout.pdf
dang | a day ago
swyx | 18 hours ago
giovannibonetti | a day ago
On the other hand, if the traffic follows a very predictable pattern, a custom implementation can be much more efficient. Specially nowadays machine learning can find much better solutions through reinforcement learning. And AI cluster data flow is much more predictable than what goes over the internet as a whole.
cma | 19 hours ago
smj-edison | 18 hours ago
scotty79 | 13 hours ago
cma | 8 hours ago
https://arxiv.org/abs/2304.01433
scotty79 | 7 hours ago
Basically they use lens array (one lens per input fiber) to focus each specific beam on its assigned mirror and mirror directs the light to the chosen output that also each has a lens in front of the fiber going out.
adastra22 | a day ago
netsharc | 10 hours ago
jMyles | 23 hours ago
wmf | 23 hours ago
Animats | 23 hours ago
The core idea: When a message arrives at the sender’s transport module, Homa divides the message into two parts: an initial unscheduled portion (the first RTTbytes bytes), followed by a scheduled portion. The sender transmits the unscheduled bytes immediately, using one or more DATA packets. The scheduled bytes are not transmitted until requested explicitly by the receiver using GRANT packets.
So it sends blind for short requests, then needs a go-ahead from the receiver. That's reasonable when the main application is a remote procedure call. It's reminiscent of QNX's networking protocol, which is also single packet message request/response but can also handle arbitrarily long messages.
What makes this work today is that per-packet processing overhead in hardware switches is low vs. per-byte overhead. In early software driven switches, per-packet overhead tended to dominate, and sending small packets was very inefficient. In modern hardware switches, where FPGAs are doing the processing, the per-packet overhead is low enough that small packets are not inefficient.
It's amusing that web stuff is so bloated today that any transaction under 1MB is considered "small". So this is not a suitable protocol for open web use.
[1] https://people.csail.mit.edu/alizadeh/papers/homa-sigcomm18....
Procrastes | 21 hours ago
Also, thanks for the clear summary and context.
wmf | 21 hours ago
throw0101c | 22 hours ago
> Link Level Flow Control: InfiniBand uses a credit-based algorithm to guarantee lossless HCA-to-HCA communication. RoCE runs on top of Ethernet. Implementations may require lossless Ethernet network for reaching to performance characteristics similar to InfiniBand. Lossless Ethernet is typically configured via Ethernet flow control or priority flow control (PFC). Configuring a Data center bridging (DCB) Ethernet network can be more complex than configuring an InfiniBand network.[19]
* https://en.wikipedia.org/wiki/RDMA_over_Converged_Ethernet
> A sending station (computer or network switch) may be transmitting data faster than the other end of the link can accept it. Using flow control, the receiving station can signal the sender requesting suspension of transmissions until the receiver catches up. Flow control on Ethernet can be implemented at the data link layer.
* https://en.wikipedia.org/wiki/Ethernet_flow_control
lokar | 22 hours ago
When a link becomes saturated working out how to manage that is a hard problem.
I have only used RoCE once at scale, it was really finicky. We would get big waves to pause frames that stalled everything.
wmf | 22 hours ago
teraflop | 22 hours ago
Consider a very simple topology:
Say hosts A and B are both sending data to C, as fast as they can, via switches S1 and S2 (which are connected via a high-speed link). And say the sum of these two flows is more than the capacity of the link to C.S2 is receiving packets destined for C faster than it can forward them, but sending an Ethernet pause frame from S2 to S1 is not a very productive way to alleviate the situation, because it also disrupts any traffic that would be bound for D. It just moves the bottleneck elsewhere and causes collateral damage.
stingraycharles | 20 hours ago
I agree that Ethernet flow control is insufficient, but given that NVidia in an act of brilliant foresight acquired Mellanox a decade ago, I’m fairly certain that this is how all these AI clusters are actually deployed, not using TCP, and maybe not even using Ethernet but infiniband instead.
etc-hosts | 15 hours ago
if not, you get RoCE over ethernet.
spwa4 | 11 hours ago
For IB the answer is generally a centralized controller, a lossless flow control algorithm and the credit system with subnet coordination. Ie. the packets that would be dropped due to buffer problems never get sent because the originating node knows it's out of credits. This means packet drops should be nonexistent and so the protocols make the assumption only one packet in a zillion ever gets dropped. On the down side: the tiniest of problems on IB networks have incredible impact (but it's realistic to just not have those)
IB doesn't interface with others. Yes, there's an ethernet bridge but it's not going to work nearly as well once you insert that bit of kit.
In ethernet packets are dropped when congestion is experienced, and you configure through priority queues whose packets get dropped first (you configure priority), and "stuff will sort itself out" like on the internet. RoCE tries to construct IB semantics from this, but there's limits. On the plus side: this is how the internet works and so the whole actually works somewhat reasonably over internet or other long-distance links (and tunnels into Neoclouds and/or Clouds). On the downside it's never going to match IB in speed (it's still going to get to 90% or so link speed, 98% if you just naively check the link utilization and ignore the effect retransmissions will have)
throw0101a | 9 hours ago
Worth noting that Nvidia/Mellanox is currently the only vendor for IB, but that has not always been the case. I had Intel IB switches for the backend network of my Isilon at my last job (we starting using them when they were 'just' Isilon, then EMC Isilon, and now Dell-EMC Isilon); only post-Dell did they start going Ethernet-only, and now with AI/ML they're re-offering IB as an option (LOL).
So Nvidia is kind of being 'rewarded' for sticking with a technology that everyone gave up on (I guess HPC wasn't big enough to bother with).
Besides IB, the other options being used in Top 500 are Slingshot (Eth-base?) and Omni-path (Category: Interconnect Family):
* https://www.top500.org/statistics/list/
Veserv | 22 hours ago
1. No way to detect whole RPC loss. Since there is no outer connection state, if every packet in the send-side of a RPC is lost then there is no way for a server to detect that it should issue a resend. The RPC is just lost to the ether. This affects small messages, like messages that fit in a single packet, more since there are fewer packets in the send-side.
2. Related to the above, there is no builtin encryption support. So, if you want encryption then you need to layer it either above or below.
3. Benchmarked performance is awful. The 60 kB average message case in [2] Table 4 takes 5(!) hyperthreads to average 20 Gbit/s. That is just 4 Gbit/s per hyperthread. Even a totally naive one-packet per system call network protocol design and implementation should get to ~8 Gbit/s per hyperthread. 30 Gbit/s per hyperthread is easy with just a little focus on performance.
4. Despite all the performance design problems in QUIC (though still faster than Homa) it already solves basically every problem Homa is trying to solve in a much cleaner way. Stream IDs correspond to RPC IDs. Stream Max corresponds to Grants. Multiple streams under single Client allows prioritization.
Except you do not randomly lose entire messages. You can compact small messages into packets. You get more precise RTT time allowing more accurate pacing/congestion calculations. You get builtin encryption. It survives ossified middleboxs. It has multiple ack frames/packets reducing ack overhead.
The only real difference is that Homa uses explicit receiver Resend instead of implicit sender Resend. Except that actually consumes significantly more receiver resources in non-trivial loss scenarios especially due to Resend packets only supporting a single Resend span. It also incurs higher latency and has higher requirements on the entire lossy network path due to all the extra transiting data that you do not want to lose.
And those are just serious problems off the top of my head after reloading the RFC into my head. I can come up with some more if needed.
[1] https://github.com/johnousterhout/homa-rfc/blob/main/draft-o...
[2] https://www.usenix.org/system/files/atc21-ousterhout.pdf
KerrAvon | 21 hours ago
Veserv | 21 hours ago
However, software QUIC implementations are generally much slower than software TCP implementations. This is not due to any fundamental protocol design limitations, it is just that most/every QUIC implementation is poorly implemented for maximizing performance.
You can, of course, do multiple times more throughput than QUIC or TCP with a protocol better designed for performance, but that is likely orthogonal to your question.
HackerThemAll | 10 hours ago
bold claim. Go tell that to Cloudflare, for example. And/or show us your implementation.
Veserv | 9 hours ago
The official MsQuic benchmarks [2] tuned for the standard QUIC benchmark [3] with massively over-provisioned hardware does not even reach 8 Gbit/s on large sends with negligible RTT and loss.
If the highest performing public implementation can not even reach a paltry 8 Gbit/s it is fair to say they are all poorly performing.
As for what I would consider “okay” that would be around 30 Gbit/s per core, bottlenecking on decryption (assuming you do not use the allowed null encryption for benchmarking), or bottlenecking on ~1/3 of the bulk 1 Kbyte memcpy() throughput on your hardware, whatever comes first. “Good” is on the order of 100 Gbit/s per core or the same external limiting factors.
[1] https://www.sciencedirect.com/science/article/pii/S014036642...
[2] https://microsoft.github.io/msquic/
[3] https://datatracker.ietf.org/doc/html/draft-banks-quic-perfo...
wmf | 21 hours ago
2. TCP doesn't have encryption either. They should probably show DTLS working though.
4. You can't just say QUIC is faster without testing it.
Veserv | 18 hours ago
That is a sender-driven resend and as such relies on the sender identifying the “send lost” condition. However, the entire protocol is designed around not doing that and thus you can only feasibly rely on “should have received a reply by now” as your timeout.
But Homa is intended to be a RPC protocol. So you send to server, wait for server to process command, then wait for server to send reply. Your timeout depends on the variable and heterogeneous command processing time.
Even if you were able to give a separate correctly tuned timeout for every possible RPC that is still awful. Any RPC with long processing time should not trigger the timeout until the expected reply time, but if that is far larger than the RTT then you are waiting a tremendous amount of time.
For instance, a select query that is only a few bytes long (and thus fits in one packet) could take seconds on a large database even though the database server is physically nearby and only microseconds away. In that case you would have to wait seconds before timing out instead of just microseconds like a ack-based design could achieve. The worst-case transport latency becomes application-specific instead of related to the application-independent transport parameters like RTT. That is troubling.
2. TCP can be forgiven for that given it predates asymmetric cryptography and even DES. Not considering it relevant in a new protocol this side of the millennium is much less reasonable.
4. MsQuic at 7.5 Gbit/s: https://microsoft.github.io/msquic/
wmf | 17 hours ago
Veserv | 16 hours ago
And then in basically every other respect it is worse including compatibility.
boesboes | 14 hours ago
tcdent | 20 hours ago
MrDrMcCoy | 19 hours ago
vlovich123 | 19 hours ago
Because the NSA and similar organizations will 100% infiltrate the physical infrastructure. It’s easier and harder if not impossible to detect.
Veserv | 19 hours ago
locknitpicker | 15 hours ago
This is self-evident, isn't it?
Other than the buzzword factor, why do you think name-dropping AI makes a difference?
preisschild | 11 hours ago
swyx | 18 hours ago
ContinuityLab | 17 hours ago
thelastgallon | 14 hours ago
londons_explore | 14 hours ago
That is relevant for musks in-space computation. There, latency between compute nodes, due to being hundreds of kilometres apart, cannot be low.
This pretty much limits the size of AI models which can be trained on a space datacenter to a single satellite to avoid incurring huge latency overhead both in training and inference.
With the smartest models getting bigger and bigger, that really doesn't look good for space based compute.
rwmj | 13 hours ago
musicale | 14 hours ago
jabl | 12 hours ago
Blog post: https://openai.com/index/mrc-supercomputer-networking/
Paper: https://cdn.openai.com/pdf/resilient-ai-supercomputer-networ...
Spec: https://www.opencompute.org/documents/ocp-mrc-1-0-pdf
MRC is an extension of RoCE that does out-of-order packet delivery, packet spraying across all available paths using SRv6 routing, combination of ECN and packet trimming for congestion control (so sender-based CC like TCP and unlike Homa).
throw0101c | 8 hours ago
> Today's show puts the “heavy” in Heavy Networking with a discussion about Multipath Reliable Connection (MRC). MRC is about getting an even distribution of RDMA traffic across equal cost multipath links in massive data centers, and doing so over plain old lossy Ethernet. Ethan, Drew, and guests discuss how package spraying, resilience to link and fabric failures, and congestion control allow MRC to be a specialized transport for AI workloads.
> Our guests today are from Broadcom: Rip Sohan, Distinguished Engineer; Eric Davis, Master Engineer; and Eric Spada, Technical Director and Distinguished Engineer.
* https://www.youtube.com/watch?v=AQhoTd9rf60
up2isomorphism | 7 hours ago