Netflix's Distributed Counter Abstraction
netflixtechblog.com
netflixtechblog.com
Sources:
Netflix doesn't care though because they just want people watching the newest thing they shovel at us. They go out of their way to hide a lot of older content from users to keep people's attention on the new stuff. Half of the categories they show you are just new ways to show you the latest things they're pushing (recently added, trending, top 10, new on netflix, your next watch, top picks for you, we think you'll love these, etc.) and every other category just repeats the same shows.
Who cares if some show leaves you disappointed because it got you hooked but then was canceled after 3 weeks, Netflix will just push some other newer show to keep you watching until they cancel that too.
It's also possible that savings on infrastructure costs, reduced engineering hours spent on debugging, etc give them a quiet competitive advantage
Most engineering is not visible to the end-user
And the best part... you chat with the support and they try to force a how to use our webpage shit on you when you want to report an annoying UI issue. It's really mind boggling how a congregation of idiots with their heads up their asses runs that place.
I would also like to point out that those "savings on infrastructure costs" seemingly do not benefit their users: They have repeatedly increased prices over the past few years.
I'm also unsure whether they are using their competitive advantage properly. Anecdotally, it used to be that almost everyone I knew was only streaming on Netflix. These days, the streaming services that people use seem much more diverse.
An example (in this case ONE PIECE on Netflix Japan region): My Netflix language is set to English, but I am not in an English speaking country. And yet, all shows have an English description and title, as well as episode lists with descriptions also in English. And yet, the program itself is not in English, nor does it have English audio options or subtitles. And the UI does not indicate if there are English subtitles until I play an episode and open the subtitles menu. This is compounded to be even worse when there are so many shows that are only half subtitled (different rights holders I assume). Why does the UI lie to me by showing everything about the series/movie in my native language except the actual content itself? This is really common in non-English speaking regions, and it looks like a basic engineering failure of looking up their global database content in MY language while ignoring whether the actual CONTENT is available in my language. I suppose this could be a UX issue, but it also looks like an engineering one to me. And aren't those intertwined to some extent anyway?
Am I really alone to notice that each year, our collective UX across consumer (SaaS) software gets worse?
The engineering was very thoughtful and the internal tooling was really good.
One underappreciated aspect of simply engineered systems is that they are simple to think about, simple to scale up, and simple to debug when something goes wrong.
EVCache is a disaster. The code base has no concept of a threading model. The code is almost completely untested* too. I was on call at least 2 time when EVcache blew up on us. I tried root causing it and the code is a rats nest. Avoid!
From the looks of it, each module has plenty of tests - and the codebase is written in a spring/boot style, making it fairly intuitive to navigate.
But Momento exists now. It solves every problem EVCache was supposed to solve.
There are other options too. They should retire it by now.
Seems to be cloud hosted only.
EVCache definitely has some sharp edges and can be hard to use, which is one of the reasons we are putting it behind these gRPC abstractions like this Counter one or e.g. KeyValue [1] which offer CompletableFuture APIs with clean async and blocking modalities. We are also starting to add proper async APIs to EVCache itself e.g. getAsync [2] which the abstractions are using under-the-hood.
At the same time, EVCache is the cheapest (by about 10x in our experiments) caching solution with global replication [3] and cache warming [4] we are aware of. Every time we've tried alternatives like Redis or managed services they either fail to scale (e.g. cannot leverage flash storage effectively [5]) or cost waaay too much at our scale.
I absolutely agree though EVCache is probably the wrong choice for most folks - most folks aren't doing 100 million operations / second with 4-region full-active replication and applications that expect p50 client-side latency <500us. Similar I think to how most folks should probably start with PostgreSQL and not Cassandra.
[1] https://netflixtechblog.com/introducing-netflixs-key-value-d...
[2] https://github.com/Netflix/EVCache/blob/11b47ecb4e15234ca99c...
[3] https://www.infoq.com/articles/netflix-global-cache/
[4] https://netflixtechblog.medium.com/cache-warming-leveraging-...
[5] https://netflixtechblog.com/evolution-of-application-data-ca...
So you put 200us network with 30us response time and get about 250us average latency. Of course the P99 tail is closer to a millisecond and you have to do things like hedges to fight things like the hard coded eternity 200ms TCP packet retry timer ... But that's a whole other can of worms to talk about.
[1] https://github.com/Netflix-Skunkworks/service-capacity-model...
For plugging into other apps they may only need a small slice of EVCache; just the fetch from local-then-far, copy sets to multiple zones, etc. A greenfield client with the same backing store could be trivial to do.
That all said I wouldn't advise people copy their method of expanding cache clusters: it's possible to add or remove one instance at a time without rebuilding and re-warming the whole thing.
I can see someone setting up a huge number of counters then leaving...and in a hundred years their counters are taking up TB of space and thousands of requests-per-second.
If you have high cardinality metrics, it can still be really painful, although I think you will feel the pain initially and it won't take years. Usually these systems have a way to inspect what metrics or counters are using the most resources and then they can be reduced or purged from time to time.
- Write counter changes to a Kafka topic with many partitions. The partition key is derived from the counter name.
- Use Kafka connect to push all counter events to S3 for audit and analysis.
- Write a Kafka consumer that reads events in batches and updates a persistent store with the current count.
- Pick a good Kafka message lifetime to ensure that topic size is kept under control, but data is not lost.
This gives us:
- Fast reads (count is precomputed, but potentially stale)
- Fast writes (Kafka)
- Correctness (every counter is assigned exactly one consumer)
- Durability (all state is in Kafka or the persistent store)
- Scalable storage and compute requirements over time
If I were to really go crazy with this, I would shard each counter further and use CRDTs to compute the total across all shards.
1. What was the count for counter X between times T1 and T2? 2. "I am going to re-run my batch job again from yesterday. Adjust the increments for this window and re-compute the final count".
Although the #2 use-case requires lot of other nuances around Recounting, which we allude to but don't expand upon in the article (adjustable retention, multiple rollup checkpoints per counter, pushing back accept-limit for backfills etc.)
I suspect if I were to recursively ask "why?", we may eventually wind up at some triviality (to the actual business/customer) that could have easily gone another way and obviated the need for this abstraction in the first place.
Just thinking about the raw information, I don't see how the average streaming media consumer produces more than a few hundred kb of useful data per month. I get that the delivery of the content is hard, but I dont see why we need to torture ourselves over the gathering of basic statistics.
But what in the world are they using a global counter for that needs "low millisecond latencies"? I don't see a single example in the entire article, and I can't think of any.
I don't see how any of those would suffer if the numbers took seconds to update instead of milliseconds.
We're talking about having updated numbers in milliseconds, right? Not just "the database responds in a reasonable amount of time" because that's been solved many many times over and the article specifically says "this category requires near-immediate access to the current count at low latencies".
> some initial use cases related to interactive titles
Maybe 1 second of latency for a group interaction? That's still orders of magnitude more slack. And I'd expect only moderate accuracy requirements.
I would be more interested in how a higher traffic video company like Pornhub handles things like this.
Or even simpler, Roaring Bitmaps: https://pncnmnp.github.io/blogs/roaring-bitmaps.html
The queue solution is pretty elegant.
How many times has Netflix been entirely down over the years?
Seems it's not "nonesense".