We built a self-healing system to survive a concurrency bug at Netflix
pushtoprod.substack.com
pushtoprod.substack.com
30 years ago or so I worked at a tiny networking company where several coworkers came from a small company (call it C) that made AppleTalk routers. They recounted being puzzled that their competitor (company S) had a reputation for having a rock-solid product, but when they got it into the lab they found their competitor's product crashed maybe 10 times more often than their own.
It turned out that the competing device could reboot faster than the end-to-end connection timeout in the higher-level protocol, so in practice failures were invisible. Their router, on the other hand, took long enough to reboot that your print job or file server copy would fail. It was as simple as that, and in practice the other product was rock-solid and theirs wasn't.
(This is a fairly accurate summary of what I was told, but there's a chance my coworkers were totally wrong. The conclusion still stands, I think - fast restarts can save your ass.)
Also the story isn't that they couldn't just that they were measuring the actual failure rate not the effective failure rate because the device could recover faster than the failure caused actual issues.
Each running process had a backup on another blade in the chassis. All internal state was replicated. And the process was written in a crash only fashion, anything unexpected happened and the process would just minicore and exit.
One day I think I noticed that we had over a hundred thousand crashes in the previous 24 hours, but no one complained and we just sent over the minicores to the devs and got them fixed. In theory some users would be impacted that were triggering the crashes, their devices might have a glitch and need to re-associate with the network, but the crashes caused no widespread impacts in that case.
To this day I'm a fan of crash only software as a philosophy, even though I haven't had the opportunity to implement it in the software I work on.
Our half-day workaround implementation was the same thing, just cycle the cluster regularly automatically.
Since we're running on AWS, we just double the size of the cluster, wait for the instances to initialize, then rapidly decommission the old instances. Every 2 hours.
It's shockingly stable. So much so that resolving the root cause isn't considered a priority and so we've had this running for months.
I've seen multiple issues solved like this after engineering teams have been cut to the bone.
If the cost of maintaining enough engineers to keep systems stable for more than 24 hours, is more than the cost of doubling the container count, then this is what happens
x = time it takes to switchover
y = length of the cycles
x/y = % increase in cost
For us, it's 15 minutes / 120 minutes = 12.5% increase, which was deemed acceptable enough for a small service.
There, formalized the approach, so you can't call it terrible anymore.
If you are looking for an apt metaphor, Stalin sort might be more in line with what's going on here. Or maybe "ostrich algorithm".
Some collectors may need to do this, but there are several collectors that don't. EpsilonGC is a prime example of a GC that doesen't need to prove anything
I mean, I interpret your comment to be a joke, but you could've made it a bit more obvious for people not familiar with the latest fancy in Java world.
Pragmatically, restart the service periodically and spend your time on more pressing matters.
On the other hand, we fully understand the reason for the fault, but we don't know exactly where the fault is. And it is, our fault. It takes a certain kind of discipline to say "there are many things I understand but don't have the time to master now, let's leave it."
It's, mostly, embarrassing.
Then you realize it's a paper idol and the best you can do is suck less than the average.
Thanks for playing Wing Commander!
captain america voice I got that reference :-)
Not OP but this is a somewhat normal case of making a tradeoff? They aren't able to repair it at the moment (or rather don't want/can't allocate the time for it) and instead trade their ressource usage for stability and technical debt.
but, as a discipline, engineers manage to encourage the ascent of the least engineer-ly (or, perhaps, "hacker"-ly) among them ("us") ...-selves... through their sui generis combination of learned helplessness, willful ignorance, incorrigible myopia, innate naïvete, and cynical self-servitude that signify the Institutional (Software) Engineer. coddled more than any other specialty within "the enterprise", they manage to simultaneously underplay their hand with respect to True Leverage (read: "Power") and overplay their hand with respect to complices of superiority. i am ashamed and dismayed to recall the numerous times i have heard (and heard of) comments to the effect of "my time is too expensive for this meeting" in the workplace... every single one of which has come not from the managerial class-- as one might reasonably, if superficially, expect-- but from the software engineer rank and file.
to be clear: i don't think it's fair to expect high-minded idealism from anyone. but if you are looking for the archetypical "company person"... engineers need look no further than their fellow podmates / slack-room-mates / etc. and thus no one should be surprised to see the state of the world we all collectively hath wrought.
I don't know why my senses tell me that this is wrong even if you can afford it
The fix is also hiding other issues that show up. So it degrades over time and eventually you’re stuck trying to solve multiple problems at the same time.
As a Director of Engineering at my last startup, I had an "all hands on deck" policy as soon as any concurrency bug was spotted. You do NOT want to let those fester. They are nondeterministic, infrequent, and exponentially dangerous as more and more appear and are swept under the rug via "reset-to-known-good" mitigations.
You should do proper risk assessment, such bug may be leveraged by an attacker, that may actually be a symptom of a running attack. That may also lead to data corruption or exposure. That may mean some part of the system are poorly optimised and over-consuming resources, maybe impacting user-experience. With a dirty workaround, your technical debt increases, expect more and more random issues that requires aggressive "self-healing".
You don’t need wizards in your team anymore.
Something seems off in the instance? Just nuke it and spin up a new one. Let the system debugging for the Amazon folks.
If you do look into the Oracle dba handbook, scheduled index rebuilds are somewhat recommended. We do it on weekends on our Oracle instances. Otherwise you will encounter severe performance degredation in tables where data is inserted and deleted at high throughput thus leading to fragmented indexes. And since Oracle 12g with ONLINE REBUILD this is no problem anymore even at peak hours.
I mean, you probably know this, but sooner or later this attitude is going to come back to bite you. What happens when you need to do it every hour? Every ten minutes? Every 30 seconds?
This sort of solution is really only suitable for use as short-term life-support; unless you understand exactly what is happening (but for some reason have chosen not to fix it), it's very, very dangerous.
Once it's happening every 30 seconds, then they have up to 120 opportunities per hour, and it'll be fixed that much quicker!
"This sparked and interesting memory for me. I was once working with a customer who was producing on-board software for a missile. In my analysis of the code, I pointed out that they had a number of problems with storage leaks. Imagine my surprise when the customers chief software engineer said "Of course it leaks"
He went on to point out that they had calculated the amount of memory the application would leak in the total possible flight time for the missile and then doubled that number. They added this much additional memory to the hardware to "support" the leaks. Since the missile will explode when it hits it's target or at the end of it's flight, the ultimate in garbage collection is performed without programmer intervention."
What always bothers me, is when (note, I'm not saying this is the case for the grandparent comment, but it's implied) people don't understand what exactly is broken, but just reboot every so often to fix things. :0
For a lot of bugs, there's often the component you see (like the obvious resource leak) combined with subtle problems you don't see (data corruption, perhaps?) and you won't really know until the problem is tracked down.
The trick is to not tell your manager that your bandaid works so well, but that it barely keeps the system alive and you need to introduce a proper fix. Been doing this for the last 10 years and we got our system so stable that I haven't had a midnight call in the last two years.
The problem is that you merely borrowed yourself some time. As time goes on, more inefficiencies/bugs of this nature will creep in unnoticed, some will perhaps silently corrupt data before it is noticed (!), and it will be vastly more difficult at that point to troubleshoot 10 bugs of varying degrees of severity and frequency all happening at the same time causing you to have to reboot said servers at faster and faster intervals which simultaneously makes it harder to diagnose them individually.
> It's shockingly stable.
Well of course it is. You're "turning it off and then on again," the classic way to return to a known-good state. It is not a root-cause fix though, it is a band-aid.
Also, great point about (depending on your architecture) losing the ability to do things like cache results
You also want some machines to reboot much more frequently than others, so you catch boot issues before they affect your entire fleet.
Also, sometimes the management API goes out due to a bug/networking issue/thundering herd
i.e.
if bad(t) = fraction of bad instances at time t
and
bad(0) = 0
then
d(bad(t))/dt = -0.05 * bad(t) + 0.01 * (1 - bad(t))
so
bad(t) = 0.166667 - 0.166667 e^(-0.06 t)
Which looks a mighty lot like the graph of bad instances in the blog post.
> We created a rule in our central monitoring and alerting system to randomly kill a few instances every 15 minutes. Every killed instance would be replaced with a healthy, fresh one.
It doesn't look like they worked out the numbers ahead of the time.
That Netflix had already built a self-healing system means they were able to handle a memory leak by killing random servers faster than memory was leaking.
This post isn't about how they've managed that, it's just showing off that their existing system is robust enough that you can do hacks like this to it.
I don't remember all the details but I've still not be able to find the bug.
But this being in Elixir I "fixed it" with Task, TaskSupervisor and try/catch/rescue.
Not really a win but it is still running fine to this day.
[0] https://github.com/conradfr/ProgRadio/blob/1fa12ca73a40aedb9...
If people want to belittle something, either we aren't trying to solve the same problem (sure) or they're actively turning people away from what could be a serious advantage (more for me!)
If the cost of switching wasn't so high, I'd love to write Elixir all day. It's a joy.
I have long advocated randomly restarting things with different thresholds partly for reasons like this* and to ensure people are not complacent wrt architecture choices. The resistance, which you can see elsewhere here, is huge, but at scale it will happen regardless of how clever you try to be. (A lesson from the erlang people that is often overlooked).
* Many moons ago I worked on a video player which had a low level resource leak in some decoder dependency. Luckily the leak was attached to the process, so it was a simple matter of cycling the process every 5 minutes and seamlessly attaching a new one. That just kept going for months on end, and eventually the dependency vendor fixed the leak, but many years later.
Likely replacing HashMap with CHM would not solve the concurrency issue either, but it'd prevent an infinite loop. (Edit) It appear that part is just wrong: "some calls to ConcurrentHashMap.get() seemed to be running infinitely." <-- it's possible to happen no hashmap during concurrent put(s), but not to ConcurrentHashMap
https://netflixtechblog.com/the-netflix-simian-army-16e57fba...
Overall such a mistake alone undermines the effort/article.
(also worth noting this post seems to be discussing an event that occurred many years ago, circa 2011, so might not be a reflection of where they are today)
I can appreciate the hack to deal with this (I actually came up with the same solution in my head as reading) but if you cannot rollback and you cannot roll forward you are stuck in a special purgatory of CD hell that you should be spending every moment of time getting out of before doing anything else.
How was he managing the instances? Was he using kubernetes, or did he write some script to manage the auto terminating of the instances?
It would also be nice to know why:
1. Killing was quicker than restarting. Perhaps because of the business logic built into the java application?
2. Killing was safe. How was the system architectured so that the requests weren't dropped altogether.
EDIT: formatting
Kubernetes launched in 2014, if memory serves, and it took a bit before widespread adoption, so I’m guessing this was some internal solution.
This was a great read, and harkens back to the days of managing 1000s of cores on bare metal!
1. Killing was quicker than restarting.
If you happen to restart one of the instances that was hanging in the infinite thread, you can wait a very long time until the Java container actually decides to kill itself because it did not finish its graceful shutdown within the alotted timeout period. Some Java containers have a default of 300s for this. In this circumstance kill -9 is faster by a lot ;)
Also we had circumstances where the affected Java container did not stop even if the timeout was reached because the misbehaving thread did consume the whole cpu and none was left for the supervisor thread. Then you can only kill the host process of the JVM.
So improving uptime involves holding out a set of GPUs to swap out failed ones while they reboot. But also the whole run can just randomly deadlock, so you might solve that by listening to the logs and restarting after a certain amount of inactivity. And you have to be clever with how to save/load checkpoints, since that can start to become a huge bottleneck.
After many layers of self healing, we managed to take a vacation for a few days without any calls :)
Is software deployed regularly on this cluster? Does that deployment happen faster than the rate at which they were losing CPUs? Why not just periodically force a deployment, given it's a repeated process that probably already happens frequently.
What happens to the clients trying to connect to the stuck instances? Did they just get stuck/timeout? Would it have been better to have more targeted terminations/full terminations instead?
But as time goes by I just ask, all this work and costs and complexity, to serve files? Yeah don't get me wrong, the size of the files are really big, AND they are streamed, noted. But it's not the programming complexity challenge that one would expect, almost all of the complexity seems to stem from metadata like when users stop watching, and how to recommend them titles to keep them hooked, and when to cut the titles and autoplay the next video to make them addicted to binge watching.
Case in point, the blogpost speaks of a CPU concurrency bug and clients being servers? But never once refers to an actual business domain purpose. Like are these servers even loading video content? My bet is they are more on the optimizing engagement side of things. And I make this bet knowing that these are servers with high video-like load, but I'm confident that these guys are juggling 10TB/s of mouse metadata into some ML system more than I'm confident that they have some problem with the core of their technology which has worked since launch.
As I say this, I know I'm probably wrong, surely the production issues are cause by high peak loads like a new chapter of the latest series or whatever.
I'm all over the place, I just don't like netflix is what I'm saying
You could say the same thing about the entire web.
> When you're watching a movie your progress is constantly updated, posting data
This can be implemented on server side and with read requests only.
A proper comparison would be YouTube where people upload videos and comment stuff in real-time.
Even in this one sentence you're conflating two types of interaction. Surely downloading videos is yet a third, and possibly the rest of the assets on the site a fourth.
Why not just say the exact problem you think is worth of discussion with your full chest if you so clearly have one in mind?
-the entropy of the data: a video is orders of magnitude than browsing metadata.
- the compute required: other than an ML algorithm optimizing for engagement, there's no computationally intensive business domain work (throughput related challenges dont count)
- finally programming complexity, in terms of business domain, is not there.
I mean my main argument is that a video provider is a simple business requirement. Sure you can make something simple at huge scale and that is a challenge. Granted.
In this case it's not the same whether your server sends a 10 second packet and the viewer views all of it, and whether server sends a 10 second packet, but client pauses at the 5s mark (which needs client-side logic)
Might sound trivial, but at netflix scale there's guaranteed a developer dedicated to that, probably a team, and maybe even a department.
Of course, how much, depends on the service. Particularly, how much concurrent writing is happening, and do you need to update this state globally, in real-time as result of this writing. Also, is local caching happening and do you need to invalidate the cache as well as a result of this writing.
The most of the relevant problems disappear, if you can just replicate most of the data without worrying that someone is updating it, and you also don't have cache invalidation issues. No race conditions. No real-time replication issues.
Database-driven traffic is still a tiny percentage of internet traffic. It's harder to tell these days with encryption but on any given page-load on any project I've worked on, most of the traffic is in assets, not application data.
Now, latency might be a different issue, but it seems ridiculous to me to consider "downloading a file" to be a niche concern—it's just that most people offload that concern to other people.
Yet you have to design the whole infrastructure to note that tiny margin to work flawlessly, because otherwise the service usually is not driving its purpose.
Read-only assets are the easy part, which was my original claim.
I don't think this is true at all given the volume. With that kind of scale everything is hard. It's just a different sort of hard than contended resources. Hell, even that is as "easy" these days with CRDTs (and I say this with dripping sarcasm).
I'm honestly much more impressed by free apps like youtube and tiktok in terms of throughput, they have MUCH more traffic since users don't pay!
IMHO, a large amount of the complexity is all the other stuff. Account information, browsing movies, recommendations, viewed/not/how much seen, steering to local CDN nodes, DRM stuff, etc.
The file servers have a lot less complexity; copy content to CDN nodes, send the client to the right node for the content, serve 400Gbps+ per node. Probably some really interesting stuff for their real time streams (but I haven't seen a blog/presentation on those)
Transcoding is probably interesting too. Managing job queues isn't new, but there's probably some fun stuff around cost effectiveness.
They've also contributed significantly to open source tools for video processing, one of the biggest things that stands out is probably their VMAF tool for quantifying perceptual quality in video. It's probably the best open source tool for measuring video quality out there right now.
It's also absolutely true that in any streaming service, the orchestration, account management, billing and catalogue components are waaaay more complex than actually delivering video on-demand. To counter one thing you've said: mouse movement... most viewing of premium content isn't done on web or even mobile devices. Most viewing time of paid content is done on a TV, where you're not measuring focus. But that's just a piece of trivia.
As you said, you just don't like them, but they've done a lot for the open source community and that should be understood.
That said, free apps like tiktok and youtube probably face higher throughput, so the user-pays model probably means netflix is at the state of the art at high volume quality (both app experience and content) rather than sheer volume low quality or premium quality low volume markets.
I mean serving millions of customers at 8 bucks per month. Which is not quite like serving billions.
For example: calculate hash code of string, determine index, find apparent match, hashmap is modified by another thread, return value at that index which no longer matches.
I don't think that particular issue can happen with Java's HashMap, but there's probably some sort of similar goofiness.
> Rolling back was cumbersome
It's a fundamental principle of modern DevOps practice that rollbacks should be quick and easy, done immediately when you notice a production regression, and ideally automated. And at Netflix's scale, one would have wanted this rollout to be done in waves to minimize risk.
Apparently this happened back in 2021. Did the team investigate later why you couldn't do this, and address it?
Then DevOps principles are in conflict with reality.
It really is a fundamental advantage against the worst kinds of this category of bug
Of course they did. And whoever though "Concurrent" meant it would work fine gets burned by it. Of course.
And of course it doesn't work properly or intuitively for some very stupid reason. Sigh
Ingeniously simple solution for this particular bug though.
*as I recall, it had to do with merging a regular Hash in the ENV with a HashWithIndifferentAccess, which as it turns out was ill-conceived at the time and had undefined corner cases (example: what should happen when you merge a regular Hash containing either a string or symbol key (or both) into a HashWithIndifferentAccess containing the same key but internally only represented as a string? Which takes precedence was undefined at the time.)
Had a similar problem but memory wise with a pesky memory leak, and the short term solution was to do nothing as instances would to do nothing.
At first I thought maybe we should add a "hack" to cycle all the pods over 24 hours old, but then I wondered if making holiday freezes behave like normal weeks was really a hack at all or just reasonable predictability.
In the end folks managed to fix the leak and we didn't resolve the philosophical question though.
To me that is blocker to my thinking. I really need to understand the impact of leaving something behind before continuing. I’d likely do everything in my power to remove that unknown from current state for the sake of sanity.
It's not explained why they couldn't write a monitor script instead to find servers having the issue and only killing those.
The main takeaway from this post. You will not encounter exactly this or similar issues at your workplace, but this single piece of advice can help you fix any issues that you do encounter.
iirc it was mostly folks with websocket issues, but fixing the upstream was harder
10 years later and specific software has gotten better, but this type of problem is certainly still prevalent!
1 - https://www.theregister.com/2020/04/02/boeing_787_power_cycl...
(not that i don't get the sarcasm)
For Boeing, it's probably something fairly simple actually, but they don't want to fix it because their software has to go through a strict development process based on requirements and needing certification and testing, so fixing even a trivial bug is extremely time-consuming and expensive, so it's easier to just put a directive in the manual saying the equipment needs to be power-cycled every so often and let the users deal with it. The OP isn't dealing with this kind of situation.
If you don't know why you should reboot servers/services properly instead of terminating them..
These companies have achieved vast scale because correctness doesn’t matter that much so long as it is “good enough” for a large enough statistical population, and their Devops practices and coding practices have evolved with this as a key factor.
It is not uncommon at all for Netflix or Hulu or Facebook or Instagram to throw an error or do something bone headed. When it happens you shrug and try again.
Now imagine if this was applied to credit card payments systems, or your ATM network, or similar. The reality of course is that some financial systems do operate this way, but it’s recognized as a problem and usually gets on people’s radar to fix as failed transaction rates creep up and it starts costing money directly or clients.
“Just randomly kill shit” is perfectly fine in the Netflix world. In other domains, not so much (but again it can and will be used as an emergency measure!).
Memory leaks are often "resolved" this way... until time allows for a proper fix.
I put a CloudWatch alarm at 90% CPU usage which would trigger a reboot (which completed way before anyone would notice a downtime).
Never had issues again.
> I’m not a real programmer. I throw together things until it works then I move on. The real programmers will say “Yeah it works but you’re leaking memory everywhere. Perhaps we should fix that.” I’ll just restart Apache every 10 requests.