Serving Netflix Video Traffic at 800Gb/s and Beyond [pdf]
nabstreamingsummit.com
nabstreamingsummit.com
Building software and productionalizing/scaling it are two very different problems, and the latter is far more difficult. Running a successful company always requires an unlimited number of very smart people who are willing to get their hands dirty optimizing every aspect of the product and business. Too many people today think that programming starts and ends at pulling a dozen popular libraries and making some API calls.
The needle keeps moving doesn’t it? A tremendous breadth of difficult problems can be effectively addressed by pulling together libraries and calling APIs today that weren’t possible before. Today’s hard problems are yesterday impossibilities. The challenge for those seeking to make an impact is to dream big enough.
Anecdotal, but most of the people I've worked with as ICs couldn't give a damn about that. They want dollarydoos.
One of the 10X-ers I know (they exist and are real), told me repeatedly how he'd much rather be doing his own thing. He hates the business needs. But income is important and that's why he's dedicated to doing it. I'm surprised at how focused and good he is given his disposition, and I want to hire him when I scale my business more. Drive and passion are sometimes just spontaneous.
An old CEO of mine even quipped that we were not family and that we were there to do a job. All true. Most of the people doing that job were only there for the money.
Most jobs that drive sales and revenue simply aren't fun or rewarding. There's lots of infrastructural glue and scaling. Tiring, boring, monotonous work. 24/7 oncall work. The money is good, though.
This is the sort of impressive work that I've never seen scale.
With that said, we are standing on the shoulders of giants. There are tons of other optimizations not mentioned in this talk where removing any one of them could tank performance. I'm giving a talk about that at EuroBSDCon next month.
https://www.wsj.com/articles/SB10001424052702304834704579401...
Is that all of Disney or just Disney+?
It doesn't seem like that would be a useful statistic if that includes completely unrelated positions (e.g. does that 20x statistic include Disney employees working at Disney Land/World serving up hotdogs? Because they probably don't contribute much to the streaming service)
Walmart's market cap per employee is probably much, much lower than Disney or Netflix, too. That doesn't mean Walmart is doing anything wrong.
It's closer to a live streaming problem than pre-encoded video like Netflix.
Having worked at Netflix I can say that the YouTube problem is much more complex.
I do think there is some temporal logic already present in Google’s algorithm which wasn’t part of the challenge.
[1] https://www.foxbusiness.com/technology/5-things-to-know-abou...
I've done a lot of video processing professionally (the server side stuff, exactly what Netflix does) and Netflix is by far the worst of all the streaming providers. They absolutely sacrifice the quality of the video to save bandwidth costs in aggregate and it shows (or more accurately it doesn't show, all the fidelity is lost).
Most companies maintain internal calculations of these sorts of things, and make rational decisions.
When you say that companies maintain internal calculations of the benefits, would you say that it’s (extremely roughly) something like: $10M benefit, need 5 core engineers + benefits + PM + testing lab etc etc -> we can spend up to $500k per eng give or take.
Or is the $10M one number (that would be held somewhat secretly internally at the company) and the salaries mostly represents where the market is? Does the (salary) market take into account the down-the-line $10M value?
Basically, could those engs negotiate to be paid more, or are they already sort of paid close to exactly what the group they’re part of generates in terms of revenue?
Thanks!
—
I see that you said $10M per person, not for the “network optimization group”. Hmm. So it would be fair to say that the engs are definitely not paid according to the value they generate..? I wouldn’t be surprised by that but just to confirm.
In terms of negotiation, it really depends on how differentiated your skills are. Short answer is that if you can convince management that it would be difficult to find other engineers who could deliver the optimizations you're delivering, yes, you have leverage.
Very highly skilled engineers in specific niches can basically price themselves like monopolists, because the company can easily figure out how much money they are leaving on the table by not hiring them. This is not like "feature work" engineers, whose value is very nebulous and unknown.
Secondly, yes $10M per employee of revenue or cash flow is pretty reasonable for similar companies. The prioritization is NOT “how many employees per $MM.” The allocation is “what opportunity is the highest $MM return per available employee.”
Where the CDN boxes go, you can't always just throw more hardware. There's a limited amount of space, it's not controlled by Netflix, and other people want to throw hardware into that same space. Pushing 800gbps in the same amount of space that others do 80gbps (or less) is a big deal.
It's still impressive. A 5x increase at that scale can be a phenomenal challenge. Where do you source the ingredients? Where do you build the factories (plural because at that scale you almost certainly have multiple locations in different geographic locales subject to different regulatory structures). Where do you hire the people? How do you manage it? What about the storage and shipping and maintenance of all the equipment and on and on? How much do you do in house how much do you outsource to partners? What happens when a partner goes belly up or can't meet your ever increasing needs?
Your comment is a great example of what the OP pointed out.
> The menu team comes up with interesting ideas like including kale in salads. The procurement team and suppliers then try to get the menu team to understand the challenges. How do you bring kale to 14,000 restaurants? As one example, when they introduced Blueberry Smoothies in the U.S., McDonald’s ended up consuming one third of the blueberry market overnight.
https://www.forbes.com/sites/stevebanker/2015/10/14/mcdonald...
I couldn't find any other source to back it up, but still wow! That's an absurd number.
In the case in animal product, there are almost certainly major operations worldwide that have been built and financed purely to serve McDonalds demand. They probably have to even build these out well before entering some markets.
They raise a lot of chickens in the US. I’ve hauled chicken nuggets or chicken breasts for McDonald’s in the past quite often.
I can’t even tell you where they grow blueberries.
Taking the software example, you can easily scale from 1 to 100 users on your own machine. You can handle thousands by moving to a shared host. Using off-the-shelf web servers and load balancers will help you serve a million+. From there on you'll have to spend a lot more effort optimizing and fixing bottlenecks to get to tens, maybe hundreds of millions. What if you want to handle a billion users? Five billion? Ten billion? It always gets harder, not easier.
Pushing the established limits of a problem takes exponentially more effort than reusing existing solutions, even though the marginal improvement may be a lot smaller. Getting from 99.9% to 99.99% efficiency takes more effort than getting from 90% to 99%, which takes more effort than getting from 50% to 90%.
You never pierce the scaling wall. It only keeps getting higher.
Each employee adds some overhead, which requires more employees... which requires more employees.
To add something useful as well besides snark, first of all, there are hard physical limits, which are sometimes well within context (you really shouldn’t try to outcompete light speed for example, relevant in some high-freq trading, infrastructure projects). Then you can try to increase headcount to any number, you won’t produce for example a better compiler. There are simply jobs that are more “serial” - the only way to win at those is to try to employ the very best of the field in a small team.
“Never underestimate the bandwidth of a station wagon full of tapes hurtling down the highway.“ -Andrew Tannenbaum
1M users and 10k employees is not in the range where you have crushingly impactful logistics.
That's a 2X increase. Now do it again and a half for a 5x. Crazy to say there's a "scaling wall" that once you "pierce" it's easy to scale up. It's the opposite, McDonald's already knows how to supply and sell X McRibs a year, there's no company that's ever sold 5X those McRibs so they have to figure it out themselves.
Anecdotally I experienced this when scaling my software product from 1 --> 10 --> 100's --> 1000's etc. of users.
Thats not to say 2x can't be a substantial challenge, as you pointed out. It gets harder (and IMO more fun) when you're at the bleeding edge of your industry.
This is absolutely not true. The closer you are to peak performance, the harder it is to scale, and the returns diminish heavily. At many major tech companies, there's a huge amount of effort into just 1% - 5% optimizations -- these efforts really require creative thinking and complex engineering (not just "scaling existing processes".) At the volumes these companies operate, even a 1% optimization is quite significant.
If you're on 100M users you're probably scaling vertically. So adding 5x more hardware shouldn't be a problem.
But when you're at 500M all of a sudden it makes sense to optimize further since the capital saved will be the same percentage(ish) but the money is worth peoples time all of a sudden.
I know that we don't care particularly about power savings in the DCs I've worked in, because they're relatively small. While bigtech will do all kinds of shenanigans to save a couple watts here and there, because it's worth it across your hundreds of thousands of servers.
If only things were that simple. https://en.wikipedia.org/wiki/Amdahl%27s_law
Building five "Netflixes" with identical content is possible; the amount of content wouldn't change (it would decrease, the cynic says); you just need parallel copies of everything (servers, bandwidth, etc).
The fun would come in syncing usernames, etc through the system.
It's an entirely different class of problem compared to "acquire resource, convert it, sell it".
Sorry, you gotta overhaul majority of your architecture and its components to scale by every 10x. It's not a single "scaling wall" to break through but it's more of a relentless stream of uphill battles. And this gets even more interesting when you reached to the point where there's no prior art for your problem, usually at hundreds of billions of users.
Making 5 things and making 5,000 things is as different as making 50,000 things and 1m things. There are always cost constraints at each level, and each design can only go so far.
I would like to hope nobody asks that. Video is the one of the, if the not the hardest data plumbing use-case on the internet.
A lot of these tricks being discussed here cannot be applied to Skype calls.
Pre-recorded video streaming is, under the hood, really just a high-volume variant of serving up static web pages. You have a few gigabytes of file to send from the server its stored on to the device that wants to play back the video. As this presentation demonstrates that isn't trivial at scale, but the core functionality of sending files over the internet is what it was designed to do from day one. Because you can generally download video across the internet faster than it can be played back its possible to build up a decent sized buffer which allows you to paper over temporary variance in network performance without the customer noticing.
Realtime video streaming has two variants. One to many Twitch style video streaming is relatively simple, since you can encode video into files and upload them to a server for people watching to download those files. This is how HLS streaming works, and most of the techniques Netflix use to optimise video delivery can also be applied here at the cost of adding latency between the event being streamed and people consuming it. That latency will often sit at about 30 seconds, and people generally find that acceptable.
Skype style realtime video streaming is much harder. You're taking video from one person's camera, and then sending it over the internet to one or more people's device. You can't do any sort of pre-processing on that, or stage the video on servers closer to the consuming users, because you have no way of generating that video until the point your users decide to start talking to each other. Because you can't pre-stage that video you need to be able to establish a network route between the people on a call, potentially in an environment where none of the participants have any open connection from the internet directly to the device they're streaming from. Slight fluctuations in network performance can potentially degrade video delivery to the point of it being unusable. The most common route to deal with that is systems that attempt to establish a direct connection (ideally over a local network) between participants, and if that doesn't work going via relay servers operated by the software provider. These servers provide a single point on the internet all parties can connect to, and then allow passing packets as if they were all on the same network.
I agree but since some optimization is theoretically impossible, perhaps it can be said that optimizable area is much smaller than other services. YouTube (many many videos, high quality, many livestreams, and supports VR) seems to most difficult service to fully optimize.
I'm not sure what you mean? Real time communication, both video and audio-only, have much lower latency requirements. You can't just buffer ahead when you have some spare bandwidth, like Netflix or YouTube can.
My point was intended to be that there's the same challenges and more - but it's not something I've thought about in depth (and certainly not had to work with), it maybe wasn't a very good characterisation because it's not the same on the other side either, no large file to serve because at the start of the call it doesn't exist yet for example, so perhaps I take it back.
In any case, it's hard to say what the 'greatest' engineering challenge is. You can make almost any kind of engineering really challenging, if you (or the market..) sets yourself a very low cost-ceiling.
It is much much harder because you can't do the "cache everything on the edge" solution. If storage was infinitely cheap and small, Netflix could run their entire business by sending a USB stick with every single movie/tv show they have on it encrypted to you, and everything would play locally. This is basically what they do with their edge servers/CDNs.
You can't do that with video calls, because the video/audio didn't exist 1 millisecond ago.
I'm not a telecommunications guy, but I had some professors back in college explain how difficult and fundamental the research of "ma-bell" was from the 60s through 80s. I'm talking Erlang, C++, CLOS circuits, etc. etc. The innovations from Bell Labs are nearly endless.
Telephone communications is one of the biggest sources of fundamental comp. sci research over the 1950s through 2000s.
A lot of innovations did come out of Bell Labs. But I'm pretty sure Erlang wasn't one of them.
Compression removes redundancy from your message. Error correcting codes introduces carefully controlled redundancies, so your message survives transmission errors, so your recipient can read it. And cryptography deals with 'scrambling' your message, so that no one else can read it (and also deal with authentication etc).
Scaling cost-effectively is.
(Just like sending a human to Alpha Centauri is hard, even if you had unlimited funds.)
If Netflix built out more slower servers, that would be acceptable scaling. I don't see any plausible scenario where that becomes too difficult. Even if they had billions of subscribers.
Eg enriching a tiny bit of uranium is 'relatively' easy. Enriching enough of the stuff to build an atomic bomb is a question of scaling. And, yes, in principle you can just use more and more centrifuges etc.
Humanity has already built some small infrastructure off-earth, like the ISS or satellite networks. Also various unmanned probes that toiled away for years in space and on Mars.
Building a habitat big enough to plausibly keep humans alive for the decades or more it would take to get to the Alpha Centauri would involve enormous scaling efforts. And that's not even touching on propulsion, yet.
However, yes, grand feats require more than just scaling.
But we're not scaling for scaling's sake. We're scaling with the number of customers.
If it was profitable to enrich some milligrams by hand, then you'd be able to get more scientists and engineers and better equipment and automate some parts and make even more profit. If you go the centrifuge route, it's not because you had to do the harder methods, it's because it's massively cheaper per gram of product.
Such a scenario is very different from "we can't do it at all, and we need to build complex infrastructure just to make it possible".
* after a certain scale
What I've read they burned a lot of money and hat large problems scaling nevertheless. Which I don't find too surprising, not because they are unable, but because it isn't easy to scale.
From my experience and from what I read scaling people roughly a power of ten is a larger change in an organisation and therefor likely a challenge. For _any_ technical process the boundaries might not be strictly a power of ten but i would say that scaling a power of a hundred is a challenge if this value is not already reached on any process in your organisation.
Scaling to - say - Paramount+ size should not be difficult if you're willing to pay AWS / Azure / GCP 10-100x what it would cost to serve it yourself (which in many cases actually makes sense).
It's possible at Netflix's size, they couldn't just run on AWS anymore. Though, given enough lead time and a realistic growth curve - I'm sure it's feasible.
Obviously scaling manufacturing is not a solved problem like (realistically) scaling network and compute usage.
Today I rarely use Netflix, and I wouldn’t pay for it. Periodically I open the app and add to my list those titles that are immediately visible which I know I want to watch, like comedy specials of comedians I’m familiar with, but I don’t inspect further, because if I haven’t already heard about a title from some external sources I trust, it’s not worth my time to check. That list just grows and grows, though titles are often removed as they become unavailable, but I never prune the list, because when I remove one title, I’m taken all the way back to the beginning of the list. Trying is just a wast of time.
It is baffling to me that anyone could interact with Netflix’s current UI and conclude that it was anything but a raging dumpster fire.
Jargon BS is invading people's heads and it has to stop.
Admitedly, it follows the same pattern as Netflix, but I like how it's more responsive and feels way simpler/lighter.
This is all in stark contrast to services like HBO Max and Disney+ which still stutter and crash multiple times a day. Amazon for some reason treats every season of a TV show and HD/SD versions of movies as independent items in their library. I still haven't been able to download a HBO Max video for offline viewing on iOS without the app crashing on me at 99%.
The problems you mention with Netflix are real, but they have more to do with the business side of things. Netflix recommendations seem crap because they don't have a lot of third party content to recommend in the first place. Their front page layout is optimized to maximize repetition and make their library seem larger. They steer viewers to their own shows because that's what the business team wants. None of these are problems you can fix by reassigning engineers.
Wait, what? Netflix is the absolute worst at this. Every time I log in the interface is different! Netflix could not care less about users having a consistent seamless experience.
But as far as performance goes, I totally agree with you. The performance is impressively good and noticeably better than the other streaming apps I use.
The UX is just so bad in so many ways (UI churn, autoplay, useless ratings, useless categories, recaps that can be watched exactly once, and so on...) it mostly ruins the app for me. The actual video quality is great though.
- by the time HBO Max finish loading, I've already lost interest
- Amazon Prime constantly gives me errors, and it's often hard to find what you paid for and what you have to pay for
- Paramount+ often restart episode from beginning instead of resuming.
- Many leave shit in your queue with a few seconds left for you to "Continue Watching". I still have shows in Paramount+ that I've finished months ago in the queue, and there is no way to delete them without watching end credits. - HBO Max only allows you FF in small fixed intervals
- Plex...used to be okay, now it's pushing its streaming services and works very bad offline
- Apple TV has awful offline experience compared to netflix in terms of UX
Nah, I will take netflix constantly changing rows over shit others do.
The issue is the business polices surrounding it. The UI itself is user-hostile.
Their recommendations used to be world class until they got rid of the user feedback mechanisms.
I don't know how it is in the US, but in Japan, it's even worse than that. Japanese dub and English with Japanese subs are different items (although the UI has a way to choose between audio and subtitle channels), and Japanese subtitles are burned into the video.
Very often, despite all the config being in Japanese, it will complain in Japanese, when starting a video, that there isn't an English dub for the show, as if I had asked for it, which I never have.
In related info for Anime that is definitely not an English dub, it shows the names of the US voice actors and their related works more often than the Japanese voice actors.
Multi-season shows in the Prime Video app might be either separate items or not, depending on the title. Through the Fire TV home, though, it's a huge mess. You may have Season 1, 2, 3 under an item for $Show, with Season 2 being the second half of season 1, but on netflix and season 3 really being season 2, on netflix as well. At the same time, the real season 2 is also a separate item which has 2 seasons, with season 1 being really season 2 on Prime and season 2 being... season 2 on another service. Confused?
It's often super hard to get those non-season 1 items as search results too...
Things are however a little better if only using the Prime video app.
Which devices are you referring to? I’ve only used the PC and mobile interfaces both of which are quite pleasant.
Hulu did that big redesign, and it's extremely pretty to look at, but even after a few years of trying to use it, I still struggle to do anything other than "resume episode". Finding the previous episode, list episodes, etc is always an exercise in randomly clicking, swiping, long pressing, waiting for loading bars, etc.
One thing Netflix really got right as well: the "Watch It Again" section. So many times I want to rewatch the episode I just "finished" (because either my wife finished a show when I leave the room, the kids fell off the table, I fell asleep or wasn't paying attention, etc), and every other platform makes this extremely difficult to find.
Back to Hulu--the only way I know how is the search feature, which is a PITA with a remote.
It's just that they succeeded in streaming market with low competition and great success bring in lot of post facto justifications on how outrageously great Netflix tech infra is.
I mean it may be excellent for their purpose but to think their solution can be industry wide replicated seems not true to me.
In my personal experience lots of companies (admittedly all large companies, but many of which sell their services / software / hardware to smaller companies) have a use for serving hundreds of Gbps of static file traffic as cheaply as possible. And the slides for this talk seem exactly on the money (again from my experience slinging lots of static data to lots of users).
But it's like, every discussion today must end with something about the pay and head count of engineers.
Strictly speaking, the Internet was supposed to help some servers survive and continue working together despite some others being destroyed by a nuke. That is more-or-less the case today: we see how people use VPNs to route around censorship. Whether you were supposed to stream TikTok videos directly from the phones of their authors or through a centralized data hose - i'm not sure that was ever the grand idea.
Also "decentralized" and "monetize" don't go well together because innovation is stimulated by profit margins and rent-free decentralized solutions by definition have those margins equal to zero (otherwise the solution is not decentralized enough).
So it turns out decentralized multi routed is not a good solution for video streaming.
Good quality, barely any buffering.
The niche content may be too difficult for a "live" streaming experience.
And I'm not sure most people actually care if their home hardware is being used for whatever by the service they're using, or else there'd be pushback on electron apps from more than just HN.
The sense I always got from Netflix's P2P work was that it was heavily tied into the political battles wrt the BS arguments that Netflix should pay for peering with tier 2 ISPs. Did this work there continue much after that problem went quieter?
Works amazingly well for watching something front to back if your download speed is fast enough; you'd never know it wasn't being streamed. The hardest part is finding a good torrent for what you want to watch. Ironically the Netflix catalog is one of the most easily available to pirate since people rip it directly from web.
Umm, those solutions exist (from places like AWS and Azure) because Netflix was able to do it without them. The cloud platforms recognized that others would want to build their own streaming services, so they built video streaming offerings.
You have the cart in front of the horse. The out-of-the-box solutions of today don’t exist without Netflix (and YouTube) building a planet scale video solution first.
And your timeline is all wrong too. Netflix didn't even engage with the ISPs about bandwidth until long after moving out of our own datacenter. We started the OpenConnect program specifically to make it easier for ISPs, there was no bullying. The spat you're thinking of is that Comcast didn't want to adopt the OpenConnect but also didn't want to appropriately peer with other networks to give their customers the advertised speeds.
And hardware cost is a hugely relevant parameter. Being efficient with hardware is the difference between profitable streaming at that scale and not profitable.
Someone still has to do the R&D for edge cache? These slides are about Open Connect - their own edge cache solution that gets installed in partners racks (i.e. ISPs and Exchanges). Before things that Netflix and Nginx implemented in FreeBSD, hardware compute power was wasted on various things they discuss in slides.
Yes, you can throw money at the problem and buy more hardware.
How do you think this compares in complexity with real time distributed transactions that spans across several financial partners across the globe?
Like, from work, hosting postgres. At this point, I very much understand why a consultant once said - "You cannot make mistakes in a postgres 10GB or 100GB and a dozen transactions per second in size". And he's right, give it some hardware, don't touch knobs except for 1 or 2 and that's it. The average application accessing our postgres clusters is just too small to cause problems.
And then we have 2 postgres clusters with a dataset size of 1TB or 2TB peaking at like 300 - 400 transactions per second. That's not necessarily big or busy for what postgres can do, but it becomes noticeable that you have to do some things right at this point and some patterns just stop working hard.
And then there are people dealing with postgres instances 100 - 1000x bigger than this. And that's becoming tangibly awesome and frightening by now, using awesome in a more oldschool way there.
I'm sure there are many teams that could design such a network with nearly unlimited resources, but it is entirely different when you have profit margins.
Serving large content has been solved for decades already. It's much easier and reliable to serve from multiple sources, each at their maximum speed. Want more speed ? Add another source. Any client can be a source.
Netflix artificially restrains itself by only serving from their machines. It is a very nice engineering feat, but is completely artificial. As a user it feels weird to think of them highly when they could just have gone the easier road.
Just use Nginx and a backend lang of your choosing.
I wouldn't bother. Unless you use storage at the CDN - which is probably very not cost effective for you.
Bittorrent is built towards "offline" viewing. Try Peertube for a stack that is more built for streaming and has bittorrent sharing built-in (actually webtorrent, because the browser doesn't speak raw TCP or UDP, but the idea is the same)
We ended up switching to Fastly for CDN. There's something hidden here though that becomes a problem at Netflix size. We were willing to pay the cloud provider tax, and we didn't dig down into kernel level or storage optimizations because off the shelf was good enough. At Netflixes scale, that adds up to millions of extra server hours you have to pay for if you don't do the 5% optimizations outlined in the article.
The solution I'm talking about is bittorrent. The more people watch your content, the less your servers bear load. That is using the internet to its best potential instead of reverting back to the centralized model of the big shopping mall and its individual users.
They work like garbage in practice for mainstream users.
By creating this optimized system, it makes serving that much video profitable.
I don't think that's a popular sentiment about Netflix. Twitter, Reddit, Facebook, yes, but Netflix, YouTube, Zoom, not so much.
Is this claim based on some example I should know? Countless companies never achieve product/market fit, but very few I can think of fail because they weren't able to handle all their customers.
In a way it's the opposite. Things scale up by removing complex interactions. Software gets faster by solving the same problem with less steps.
Any beginner can write to much code and pile up too many software and hardware components.
It takes knowledge to simplify things.
I'd take it a step further, the beginner has no choice but to write too much code/cruft. By definition they lack the experience required to know the options before them, many of which are simpler.
Removing complexity and saving CPU cycles does not feel as good to beginners.
That’s one of the fundamental problems in todays Open Source world: everything’s being optimized for a very rare case.
It's the kind of thing that people run at home on a raspberry pi, docker container or linux server and it consumes almost no resources.
But at our organization this needs to scale up to millions of users in an extremely reliable way. It turns out this is incredibly hard and expensive and takes a team of people and a bucket of money to pull it off correctly.
When I tell people what I work on they only think about their tiny implementation of it, not the difficulty of doing it at an extreme scale.
Believe me or not, I was in a company doing web file streaming in 2009 using Nginx, sendfile and SSL offloading on the NIC.
It was installed by one dude. A standard Linux distro, standard kernel and no custom software. Just compile the SSL offloading kernel module once.
It's strange that you assumed otherwise.
I call them asinine because these kinds of architectures aren't written out in one day (let alone 60 minutes) and it's silly to not build an evolving setup (unless you are replacing a legacy system in a company with scale, but then again you have more than 60 minutes...). I wish system design interviews would not just think about if you know how to write high scale designs, but if you know how to build tiny designs that minimize footprint before scale and give maximum flexibility for scaling when its indicated ...
It's such a previous generation thing to be angry at the way that modern development is done. Of _course_ dev is now web heavy and people are pulling in all sorts of libraries to make their lives easier.
But pretending that those same devs wouldn't also be capable of development on lower/less abstracted levels (well sure, maybe not _all_ of them) is insulting.
- They are not using OS page cache or any memory caching for that, every request is served directly from disks. This seems possible only when requests are spread between may NVMe disks since single high-end NVMe like Micron 9300 PRO has max 3.5GB/s read speed (or 28Gbps) - far less than 800Gbps. Looks like it works ok for long-tail content but what about new hot content everybody wants to watch at the day of release? Do they spread the same content over multiple disks for this purpose?
- Async I/O resolves issues with nginx process stalling because of disk read operation but only after you've already opened the file. Depending on FS / number of files / other FS activities, directory structure opening the file can block for significant time and there is no async open() AFAIK. How they resolve that? Are we assuming i-node cache contains all i-nodes and open() time is insignificant? Or are they configuring nginx() with large open file cache?
- TLS for streamed media was necessary because browsers started to complain about non-TLS content. But that makes things sooo complicated as we see in the presentation (kTLS is 50% of CPU usage before moving to encryption offloaded by NIC). One has to remember that the content is most probably already encrypted (DRM), we just add another layer of encryption / authentication. TLS for media segments make so little sens IMO.
- When you relay on encryption or TCP offloading by NIC you are stuck with that is possible with your NIC. I guess no HTTP/3 over UDP or fancy congestion control optimization in TCP until the vendor somehow implements it in the hardware.
But do we really need such protection for a TV show?
"Metadata" in HLS / DASH is a separate HTTP request which can be served over HTTPS if you wish. Then it can refer to media segments served over HTTP (unless your browser / client doesn't like "mixed content").
DRM may be mandated by the content owners. TLS gives Netflix customers privacy against their ISP snooping what they're watching.
What you watch can be a very private thing, especially for famous people.
https://www.blackhat.com/docs/eu-16/materials/eu-16-Dubin-I-...
https://americansforbgu.org/hackers-can-see-what-youtube-vid...
I believe our TLS initiative was started before browsers started to complain, and was done to protect our customer's privacy.
We have lots of fancy congestion optimizations in TCP. We offload TLS to the NIC, *NOT* TCP.
There is a Netflix Tech Blog from a few years ago that talks about this better than I could: https://netflixtechblog.com/content-popularity-for-open-conn...
How is this possible? If TCP is done on the host and TLS on the NIC data will need to pass through the CPU right? But the slides show cpu fully bypassed for data
Modern NICs use packet descriptors that allow you to more or less say take N bytes from this address, then M bytes from some other address, etc to form the packet. So the kernel is going to make the tcp/ip header, and then tell the nic to send that with the next bytes of data (and mark it for TLS however that's done).
https://www.freebsd.org/cgi/man.cgi?query=sendfile&sektion=2
Given one can specify arbitrary offsets for sendfile(), it's not clear to me that there must be any kind of O(k > 1) relationship between open() and sendfile() calls: As long as you can map requested content to a sub-interval of a file, you can co-mingle the catalogue into an arbitrarily small number of files, or potentially even stream directly off raw block devices.
My own testing on single socket systems that look rather similar to the ones they are using suggests it is much easier to push many 100 Gbit interfaces to their maximum throughput without caching. If your working set fits in cache, that may be different. If you have a legit need for sixteen 14 TiB (15.36 TB) drives, you won't be able to fit that amount of RAM into the system. (Edit: I saw a response saying they do use the cache for the most popular content. They seem to explicitly choose what goes into cache, not allowing a bunch of random stuff to keep knocking the most important content out of cache. That makes perfect sense and is not inconsistent with my assertion that hoping a half TiB cache will do the right thing with 224 TiB of content.)
TLS is probably also to keep the cable company from snooping on the Netflix traffic, which would allow the cable company to more effectively market rival products and services. If there's a vulnerability in the decoders of encrypted media formats, putting the content in TLS prevents a MITM from exploiting that.
From the slides, you will see that they started working with Mellanox on this in 2016 and got the first capable hardware in 2020, with iterations since then. Maybe they see value in the engineering relationship to get the HW acceleration that they value into the hardware components they buy.
Disclaimer: I work for NVIDIA who bought Mellanox a while back. I have no inside knowledge of the NVIDIA/Netflix relationship.
The NUMA work they did, I remember being in a meeting with them as a Linux Developer at Intel at the time. They bought NVMe drives or were saying they were going to buy NVMe drives from Intel which got them access to "on the ground" kernel developers and CPU people from Intel. Instead of talking about NVMe they spent the entire meeting asking us about howt the Linux kernel handles NUMA and corner cases around memory and scheudling. If I recall correctly I think they asked if we could help them upstream BSD code for NVMe and NUMA. I think in that meeting there was even some L9 or super high up NUMA CPU guy from Hillsborough they some how convinced to join.
The conversation and technical discussion was quite fun, but it was sort of funny to us at the time they were having to do all this work on the BSD kernel that was solved years ago for linux.
Technical debt I guess.
If someone was going to do a similar comparison now the results could be different.
The funny thing is the rest of Netflix runs on Ubuntu, only those edge CDN runs on BSD.
I don’t think you can dismiss BSD faster than Linux (or make any claims about the relative speed of different OSes) just because big companies run Linux. There are other costs involved and optimisations that can be shared if your edge serving stack is as similar as possible to the non-edge serving stack (that you have many more engineers developing for).
All you can conclude is that with enough optimisation, Linux can be made to perform well enough for it to not be worth replacing (yet). Because replacing Linux would require replicating all the custom software and optimisations made to it for whatever other platform you pick.
Here's a Facebook engineering blog post about how they left NUMA behind. https://engineering.fb.com/2016/03/09/data-center-engineerin...
Well, not on Epyc generation 1. Those have four NUMA segments in each socket.
Also those Xeon Platinum 9200 processors Intel made as an attention grab.
We invested in NUMA back when Intel was the only game in town, and they refused to give enough IO and memory bandwidth per-socket to scale to 200Gb/s. Then AMD EPYC came along. And even though Naples was single-socket, you had to treat it as NUMA to get performance out of it. With Rome and Milan, you can run them in 1NPS mode and still get good performance, so NUMA is used mainly for forward looking performance testbeds.
They have 9 chips on what is essentially a tiny, high-density motherboard. Effectively they are 8-socket server boards that fit in the palm of your hand.
The dual-socket version is effectively a 16-socket motherboard with a complex topology configured in a hierarchy.
Take a look at some "core-to-core" latency diagrams. They're quite complex because of the various paths possible: https://www.anandtech.com/show/16214/amd-zen-3-ryzen-deep-di...
Intel is not immune from this either.Their higher core-count server processors have two internal ring-bus networks, with some cores "closer" to PCIe devices or certain memory buses: https://semiaccurate.com/2017/06/15/intel-talks-skylakes-mes...
Incidentally, I saw some of their job posts yesterday. If you think this presentation was cool, and you want to work with some competent yet humble colleagues, check these out:
CDN Site Reliability Engineer https://jobs.netflix.com/jobs/223403454
Senior Software Engineer - Low Latency Transport Design https://jobs.netflix.com/jobs/196504134
The client side team is hiring, too! (This is my old team.) Again, it's full of amazing people, fascinating problems, and huge impact:
Senior Software Engineer, Streaming Algorithms https://jobs.netflix.com/jobs/224538050
That last job post has a link to another very deep-dive tech talk showing the client side perspective.
I am not a Netflix subscriber but I dont think Netflix does live streaming for anything much if at all.
The Ultra Low Latency seems to suggest Netflix is exploring this idea. Which could be Live Sport or some other shows.
I’m on mobile and there does not seem to exist a direct link. Search for: “Case Study: Serving Netflix Video Traffic at 400Gb/s and Beyond”
But maybe DDR5 will come out by then and get this team busy again lol.
It is interesting that this work is happening on FreeBSD, and potentially with diverging implementations than Linux. Linux programs seem to be moving towards userspace getting more power, with things like io_uring and increasing use of frameworks like DPDK/SPDK. This work is all about getting userspace out of the way, with things like async sendfile and kernel TLS. That's pretty neat!
“~200GB/sec of memory bandwidth is needed to serve 800Gb/s” and “16x Intel Gen4 x4 14TB NVME”. So each NVMe drive would need to serve 12.5GB/s which is more than the 8GB/s limit for PCIe 4.0 x4. Also popular content would need to be on every drive, drastically lowering the total content stored.
Also see drewg’s comment on this for a different reason: https://news.ycombinator.com/item?id=32523509
1. https://www.servethehome.com/nvidia-connectx-7-shown-at-isc-...
2. https://wccftech.com/amd-epyc-7004-genoa-32-zen-4-core-cpu-s...
Simply following trends and doing what everyone else does leads to mediocre results and the assembly line type of work that most software development has become.
Although I would assume the CPU being the bottleneck with 8x the required processing power.
Hopefully that means Netflix now has the incentive to bump bitrate to 20Mbps or higher.
OpenBSD had support for them like 22 years ago.
https://www.google.com/search?client=firefox-b-d&q=SSL+accel...
Now we have TLS1.2/TLS1.3 offload getting built into the PCI-E 4.0 100/200/400GbE (whatever speed) NIC.
The Mellanox (and Chelsio) NICs do in-line crypto. They DMA down the network packets as they would have to do anyway in order to send them. Then they encrypt them on their way out onto the wire. This reduces memory bandwidth requirements by roughly 50%.
Have you examined other NIC vendors? (Chelsio?)
We looked at Chelsio (as T6 was available well before CX6-DX). However, the CX6-DX offers a killer feature not available on T6. The CX6-DX can remember the crypto state of any in-order stream, while the T6 cannot. That means that the TCP stack can send, say, 4K of a TLS record, wait for acks, and come back 40ms later and send the next 4K and DMA just the requested 4K from the host. The T6 cannot remember the state, and would need to DMA the first 4K (which was already sent) in order to re-establish the crypto state, and then DMA the requested 4K. This could run the PCIe bus out of bandwidth. The alternative is to make TCP always chunk sends at the TLS record size, but this was horrible for streaming quality.
As much as I enjoy the results of the work, I'm always a bit curious how the sausage is made. Is pushing the hardware limits your primary job or something you do periodically? How do you go about selecting the gear you use? How much do you work with the vendors? (etc etc) I'd really enjoy a behind the scenes blog post or something wrt this serving absurd amounts of traffic from a single box.
But I do plenty of other things as well, including fixing random kernel bugs. You can read the git log of the FreeBSD main branch to see some of the things I've been working on..
This part I don't get. How about DRM? Unless Netflix pre-DRM all contents for all user?
Netflix's DRM is sufficiently good that the Reddit Piracy subreddit has spent the last three months moaning that they have no access to 4K Netflix rips, at least for weeks or months after the content comes out.
Netflix's DRM and key management systems do what they care about pretty well at this point, which is protect the initial airing of popular shows.
It only really matters that this key is unique per package, not per user, because once even a single user can compromise the trusted execution environment and extract either the key or the plain video stream, that piece of content is now pirated anyway. So, key reuse against the same content probably isn't really a major part of the threat model - this attacker could share the key with others, but they might as well share the decrypted content instead.
I totally get the cost, convenience, and supply chain risk-value in commodity stuff that you can just go out and buy, but once you're bound to a single network card, this advantage starts to go away, and it seems like you're fighting with the entire system topology when it comes to NUMA, no? Why not a "TCP file send accelerator" instead of a whole computer?
The bottleneck is at your NIC anyways, so seems like there would be a market for NIC that can directly read from disk into NIC's working memory
Ah, so this is why everything stutters / falls apart when you switch subtitles on or off -- it has to access a whole different file and resume at the same place in that file I assume? I would think you would want the (verbal) audio separated out in a different file so it can be swapped out on the fly without re-initializing the video stream, and same thing with subtitle files? I'm just making some assumptions based on the behavior I've seen but would be cool to know how this works.
I've never seen this bad behavior myself. Do you mind sharing the client you're using?
Future technology advances increasingly looks like this complex work integrating hardware, OS fixes, team collaboration. People and teams and companies working together, and contributing to shared resources like FreeBSD. Tolerating mistakes at scale, giving credit where credit is due, and all the other things that make respect real, which creates a space to get things done.
Most of us will never get close to these opportunities or contexts, but still it helps us advance our own technique/culture to observe and model your story. And perhaps you'll help new collaborators find you. All the best.
In 2021 somebody submitted a patch for io_uring support in nginx:
https://mailman.nginx.org/pipermail/nginx-devel/2021-Februar...
I'm not sure if there has been further progress on it so far. In one comment feedback is "it doesn't seem to make the typical nginx use case much faster" [at that time].
But I find this interesting, because io_uring can make almost all things async that can't be used async so far in Linux (open(), stat(), etc) and thus in nginx.
Would io_uring integration in nginx be relevant for you?
* https://papers.freebsd.org/2021/eurobsdcon/gallatin-netflix-...
The reason companies like Sony and Apple pick FreeBSD is because they get an open source POSIX-compliant OS they can drastically modify down to the kernel level without having to open source their modifications.
Who knows what plans Netflix had in mind when this decision was made? Maybe they simply commit to FreeBSD because they do not want to release their code under the GPL. I certainly never do, and I will never ever make any improvements to any GPL code I use. I think it is an awful license. Maybe Netflix does, too? That would be directly relevant to the license of the OS even if they never release their customized OS, if they have a customized version at all.
FreeBSD (and other BSDs) are also a much simpler platform with orders of magnitude less "churn" taking place - if you are doing an engineering experiment that involves low-level integration with OS-level or kernel-level subsystems (eg kernel TLS f.ex) then it's a lot simpler to start from BSD, and you can very safely assume that your code will pretty much continue working in future releases as well. It's a far more stable platform to build a product around, Linux is "constantly on fire all the time" in comparison to the absolutely sedate pace of BSD development. It was built over the last 50 years (the oldest code in the FreeBSD repos dates to 1979) and it mostly Just Works Like You'd Expect. The good old 'cathedral vs the bazaar' tradeoff. Well, if you are building an annex, at least you know the cathedral isn't going to collapse under you next week, and by the time you get your project finished the bazaar might have moved on to something completely different and abandoned the thing you needed. But you don't get docker and the other New Things either.
Sony likes it because it allows them to release proprietary software off an open codebase. The conceptual divide is - GPL protects end user freedoms, at the cost of developer freedoms (eg linking to code with GPL-incompatible licenses, or not distributing the source to proprietary extensions to a GPL'd product). While BSD/MIT license it's the other way around, they're "here's this code, you can use it however you want", and sometimes that means doing things that users might consider "evil".
It makes me wonder what the next hardware revolution will be. It seems like most resource intense applications are bottlenecking at transferring memory. UE5's nanite tech hinges on the ability to transfer memory directly from disk to GPU, Netflix built specific hardware to avoid copying memory between userspace and hardware, and I wonder how much other performance we're missing out on because we can't transfer memory fast enough.
How much faster could AI training be if we could get memory directly from disk to the GPU and avoid the CPU orchestrating it all? What about video streaming? I have a feeling these processes already use some clever tricks to avoid unnecessary trips through the CPU, but it will be interesting to see which direction hardware goes with this in mind.
1: https://developer.nvidia.com/gpudirect
2: https://www.servethehome.com/what-is-a-dpu-a-data-processing...
In synchronous IO paradigms, this is all managed by the application with whatever logic the app author wants to implement. You can report the errors, ping monitoring, whatever.
But with this async thing, what's the API for that? Do you have to write kernel code that lives above the driver to implement devops logic? How would one even document that?
+1 for the technical wizardry, but seems like it's going to be a long road from here to an "OS" feature that can be documented.
If there is a connection RST, then the data buffered in the kernel is released (either freed immediately, or put into the page cache, depending on SF_NOCACHE).
sendfile_iodone() is called for completion. If there is no error, it marks the mbufs on the socket buffer holding the pages that were recently brought in as ready, and pokes the TCP stack to send them. If there was an error, it calls TCP's pru_abort() function to tear down the connection and release what's sitting on the socket buffer. See https://github.com/freebsd/freebsd-src/blob/main/sys/kern/ke...
But I believe that async is an anti-pattern. From the article:
* When an nginx worker is blocked, it cannot service other requests
* Solutions to prevent nginx from blocking like aio or thread pools scale poorly
Nothing against nginx (I use it all the time, it's great) but I probably would have used a synchronous blocking approach. The bottleneck there would be artificial limits on stuff like I/O and the number of available sockets or processes.So.. why isn't anyone addressing these contrived limits of sync blocking I/O at a fundamental level? We pretend that context switching overhead is real, but it's not. It's an artifact of poorly written kernels from 30+ years ago (especially in Windows) where too many registers and too much thread state must be saved while swapping threads. We're basically all working around the fact that the big players have traditionally dragged their feet on refactoring that latency.
And that some of the more performant approaches like atomic operations using compare and swap (CAS) on thread-safe queues beat locks/mutexes/semaphores. And that content-addressable memory with multiple busses or even network storage beats vertical scaling optimizations.
So I dunno, once again this feels like kind of a drink-the-kool-aid article. If we had a better sync blocking foundation, then a simple blocking shell script could serve video and this whole PDF basically goes away. Rinse, repeat with most web programming too, where miles-long async code becomes a single deterministic blocking function that anyone can understand.
I'm kind of reaching the point where I expect more from big companies to fix the actual root causes that force these async workarounds. I kind of gave up on stuff like that over the last 10 years, so am behind the times on improvements to sync blocking kernel code. I'd love to hear if anyone knows of an OS that excels at that.
https://www.slideshare.net/facepalmtarbz2/new-sendfile-in-en...
> but I probably would have used a synchronous blocking approach.
Well, send a patch, then.
Then Varnish is probably more your style. (A discussion between phk and drewg would be fascinating to watch.)
We pretend that context switching overhead is real, but it's not.
This sounds crackpot to be honest. Linux has put a lot of effort into optimizing context switching (that's why they have NPTL instead of M:N) and I assume FreeBSD has as well.
...this whole PDF basically goes away
Sync vs. async doesn't solve any of the NUMA or TLS issues that this whole PDF is about.
You can deal with that by having lots of threads or processes doing the work, but coordinating many threads can be difficult. With async sendfile, you let the kernel manage all of this without needing userspace to do anything, and the kernel is well organized to manage pushing data to the right place when disk i/o completes.
If you just want to write code in blocking style, check out Erlang, but know that under the hood it massages most of the blocking into aggregated select/kqueue/epoll/etc calls.
And here are the slides explaining it: https://www.slideshare.net/facepalmtarbz2/new-sendfile-in-en...
There are video of various talks by Gleb Smirnoff explaining all this magic on YouTube.
The feature is fully documented in `man 2 sendfile`, it was part of the patch that did the work.
I know they use a number of techniques like kernel bypass to get the lowest latency possible, but maybe they have explored some solution to this problem as well.
never be afraid to admit that you don’t know something. guessing wrong is a much worse look than not answering at all.
No exchange should allow trades to complete, I would argue in any time less than 15 minutes, and each trade should have a random 1-15 minute delay pad on top of that.
The HFT access only serve the larger financial firms, and are used to do frontloading and other basically-illegal tricks. It provides anti-competitive advantage to large firms for markets that are supposed to open access/fair trading. And of course it leads to AI autotrading madness.
I get that it keeps a lot of tech people very well compensated, but it is either in the service of unregulated fraud at worst and unfair advantage at best.
I would highly recommend you to read this book ( https://press.princeton.edu/books/hardcover/9780691211381/tr... ) or if it's too long form, read this article ( https://www.thediff.co/p/jane-street )
I am not saying HFT's are in charitable business but they do serve an important role in the financial markets.
Additionally, the National Electric Code in the US specifies that continuous load should not exceed 80% of given circuit/breaker capacity.
So with dual 800 power supplies at "max" 80% load that's "only" 640 watts for one of these 2U servers. For 208V power that's only 3 amps. High density (for sure) compared to the old days but not as ridiculous as it may seem.
- 16 x 25W SSDs
- 2 x 225W CPUs
- On top of that, add RAM, cooling, etc.
Honestly, it's still manageable. I doubt they'd put 10 of those in a single rack (you'd need an ISP that would want serve 2.2M subscribers in peak from a single location, not necessarily desirable on their side); but if the site is getting full, you'd feel the (power) pressure (slowly) mounting.
They absolutely wouldn't want to concentrate them. The entire purpose is to reduce ISP network load and get as close to the customer eyeballs as possible. I don't have any experience with these but I imagine ISPs would install them at their peering "hubs" in major cities - in my experience the usual suspects like Chicago, NYC, Miami, etc.
[0] - https://openconnect.zendesk.com/hc/en-us/articles/3600345383...
Do we know if the rates that these hosts serve actually make it into production? Or do they derate the amount they serve from a single host and add others?
I regret that I've crashed boxes doing hundreds of Gb/s. Thankfully our stack is resilient enough that customers didn't notice.
https://news.ycombinator.com/item?id=28584738
This was the video which was posted back then alongside the slides: https://www.youtube.com/watch?v=_o-HcG8QxPc
In the case of Netflix it is in the ISPs best interest to let them push as much traffic to their customer eyeballs as possible. After all, it's much "easier" and cheaper to build out your internal fabric and network (which you have to do anyway for the traffic) than it is to buy and/or acquire transit to "the internet" for this level of traffic.
[1]: https://xrdocs.io/8000/tutorials/8201-architecture-performan...
I stand corrected on "always line rate all the time in any circumstance" but by your math and my general point < 1 Tbps from one of these appliances across multiple 100G ports isn't problematic in the least from a hardware standpoint - especially for the Netflix traffic pattern with relatively full (if not max MTU) packets.
The slide deck background though: At least half of the products in the slide deck template are no longer on Netflix...
Besides that, it sees like all of the heavy lifting here is done by Mellanox hardware...
We chose to use FreeBSD, and have contributed our code back to the FreeBSD upstream to make the world a better place for everyone.
I don't get that with Netflix, I've occasionally had it crash out 'sorry this could not be played right now' (which is a weird bug itself - because it always loads fast & fine when I immediately press play on it again) but never such slow loading or pausing.
It's a win-win. ISPs don't have to use their peering and/or transit bandwidth to upstream peers and users get a much better experience with lower latency, higher reliability, less opportunity for packet loss, etc.
Edit: Whoops, apparently this tab has been open for four hours and of course someone already had responded to you, lol.
In the US, the fastest ISP for Netflix usage seems to be Comcast (https://ispspeedindex.netflix.net/country/us ), with an average speed of 3.6Mbps. That would serve an average of 222k simultaneous customers on a single server.
Surely they transcode on some server? Maybe they just mean they don't do it on the same server that is serving bits to customers?
Switch from TCP to UDP protocol for even faster video delivery send speed?
Why use TCP for video content? Packet loss but then one would use http3 like UDP tech?
Have a wonderful day. =)
I've never wanted to build a single box pulling that much traffic when 20 would do.
It is better for me to pirate their content, play it with plex and be happy. I pay for Netflix and it is absurd.
I think the best years are over for Netflix. The hard awakening is here to make content that the users want and they are a movie/tv content company, not primarily a „tech company“.
Netflix has all the bandwidth data and metrics, but this is not working since ages. Maybe a more basic setup on their end would bring better results. Focus more and delivery, not 10 different UI versions, AB Tests, Batch Job Workflows and so on. They Post on their engineering blog how they test multiple TVs, multiple Profiles in encoding, great things, but if the basics don’t work ... well what is it good for.
I think they lost their focus.
That's largely what many staff+ engineers have to do, even in otherwise healthy organizations. "Staff" isn't a glorified, autonomous, and stress-free version of senior at most companies. There's nothing wrong with staying at the senior level indefinitely provided (1) the pay and other factors are keeping up with your contributions and (2) the staff+ and management folks are being effective umbrellas for the politics and messy uninteresting details behind interesting problems like this.
Plus, not every workload is: read from disk -> encrypt -> send to nic