The End of TCP for AI Clusters — John Ousterhout, Stanford
Description
Split a job across several nodes, let their GPUs compute, then have them exchange a little metadata before the next round. While that exchange happens every GPU sits idle, and the whole cluster waits on the slowest one. Back when a compute phase ran five seconds and the exchange took a few milliseconds, nobody cared. Agentic inference has pushed those compute phases down to milliseconds, so the synchronization now costs about what the work costs. John Ousterhout asked the room who had reason to believe small message latency was capping their throughput, and said afterward that more hands went up than he expected. His argument is that AI traffic has quietly changed shape: it used to be gigabytes of gradients where only throughput mattered, and it is increasingly tiny coordination messages, checking whether an entry sits in a distributed KV cache or clearing a barrier, where the 99th percentile is what bites. The structural problem he identifies is that congestion control has always been the sender's job, and the sender is at the wrong end of the network. Queues build at the last hop into the receiver, so the sender learns about them secondhand, through roughly one bit of information, several round trips late. It has been known for twenty years and it still oscillates. Homa, built at Stanford out of Behnam Montazeri's dissertation, inverts it: the receiver runs congestion control, because it knows from the first packet exactly how much more is coming, and it releases the rest by issuing grants. Messages replace byte streams, so short ones can overtake long ones instead of queueing behind them. Speaker info: - https://x.com/johnousterhout - https://web.stanford.edu/~ouster/cgi-bin/home.php Timestamps: 0:00 - Why latency is becoming the metric that matters 2:19 - The old workload: gigabytes and throughput 3:09 - The new workload: metadata and coordination 4:24 - How one slow exchange stalls every GPU 6:07 - Incast, and where the queue actually builds 7:10 - Why conge
Summary
Generated by claude-sonnet-4-5-20250929At-a-Glance
- Verdict: Watch fully
- Core thesis: AI workloads are shifting from throughput-dominated large transfers to latency-sensitive small messages (especially for inference/agentic systems), and TCP/RDMA architectures fundamentally cannot handle this because sender-side congestion control causes 10x+ tail latency compared to HOMA's receiver-driven design.
- Why it matters: If Ken's agent systems involve distributed inference, KV cache lookups, or barrier synchronization across nodes, tail latency is already or will soon throttle GPU utilization; HOMA offers a production-ready kernel module that delivers ~13x lower P99 latency for small messages and 2x improvement even for large transfers.
- Best use: Use this to evaluate whether network tail latency is a hidden bottleneck in distributed agent orchestration or inference workloads, and to understand why receiver-side congestion control is architecturally superior for data center messaging patterns.
Executive Summary
Professor John Ousterhout argues that AI workloads are undergoing a structural shift. Training workloads remain dominated by massive multi-gigabyte transfers where throughput is king, but inference and agentic workloads are becoming increasingly granular, requiring frequent small message exchanges for metadata queries (e.g., checking distributed KV caches) and barrier synchronization between compute phases. As agentic systems push token generation into millisecond-scale compute cycles, even millisecond-range tail latency in network synchronization can waste significant GPU cycles, because GPUs sit idle waiting for the slowest message in an incast to complete before moving to the next computation round.
The root cause of high tail latency is incast congestion: multiple senders transmit to one receiver simultaneously, packets queue at the last-hop switch egress port, and short messages get stuck behind long ones. TCP and RDMA (RoCE) rely on sender-side congestion control, where senders must infer congestion from ECN markings or packet loss, then adjust transmission rates over multiple round trips. This introduces control lag, prevents stable convergence, requires queue buildup to signal congestion (which means latency is already occurring), and cannot prioritize messages because TCP/RDMA treat data as an undifferentiated byte stream with no message boundaries. Head-of-line blocking is unavoidable: a short message queued behind two large messages in a TCP stream cannot bypass them.
Ousterhout presents HOMA, a clean-slate transport protocol developed at Stanford (originally a PhD dissertation, now his full-time project as a semi-retired professor). HOMA is message-based (the fundamental unit is a request/response RPC pair), uses receiver-side congestion control, and exploits modern switch priority queues (typically 8 per egress port). The receiver knows message length from the first packet, granting the receiver complete congestion visibility and enabling it to pace grant packets that tell senders when to transmit scheduled packets. Short messages use higher-priority queues and bypass long-message queues. This implements shortest-remaining-processing-time-first (SRPT) scheduling. In benchmarks mixing message sizes from 50 bytes to 1MB, HOMA achieves ~13x lower P99 latency for short messages (<100µs vs. >1ms for TCP) and nearly 2x better latency even for the longest messages (due to run-to-completion scheduling vs. TCP's fair scheduling). HOMA is available as a Linux kernel module on GitHub, currently being upstreamed, and Ousterhout personally offers hands-on support for production trials.
Key Takeaways
- Claim: AI inference and agentic workloads are shifting from throughput-bound large transfers to latency-sensitive small message exchanges, especially for metadata coordination and barrier synchronization. | Evidence: Ousterhout cites distributed KV cache lookups and barrier synchronization at the end of compute phases as examples; agentic workloads are pushing compute cycles down to millisecond timescales, where millisecond-range synchronization latency wastes GPU cycles. Audience poll showed multiple attendees already experiencing small-message latency impacting throughput. | Implication: If Ken's agent orchestration involves distributed inference, multi-node coordination, or cache lookups across nodes, tail latency is likely or will become a bottleneck as workloads scale or become more granular. | Caveat: Training workloads remain dominated by massive transfers; the shift is most pronounced in inference/agentic systems, not universal across all AI workloads.
- Claim: TCP and RDMA suffer from high tail latency in mixed workloads because they use sender-side congestion control, which has control lag, oscillates between over- and under-sending, and requires queue buildup to detect congestion. | Evidence: Senders detect congestion via ECN markings or packet loss, then adjust rates over multiple round trips; by the time rates stabilize, network conditions have changed (new transmissions start, old ones finish). The research community has known this for 20+ years with incremental improvements but no fundamental fix. ECN requires queues to fill before signaling, meaning latency is already occurring. | Implication: Do not assume incremental TCP/RDMA tuning will solve tail latency for Ken's workloads; the sender-side control loop is architecturally unsuited to data center incast patterns.
- Claim: TCP and RDMA treat data as a byte stream with no message boundaries, causing head-of-line blocking and preventing prioritization of short messages. | Evidence: A series of messages serialized into a TCP stream has no color-coding; TCP cannot know message size until complete, cannot prioritize short messages, and a small message queued behind two large messages cannot bypass them in the stream. | Implication: Any system Ken builds that multiplexes large and small messages over TCP/RDMA connections will experience unavoidable head-of-line blocking; architectural change (not tuning) is required.
- Claim: HOMA's receiver-side congestion control gives the receiver complete information about congestion and enables sub-round-trip response to incast, avoiding oscillation and queue buildup. | Evidence: The receiver sees the first packet and knows total message length, so it knows exactly how much data is incoming across all concurrent transfers. It sends grant packets to pace transmission, delaying grants to reduce congestion and favoring shorter messages. This eliminates the multi-round-trip adjustment loop required by sender-side control. | Implication: Receiver-driven control is not an optimization; it is a fundamental architectural advantage for data center workloads where congestion occurs at the last hop to the receiver.
- Claim: HOMA uses message boundaries and modern switch priority queues (typically 8 per egress port) to implement SRPT scheduling, allowing short messages to bypass long-message queues. | Evidence: HOMA marks short messages with higher priority queue tags; in incast scenarios, long messages queue in the lowest priority queue while short messages immediately bypass them. HOMA sends only the first few 'unscheduled' packets of each message, then waits for receiver grants for 'scheduled' packets. | Implication: If Ken's infrastructure has modern switches, HOMA can deliver low tail latency without requiring network hardware upgrades; check switch priority queue capabilities before deployment. | Caveat: Requires modern data center switches with priority queue support (common in current infrastructure but not universal in older deployments).
- Claim: HOMA achieves ~13x lower P99 tail latency for short messages (<100µs vs. >1ms for TCP) and nearly 2x better latency even for large messages in mixed workloads. | Evidence: Benchmark with message sizes from 50 bytes to 1MB: TCP P99 for short messages >1ms, HOMA <100µs. For 1MB messages, HOMA is ~2x faster. The large-message improvement comes from run-to-completion scheduling (HOMA completes one message before starting another) vs. TCP's fair scheduling (time-slices across all active flows). | Implication: HOMA is not a niche short-message optimizer; it improves both tail and median latency across the entire message size spectrum, making it a viable TCP replacement for Ken's workloads. | Caveat: Benchmark is a controlled lab environment (multiple machines exchanging variable-length messages); production gains depend on workload characteristics and network topology.
- Claim: HOMA is production-ready as a Linux kernel module, available on GitHub, currently being upstreamed into the Linux kernel, and Ousterhout personally offers hands-on support for deployment. | Evidence: Ousterhout is semi-retired from Stanford and spending 100% of his time on HOMA; he offers to help with setup, answer questions, and provide bug fixes. Email and GitHub links provided in the talk. | Implication: HOMA is not vaporware or a research prototype; if Ken identifies tail latency as a bottleneck, he can trial HOMA with direct support from the creator, though kernel-level changes require careful ops planning. | Caveat: Kernel module status means it requires kernel-level deployment and compatibility testing; upstreaming is in progress but not yet complete, so production adoption carries integration risk.
Detailed Brief
Why Sender-Side Congestion Control Fails for Data Centers
- Claims: Sender-side control was designed for wide-area networks where congestion can occur anywhere along the path, not data centers where congestion is localized to the last hop.; ECN marking was introduced to avoid packet loss as a congestion signal, but it still requires queue buildup to trigger marking thresholds.; Multiple senders to the same destination must all adjust rates simultaneously with only one bit of information (congestion yes/no), leading to coordination failure.
- Evidence: In old TCP, packet loss was the congestion signal; modern TCP/RDMA use ECN, where switches mark packets when egress queue length exceeds a threshold, receiver echoes the mark back to sender in acknowledgments.; Senders must gradually adjust rates over multiple round trips to find the right rate, but by then network state has changed.; The research community has published extensively on this for 20+ years with improvements but no fundamental solution.
- Implications: Sender-side control is not a bug to be fixed; it is a design choice optimized for wide-area networks that is structurally mismatched to data center incast patterns.
HOMA's Three Architectural Pillars
- Claims: Message-based abstraction: HOMA's fundamental unit is a request/response RPC pair, not a byte stream.; Receiver-side congestion control: The receiver, not the sender, decides when and how much data to request via grant packets.; Priority queue exploitation: HOMA dynamically assigns packets to switch priority queues based on message length and scheduling state.
- Evidence: Message length is embedded in the first packet, giving the receiver full visibility into incoming load.; Unscheduled packets (the first few packets of each message) are sent immediately; scheduled packets are sent only when the receiver issues a grant.; Modern data center switches have 8 priority queues per egress port; HOMA uses packet header fields to select the queue.
- Implications: HOMA is not a TCP variant; every major design decision differs from TCP/RDMA, requiring a mental model shift when reasoning about performance.
Run-to-Completion vs. Fair Scheduling
- Claims: HOMA uses run-to-completion: finish one message before starting another, rather than time-slicing across active flows.; TCP/RDMA use fair scheduling: divide bandwidth equally among all active flows, which causes long flows to progress slowly and prevents short flows from completing quickly.
- Evidence: HOMA's 2x improvement on large messages (not just short messages) is attributed to run-to-completion, though Ousterhout notes he did not have time to explain the mechanism in detail.
- Implications: HOMA's large-message performance gain suggests that even for workloads without small messages, there may be throughput benefits from switching to HOMA.
Notable Concepts & Terms
- Incast: Multiple senders simultaneously transmitting to a single receiver, causing packet queue buildup at the last-hop switch egress port; the primary cause of tail latency in data center networks.
- Tail Latency (P99): 99th percentile latency; critical for systems where the slowest message in a batch determines overall throughput (e.g., barrier synchronization where GPUs wait for all nodes to finish).
- ECN (Early Congestion Notification): Mechanism where switches mark packets when egress queue length exceeds a threshold; receivers echo the mark back to senders, signaling congestion without packet loss.
- Head-of-Line Blocking: Phenomenon where a short message queued behind large messages in a byte stream cannot bypass them, causing delay; inherent to stream-based protocols like TCP.
- SRPT (Shortest Remaining Processing Time First): Scheduling policy that prioritizes shorter messages to minimize average completion time and tail latency; implemented in HOMA via receiver grants and priority queues.
- Unscheduled vs. Scheduled Packets: HOMA's two packet types: unscheduled packets (first few packets of a message, sent immediately) and scheduled packets (sent only when receiver issues a grant).
- Run-to-Completion Scheduling: Policy where a flow completes one message before starting another, rather than time-slicing bandwidth across all active flows (fair scheduling); used by HOMA to improve large-message performance.
- RoCE (RDMA over Converged Ethernet): The underlying transport used by RDMA in most modern data centers; Ousterhout uses 'RDMA' and 'RoCE' interchangeably in the talk.
Operator Notes / Why Ken Should Care
- Instrument distributed agent systems and inference clusters to measure P99 latency for small messages (metadata queries, cache lookups, barrier syncs); if P99 > 500µs, investigate whether incast is occurring.
- If tail latency is confirmed as a bottleneck, trial HOMA on a test cluster before production; contact Ousterhout directly for setup support (email provided in talk).
- Verify that production switches support priority queues (8 queues per egress port is typical in modern data center switches); this is a prerequisite for HOMA's performance gains.
- Do not assume TCP/RDMA tuning (buffer sizes, ECN thresholds, pacing algorithms) will solve tail latency; sender-side congestion control has fundamental architectural limits for data center incast.
- Evaluate whether Ken's agent orchestration or distributed inference systems currently multiplex large and small messages over shared TCP/RDMA connections; if so, head-of-line blocking is unavoidable without protocol change.
- Monitor GPU utilization during barrier synchronization or multi-node coordination phases; idle GPU time during synchronization is a signal that network tail latency is throttling compute throughput.
- Consider HOMA adoption risk: kernel module requires kernel-level deployment and compatibility testing; upstreaming is in progress but not complete, so plan for integration effort.
- Understand that HOMA is not a short-message-only optimizer; the 2x improvement on large messages suggests benefits even for workloads without granular messaging, though Ousterhout did not explain the mechanism.
Source/Metadata
- Title: The End of TCP for AI Clusters — John Ousterhout, Stanford
- Transcript words: 6408
- Duration seconds: 1128
- Timestamp note: Timestamps were manually derived from estimated positions in the 1128-second video; chapters were not present in the transcript.
Transcript
Please welcome to the stage the Professor Emeritus at Stanford University, John Oosterhout. Good morning, it's really great to be here to talk about the network side of AI applications and in particular to make the case that latency matters and is probably going to be mattering more in the future. But I just want to say, this talk is unusual for me. I've never before given a talk where there are fog generators in the auditorium. It's just a really San Francisco experience, I guess. So it's well known that AI workloads depend on really great networking performance in order to achieve their own performance. And of course that's because the workloads are so large that they have to be distributed across machines and then you have to communicate between the machines. But what I want to talk about today is that it seems that those workloads are changing. And so I hope to do three things over the next 15 or 20 minutes. First, to convince you that in fact the workloads are changing and that whereas the workloads used to be completely dominated by large transfers, where throughput is the key metric that matters, we're seeing more and more smaller transfers where the latency is crucial. The second thing I hope to do is to convince you that legacy protocols like TCP and RDMA are poorly suited to this environment. They weren't designed for this environment and unfortunately they suffer from very high tail latency when you mix small messages with large ones. And I'll talk a little bit about why that's the case. Then third, I'd like to introduce HOMA, which is a new protocol we've developed at Stanford that actually was designed in a clean slate, redesigned to handle data center workloads like these. And in fact it does quite well on those workloads and can reduce tail latency by an order of magnitude or more. So I'll tell you a little bit about HOMA. So let's dive in. First, workloads. Historically, AI workloads have consisted of enormous transfers between machines. And that's all that really mattered. Gigabytes of data for things like weight gradients and so on. In these workloads, what you really care about is throughput. How many gigabits per second you can pump through the pipes. And these are relatively easy workloads for networks because if it takes a while to set up the connection and start the transfer, it doesn't matter. The transfers go on for so long that all that really matters is the throughput. And so in these environments, TCP and RDMA perform pretty well. By the way, when I say RDMA, what I really mean is RoCE, RDMA over converged Ethernet, which is the underlying transport that's used by RDMA for most purposes today. So, anyhow, the old workloads, big transfers, throughput matters, the legacy protocols work pretty well. However, it appears that the workloads are changing. They're becoming more granular with smaller chunks of computation and smaller exchanges of data. And this seems to be particularly true in the world of inference and also in agentic workloads. Not so much for training workloads, they're still massive transfers. And so what's happening is that more and more there are small message exchanges, typically for things like metadata and coordination, such as checking to see if a particular entry is present in a KV cache that's distributed. Or doing barrier synchronization at the end of periods of compute. And for these workloads, what really matters is latency. That is, what's the round trip time to send some small piece of data across the network, do a little bit of computation, and get a small result back again. And in fact, it isn't just latency or average latency that matters. What really matters is tail latency. That is, you'd like to know that if we send a whole lot of small messages, all of them will complete quickly. So, for example, we typically measure things like 99th percentile latency. And if we have high tail latency, that can limit the overall throughput of the system. So here's an example. Suppose a common thing is to take a workload and split it up across several nodes, which do intensive computation and using their GPUs for some period of time. And then once they've all finished their computation, you do some small exchange between the nodes, exchange data, metadata, and then it'll go on to the next round of computation. And while that exchange is happening, that synchronization is happening, the GPUs are sitting idle. So if even one of those exchanges takes a long time, it turns out the whole process stalled. You need all of those exchanges to complete before you can go on to the next phase of computation. Now if the computation phase is, say, five seconds, and it takes a few milliseconds for the exchange, not a problem. And that's historically what it's been. But now with the agentic workloads where you're trying to pump out tokens relatively rapidly at a regular rate, the periods of computation are getting down into the millisecond timescale. And if it also takes milliseconds to do that synchronization, then you're wasting a significant fraction of your GPUs resources waiting for the synchronization to occur. So I'm curious. I'd like to just do a quick audience poll here. Is there anybody here where you have reason to believe that the latency of small messages is impacting the overall throughput of your applications? If so, can you just raise your hand? See, is there anybody out there today? Actually more hands than I expected. So quite a few people out there are raising their hands. I think this problem is likely to get worse as the trends continue. So what's going on? Why is tail latency bad? Well, typically the cause is congestion resulting from incast. So incast is when several nodes all decide simultaneously to transfer data to some destination node. And if they all send large messages, the links are the same everywhere in the network. So three nodes can transfer three times as fast as one node can possibly receive. And so what happens is that packets accumulate at the last hop going to that destination in the top of rack switch at its egress port for the destination node. Then if some other node decides it wants to send a short message to that same destination, the short message gets stuck behind the long ones in the queue there. And that causes delay. In the worst case, so many packets arrive that the switch runs out of buffer space and it has to drop packets, and then there are timeouts and retransmissions that make everything even worse. So somehow we need some way to reduce the congestion in those queues. Somehow we have to get the sending nodes to stop sending so fast, so the queues don't just build up without limit. So the way this is done historically, virtually all network protocols before HOMA, including TCP and RDMA, congestion control is the responsibility of the sender. So senders somehow have to figure out that congestion is happening, and they have to slow down their rate of transmission. Now you might wonder why are senders doing it, because the congestion is way over at the other end of the data center network. How does the sender find out? Well in the old, really old days, the way they would find out is the queues would overflow and packets would get dropped. The sender would detect the packets got lost because it wouldn't get acknowledgements back, and it would assume that means there's congestion and then slow down its rate of transfer. That's really expensive. So today there are better techniques that mostly involve the switches providing information. So a top of rack switch, when it sees that the queue length for an egress port has reached some threshold, starting to fill, long before the queue overflows, it starts marking all of the packets that pass through with what's called early congestion notification, ECN marking. And so when those packets pass through to the receiver, the receiver sees the marking in the packets, and then when it communicates back to the sender next, for example, to send an acknowledgement, then it includes that marking that goes back to the sender, and now the sender sees, the sender realizes, there's congestion someplace, I've got to slow down my rate of transmission. So that's the basic idea. Unfortunately, getting this right is really hard, really hard. It's very hard for the congestion to figure out exactly how to set its rates, because it gets one bit of information. There's congestion someplace. And there are multiple senders all sending to the same destination. They're all trying to make adjustments simultaneously. How much do you cut back? And how do I know when I can ramp up again? And even worse, it's really hard to do this in a way that's stable, because there's control lag. That is, it takes time before the sender finds out that there's congestion. And in fact, using this process, it typically takes several round trips for the sender to gradually adjust its rate to get just the right rate to match the available bandwidth. But by the time you do that, in a network, things have changed. New transmissions have started, old ones have finished. And so these systems tend to never stabilize. They're constantly oscillating between sending too much and sending too little. Now, this problem's been around for a long time. It's been known in the research community for more than 20 years now. There have been tons of papers published on it. There have been some improvements made. That's undeniable. But we're still a long ways from anything that works well. And how do I know when I can ramp up again? And even worse, it's really hard to do this in a way that's stable, because there's control lag. That is, it takes time before the sender finds out that there's congestion. And in fact, using this process, it typically takes several round trips for the sender to gradually adjust its rate to get just the right rate to match the available bandwidth. But by the time you do that, in a network, things have changed. New transmissions have started, old ones have finished. And so these systems tend to never stabilize. They're constantly oscillating between sending too much and sending too little. Now, this problem's been around for a long time. It's been known in the research community for more than 20 years now. There have been tons of papers published on it. There have been some improvements made. That's undeniable. But we're still a long ways from anything that works well. And the problem isn't with the fundamental nature of doing congestion control on the sender side. It just doesn't work very well. So you end up with a lot of queue buildup. And in fact, the only way to find out that there's congestion is if there's queues. And so by that point, we're already experiencing delays. So that's a problem. There's one other problem with TCP and RDMA also is that their basic data model is a byte stream. Just a stream of bytes with no differentiation in it. So if you send a series of messages through a TCP socket, they get serialized into that stream. And on this slide, I've shown the messages appear like they have different colors in the stream. Well, there are no colors in real life. TCP has no idea where the message boundaries are. And that also makes life hard. For example, you don't know how much more data is coming. If you knew how big the message was, you'd know how much more is coming. And you can't prioritize short messages, which we'd really like to do. Get the short messages through faster. And you can end up with what's called head of line blocking, where somebody sends a series of messages to the same destination. And they send two really large ones. And then a small one after that, they get stuck behind them in that stream. And so it gets delayed. And again, you have tail latency issues. So all in all, TCP and RDMA are just not well suited to this environment. So what do we do? Well, what I'd like to do next is tell you about a new protocol called HOMA that we've developed at Stanford, which was based on a completely clean slate redesign for network transport. If you could start from scratch and rethink how you do transport for data centers, how would you do it? And it turns out in HOMA, virtually every major design decision is different from TCP and RDMA. TCP, for all the amazing things that it's done, is just not a good match to today's data centers, nor RDMA. So what HOMA does particularly well is manage a combination of large and small messages and make sure that the short messages have really low latency. So this started off as a PhD dissertation for one of my students, Benam Monteseri. And then the results were so great that I decided to make it my personal project to see if we could get it out of the lab and into production. As you may know, I'm not like most professors, and I love to code, and so I turned this into my own programming project. I created a kernel module for Linux. I'm currently working through the process of getting that upstreamed into the kernel. It's available on GitHub for download. So let me tell you a little bit about how HOMA works. I want to mention three things. First, it's message-based, not stream-based. In fact, the fundamental unit in HOMA is a remote procedure call, which consists of two things. A request message sent from a client to a server, and then a response message returned back from the server to the client. So the key thing here is that HOMA knows about message lengths. They're buried in the transport all the way down to the bottom. And this has a bunch of advantages. First, it allows us to predict the future. As soon as a receiver gets the first packet of a message, it knows exactly how much more data the sender wants to send. And that's so much more information for doing congestion control. Second, HOMA prioritizes shorter messages. It uses SRPT, shortest remaining processing time first, to try and prioritize shorter messages. And third, because messages are all independent, they're not serialized into a stream, every message is independent, shorter messages can bypass long ones, so they don't get queued behind long messages. The second thing about HOMA that's different is that it controls congestion from the receiver. And when you think about it, this makes sense, because the congestion happens primarily at that last downlink to the receiver. And so the receiver has way more information. In fact, with HOMA, as soon as it gets the first packet of a message, it knows exactly how much more is coming. So it has essentially complete information about congestion, and it can therefore respond to congestion much more quickly. So it's really important to understand how much more quickly and much more precisely. The way things work with HOMA is that when a sender has a message to send, it breaks it up into packets, but it only transmits the first few packets. Those are called unscheduled packets, to the receiver. Packets after that are called scheduled packets, and they only get transmitted when the receiver asks for them. So the receiver will send grant packets back. It will pace them out and send those back to the sender over time, telling the sender it's now time for you to send me the next chunk of data. And the receiver can delay those grants. So for example, if the receiver has ten messages that are incoming, there's no point in sending grants to all ten of them, because then you'll get congestion in the top of NIC queues. So it can use the grants to reduce congestion, and then it can also use the grants to give preference to its most favorite messages, which would be the shorter ones. So it's a way of implementing SRPT by favoring short messages. The third aspect of HOMA is that it takes advantage of the priority queues in modern switches. So modern data center switches have more than one queue at each egress port, typically eight, and they can be used in a priority mechanism where packets get transmitted preferentially from the highest priority queue. So I've shown only two queues on the slide here, but typically there's more than that. You can specify in packets, using the various fields of the packet, which queue it should go into. And so HOMA dynamically makes those choices in a way to give priority to shorter messages. So if we go back to the multicast example from a few slides ago, all of those long messages will pile up in the lowest priority queue. But if there's a short message coming, it will use a higher priority queue. And so it will immediately bypass all of the queued packets from the longer messages and get through to the destination more quickly. So how much of a difference does this make? So here's one sample benchmark that I use as part of my tuning and evaluation of HOMA. It consists of a workload of a bunch of machines on a network that are exchanging messages back and forth of different sizes, ranging from very small to very large. And on this graph you can see on the x-axis is the message length, so from about 50 bytes up to a megabyte. The y-axis shows you the round trip time for messages of that length. So this uses request and response messages that are the same length. You can see TCP in green, HOMA in blue, and the y-axis is round trip time, so lower is better. And for each protocol I've got two curves. One curve is the P50 curve, that's the median latency for messages of this length. And then P99 is the 99th percentile, i.e. tail latency for messages of this length. So I want to point out two things. First, the P99 for short messages is dramatically better for HOMA. So with TCP it's more than a millisecond tail latency, HOMA is less than 100 microseconds, about 13 times faster. Second, interestingly, you might think that because HOMA favors shorter messages, that long messages suffer and get worse performance. It turns out that's actually not the case. Even on the longest messages, HOMA is almost a factor of two better than TCP. I don't have time to explain that today, but it has to do with the fact that HOMA uses run to completion approaches, which are much more effective than the fair scheduling used by TCP. So just to wrap up, the role of short messages in AI appears to be increasing. I think it's likely that it's going to continue to increase. We'll see over the next year or two if that happens. And I just want to pose a question to you. As you're running your applications and measuring performance and seeing what the bottlenecks are, ask yourself, is high latency for short messages affecting your throughput? Second, interestingly, you might think that because HOMA favors shorter messages, that long messages suffer and get worse performance. It turns out that's actually not the case. Even on the longest messages, HOMA is almost a factor of two better than TCP. I don't have time to explain that today, but it has to do with the fact that HOMA uses run to completion approaches, which are much more effective than the fair scheduling used by TCP. So just to wrap up, the role of short messages in AI appears to be increasing. I think it's likely that it's going to continue to increase. We'll see over the next year or two if that happens. And I just want to pose a question to you. As you're running your applications and measuring performance and seeing what the bottlenecks are, ask yourself: is high latency for short messages affecting your throughput? If the answer is yes, then just know there is a solution available. If you give HOMA a try, you can probably reduce your tail latency by an order of magnitude or more. And by the way, HOMA is my life mission right now. I'm semi-retired from Stanford, and the reason I did that is so I can spend 100% of my time hacking on HOMA. So I'd be delighted to work with you and help you if you decide you want to experiment with HOMA. If you need help getting started, answer questions, bug fixes, whatever, I'd be happy to work with you to try and make you successful with it. So if that is interesting, feel free to contact me. My email is on the slide, or you can Google me and find me over the internet. So thanks very much for listening, and I hope to hear from some of you. I'll see you next time. They're becoming more granular with smaller chunks of computation and smaller exchanges of data. And this seems to be particularly true in the world of inference and also in agentic workloads. Not so much for training workloads, they're still massive transfers. And so what's happening is that more and more there are small message exchanges, typically for things like metadata and coordination, such as checking to see if a particular entry is present in a KV cache that's distributed. Or doing barrier synchronization at the end of periods of compute. And for these workloads, what really matters is latency. That is, what's the round trip time to send some small piece of data across the network, do a little bit of computation, and get a small result back again. And in fact, it isn't just latency or average latency that matters. What really matters is tail latency. That is, you'd like to know that if we send a whole lot of small messages, all of them will complete quickly. So, for example, we typically measure things like 99th percentile latency. And if we have high tail latency, that can limit the overall throughput of the system. So here's an example. Suppose a common thing is to take a workload and split it up across several nodes, which do intensive computation and using their GPUs for some period of time. And then once they've all finished their computation, you do some small exchange between the nodes, the exchange data, metadata, and then it'll go on to the next round of computation. And while that exchange is happening, that synchronization is happening, the GPUs are sitting idle. So if even one of those exchanges takes a long time, it turns out the whole process stalled. You need all of those exchanges to complete before you can go on to the next phase of computation. Now if the computation phase is, say, five seconds, and it takes a few milliseconds for the exchange, you know, not a problem. And that's historically what it's been. But now with the agentic workloads where you're trying to pump out tokens relatively rapidly at a regular rate, the periods of computation are getting down into sort of the millisecond timescale. And if it also takes milliseconds to do that synchronization, then you're wasting a significant fraction of your GPUs resources waiting for the synchronization to occur. So I'm curious. I'd like to just do a quick audience poll here. Is there anybody here where you have reason to believe that the latency of small messages is impacting the overall throughput of your applications? If so, can you just raise your hand? See, is there anybody out there today? Actually more hands than I expected. So quite a few people out there are raising their hands. I think this problem is likely to get worse as the trends continue. So what's going on? Why is tail latency bad? Well, typically the cause is congestion resulting from in-cast. So in-cast is when several nodes all decide simultaneously to transfer data to some destination node. And if they all send large messages, well, the links are the same everywhere in the network. So three nodes can transfer three times as fast as one node can possibly receive. And so what happens is that packets accumulate at the last hop going to that destination in the top of Rack Switch at its egress port for the destination node. Then if some other node decides it wants to send a short message to that same destination, the short message gets stuck behind the long ones in the queue there. And actually that causes delay. In the worst case, so many packets arrive that the switch runs out of buffer space and it has to drop packets, and then there are timeouts and retransmissions that make everything even worse. So somehow we need some way to reduce the congestion in those queues. Somehow we have to get the sending nodes to stop sending so fast, so the queues don't just build up without limit. So the way this is done historically, virtually all network protocols before HOMA, including TCP and RDMA, congestion control is the responsibility of the sender. So senders somehow have to figure out that congestion is happening, and they have to slow down their rate of transmission. Now you might wonder why are senders doing it, because the congestion is way over at the other end of the data center network. How does the sender find out? Well in the old, really old days, the way they would find out is the queues would overflow and packets would get dropped. The sender would detect the packets got lost because it wouldn't get acknowledgements back, and it would assume that means there's congestion and then slow down its rate of transfer. That's really expensive. So today there are better techniques that mostly involve the switches providing information. So a top of rack switch, when it sees that the queue length for an egress port has reached some threshold, starting to fill, long before the queue overflows, it starts marking all of the packets that pass through, with what's called early congestion notification, ECN marking. And so when those packets pass through to the receiver, the receiver sees the marking in the packets, and then when it communicates back to the sender next, for example, to send an acknowledgement, then it includes that marking that goes back to the sender, and now the sender sees, the sender realizes, oh, there's congestion someplace, I've got to slow down my rate of transmission. So that's the basic idea. Unfortunately, getting this right is really hard, really hard. It's very hard for the congestion to figure out exactly how to set its rates, because it gets one bit of information. There's congestion someplace. And there are multiple senders all sent into the same destination. They're all trying to make adjustments simultaneously. How much do you cut back? And how do I know when I can ramp up again? And even worse, it's really hard to do this in a way that's stable, because there's control lag. That is, it takes time before the sender finds out that there's congestion. And in fact, using this process, it typically takes several round trips for the sender to gradually adjust its rate to get just the right rate to match the available bandwidth. But by the time you do that, in a network, things have changed. New transmissions have started, old ones have finished. And so these systems tend to never stabilize. They're constantly oscillating between sending too much and sending too little. Now, this problem's been around for a long time. It's been known in the research community for more than 20 years now. There have been tons of papers published on it. There have been some improvements made. That's undeniable. But we're still a long ways from anything that works well. And the problem isn't with the fundamental nature of it, doing the congestion control on the sender side. It just doesn't work very well. So you end up with a lot of queue buildup. And in fact, you can see the only way to find out that there's congestion is if there's queues. And so by that point, we're already experiencing delays. So that's a problem. There's one other problem with TCP and RDMA also is that their basic data model is a byte stream. Just a stream of bytes with no differentiation in it. So if you send a series of messages, say, through a TCP socket, they get serialized into that stream. And on this slide, I've shown the messages appear like they have different colors in the stream. Well, there are no colors in real life. TCP has no idea where the message boundaries are. And that also makes life hard. For example, you don't know how much more data is coming. If you knew how big the message was, you'd know how much more is coming. And you can't prioritize short messages, which we'd really like to do. Get the short messages through faster. And you can end up with what's called head of line blocking, where somebody sends a series of messages to the same destination. And they send two really large ones. And then a small one after that, they get stuck behind them in that stream. And so it gets delayed. And again, you have tail latency issues. So all in all, TCP and RDMA are just not well suited to this environment. So what do we do? Well, what I'd like to do next is tell you about a new protocol called HOMA that we've developed at Stanford, which was based on a completely clean slate redesign for network transport. If you could start from scratch and rethink how you do transport for data centers, how would you do it? And it turns out in HOMA, virtually every major design decision is different from TCP and RDMA. TCP, for all the amazing things that's done, is just not a good match to today's data centers, nor RDMA. So what HOMA does particularly well is to manage a combination of large and small messages and to make sure that the short messages have really low latency. So this started off as a PhD dissertation for one of my students, Benam Monteseri. And then the results were so great that I decided to make it my personal project to see if we could get it out of the lab and into production. As you may know, I'm not like most professors, and I love to code, and so I turned this into my own programming project. I created a kernel module for Linux. I'm currently working through the process of getting that upstreamed into the kernel. It's available on GitHub for download. So let me tell you just a little bit about how HOMA works. I want to mention three things. First, it's message-based, not stream-based. In fact, the fundamental unit at HOMA is a remote procedure call, which consists of two things. A request message sent from a client to a server, and then a response message returned back from the server to the client. So the key thing here is that HOMA knows about message lengths. They're buried in the transport all the way down to the bottom. And this has a bunch of advantages. First, it allows us to predict the future. As soon as a receiver gets the first packet of a message, it knows exactly how much more data the sender wants to send. And that's so much more information for doing congestion control. Second, HOMA prioritizes shorter messages. It uses SRPT, shortest remaining processing time, first, to try and prioritize shorter messages. And third, because messages are all independent, they're not serialized into a stream, every message is independent, shorter messages can bypass long ones, so they don't get queued behind long messages. The second thing about HOMA that's different is that it controls congestion from the receiver. And when you think about it, this makes sense, because the congestion happens primarily at that last downlink to the receiver. And so the receiver has way more information. In fact, with HOMA, as soon as it gets the first packet of a message, it knows exactly how much more is coming. So it has essentially complete information about congestion, and it can therefore respond to congestion much more quickly than the other. So it's really important to understand how much more quickly and much more precisely. The way things work with HOMA is that when a sender has a message to send, it breaks it up into packets, but it only transmits the first few packets, those are called unscheduled packets, to the receiver. Packets after that are called scheduled packets, and they only get transmitted when the receiver asks for them. So the receiver will send grant packets back. It will paste them out and send those back to the sender over time, telling the sender, it's now time for you to send me the next chunk of data. And the receiver can delay those grants. So for example, if the receiver has ten messages that are incoming, there's no point in sending grants to all ten of them, because then you'll just get congestion in the top of RAC queues. So it can use the grants to reduce congestion, and then it can also use the grants to give preference to its most favorite messages, which would be the shorter ones. So it's a way of implementing SRPT by favoring short messages. The third aspect of HOMA is that it takes advantage of the priority queues in modern switches. So modern data set of switches have more than one queue at each egress port, typically eight, and they can be used in a priority mechanism where packets get transmitted preferentially from the highest priority queue. So I've shown only two queues on the slide here, but typically there's more than that. You can specify in packets, using the various fields of the packet, you can specify which queue it should go into. And so HOMA dynamically makes those choices in a way to give priority to shorter messages. So if we go back to the NCAST example from a few slides ago, all of those long messages will pile up in the lowest priority queue. But if there's a short message coming, it will use a higher priority queue. And so it will immediately bypass all of the queued packets from the longer messages and get through to the destination more quickly. So how much of a difference does this make? So here's a, on this slide I've got one sample benchmark that I use as part of my tuning and evaluation of HOMA. It consists of a workload of a bunch of machines on a network that are exchanging messages back and forth of different sizes, ranging from very small to very large. And on this graph you can see on the x-axis is the message length, so from about 50 bytes up to a megabyte. The y-axis shows you the round trip time for messages of that length. So this, this uses request and response messages that are the same length. You can see TCP in green, HOMA in blue, and the y-axis is round trip time, so lower is better. And for each protocol I've got two curves. One curve is the P50 curve, that's the median latency for messages of this length. And then P99 is the 99th percentile, i.e. tail latency for messages of this length. So I want to point out two things. First, the P99 for short messages is dramatically better for HOMA. So with TCP it's more than a millisecond tail latency, HOMA is less than 100 microseconds, about 13 times faster. Second, interestingly, you might think that because HOMA favors shorter messages, that long messages suffer and get worse performance. It turns out that's actually not the case. Even on the longest messages, HOMA is almost a factor of two better than TCP. I don't have time to explain that today, but it has to do with the fact that HOMA uses run to completion approaches, which are much more effective than the fair scheduling used by TCP. So just to wrap up, the role of short messages in AI appears to be increasing. I think it's likely that it's going to continue to increase. We'll see you over the next year or two if that happens. And I just want to pose a question to you. As you're running your applications and measuring performance and seeing what the bottlenecks are, ask yourself, is high latency for short messages affecting your throughput? If the answer is yes, then just know there is a solution available. If you give HOMA a try, you can probably reduce your tail latency by an order of magnitude or more. And by the way, HOMA is basically my life mission right now. I'm sort of semi-retired from Stanford, and the reason I did that is so I can spend 100% of my time hacking on HOMA. So I'd be delighted to work with you and help you if you decide you want to experiment with HOMA. If you need help getting started, answer questions, bug fixes, whatever, you know, I'd be happy to work with you to try and make you successful with it. So if that is interesting, feel free to contact me. My email is on the slide, or you can Google me too and find me over the internet. So thanks very much for listening, and I hope to hear from some of you. I'll see you next time.