8 ms·
Netflix's Distributed Counter Abstraction
- mannyv 2y agoI wonder how they're going to go about purging all the counters that end up unused once the employee and/or team leaves? 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.
- singron 2y agoThere is a retention policy, so the raw events aren't kept very long. The rollups probably compress really well in their time series database, which I'm guessing also has a retention policy. 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.
- rshrin 2y agoYes, once the events are aggregated (and optionally moved to a cost-effective storage for audits), we don't need them anymore in the primary storage. You can check the retention section in the article. The rollups themselves can have TTL if the users wish to set that on a namespace. Although doing that, they have to be fine with certain timing issues on when the rollups expire and new events are aggregated. We also have automation to truncate/delete namespaces.
- ilrwbwrkhv 2y agoLooks a bit overengineered due to Netflix's own microservices nonsense. I would be more interested in how a higher traffic video company like Pornhub handles things like this.
- oreoftw 2y agoHow would you design it to support mentioned use cases?
- klaussilveira 2y agoHyperLogLog and PostgreSQL: https://github.com/citusdata/postgresql-hll https://github.com/citusdata/postgresql-hll Or even simpler, Roaring Bitmaps: https://pncnmnp.github.io/blogs/roaring-bitmaps.html https://pncnmnp.github.io/blogs/roaring-bitmaps.html https://blog.quastor.org/p/grab-rate-limiting https://blog.quastor.org/p/grab-rate-limiting https://github.com/RoaringBitmap/CRoaring https://github.com/RoaringBitmap/CRoaring
- philjohn 2y agoHyperLogLog doesn't support exact counters though, does it? That seems to be one of the core requirements of the queue-based solution.
- rshrin 2y agoHyperLogLog (or even Count-Min Sketch) will not support some of the requirements of even the Best-Effort counter (clearing counts for specific keys, decrementing counts by any arbitrary number, having a TTL on counts etc.). For Accurate counters, we are trying to solve for multi-region read/write availability at low single-digit millisecond latency at cheap costs, using the infrastructure we already operate and deploy. There are also other requirements such as tracking the provenance of increments, which play a part.
- philjohn 2y agoHave you worked with distributed counters before? It's a hard problem to solve. Typical tradeoffs are lower cardinality for exact counters. The queue solution is pretty elegant.
- ElevenLathe 2y ago
- vlovich123 2y agoIt's a bit weird to not compare this to HyperLogLog & similar techniques that are designed to solve exactly this problem but much more cheaply (at least as far as I understand).
- zug_zug 2y agoI came here to write the same thing. Getting an estimate accurate for at least 5 digits on all netflix video watches worldwide can all be done with intelligent sampling (like hyperloglog) and likely one macbook air as the backend. And aside from the compute save the complexity and implementation time would be much lower too.
- rshrin 2y agoFwiw, we didn't mention any probabilistic data structures because they don't satisfy some of the basic requirements we had for the Best-Effort counter. HyperLogLog is designed for cardinality estimation, not for incrementing or decrementing specific counts (which in our case could be any arbitrary +ve/-ve number per key). AFAIK, both Count-Min Sketch and HyperLogLog do not support clearing counts for specific keys. I believe Count-Min Sketch cannot support decrement as well. The core EvCache solution for the Best-Effort counter is like 5 lines of code. And EvCache can handle millions of operations/second relatively cheaply.
- vlovich123 2y agoIncluding this in the blog would have been helpful although I don’t think the decrement explanation is unsolvable - just have a second field for decrements that is incremented when you want to decrement & then the final result is a sum of the two.
- rshrin 2y agoTrue. You could do decrements that way. We trimmed this article as the post is already quite long. But considering the multiple threads on this, we might add a few lines. There is also something to be said on operating data stores that support HLL or similar probabilistic data structures. Our goal is to build taller on what we already operate and deploy (like EvCache)
- millipede 2y ago> EVCache 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! * https://github.com/Netflix/EVCache https://github.com/Netflix/EVCache
- Alupis 2y agoCan you elaborate? 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.
- jedberg 2y agoI'm surprised it's still there! It was built over a decade ago when I was still there. At the time there were no other good solutions. But Momento exists now. It solves every problem EVCache was supposed to solve. There are other options too. They should retire it by now.
- tmikaeld 2y agoMomento? This? https://www.gomomento.com/ https://www.gomomento.com/ Seems to be cloud hosted only.
- jolynch 2y ago(I work at Netflix on these Datastores) 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-data-abstraction-layer-1ea8a0a11b30 https://netflixtechblog.com/introducing-netflixs-key-value-d... [2] https://github.com/Netflix/EVCache/blob/11b47ecb4e15234ca99c19eb35c38a934433e8fd/evcache-core/src/main/java/com/netflix/evcache/EVCache.java#L500 https://github.com/Netflix/EVCache/blob/11b47ecb4e15234ca99c... [3] https://www.infoq.com/articles/netflix-global-cache/ https://www.infoq.com/articles/netflix-global-cache/ [4] https://netflixtechblog.medium.com/cache-warming-leveraging-ebs-for-moving-petabytes-of-data-adcf7a4a78c3 https://netflixtechblog.medium.com/cache-warming-leveraging-... [5] https://netflixtechblog.com/evolution-of-application-data-caching-from-ram-to-ssd-a33d6fa7a690 https://netflixtechblog.com/evolution-of-application-data-ca...
- dopamean 2y agoWhy would netflix put their blog on medium?
- leakyabstxns 2y agoGiven the complexity of the system, I'm curious to know how many people maintain this service
- rshrin 2y ago5 people (who also maintain a lot of other services like the linked TimeSeries service). The self-service to create new namespaces is pretty much autonomous (see attached link in the article on "Provisioning"). The stateless layer auto-scales up and down based on attached CPU/Network-based scaling policies. The exports to audit stores can be scheduled at a cadence. The only intervention is when we have to scale the storage layer (although parts of it also automated using the same Provisioning workflow). I guess the other intervention is when we decide to change the configs (like number of queues) and trigger a re-deploy. But thats about it. So far, we have spent a very small percentage of our support budget for this.
- deleted 2y ago[deleted]
- notfried 2y agoNetflix Engineering is probably far ahead of any other competitor streaming service, but I wonder how much the ROI is on that effort and cost. As a user, I don’t see much difference in reliability between Netflix, Disney, HBO, Hulu, Peacock and Paramount+. They all error out every now and then. Maybe 5-10 years ago you needed to be much more sophisticated because of lower bandwidth and less mature tech. Ultimately, the only real difference that makes me go for one service over the other is the content.
- l33t7332273 2y agoI’ve thought this for a while, and it’s sad because I want to reward good tech. I usually think about this while I’m waiting for paramount to load for the third time because picture in picture has gone black again when I tried to full screen .
- TZubiri 2y agoThe value of netflix is probably not only on its technical prowess and app experience, but it seems they are pretty involved in content direction through metrics. ¹ ² Sources: [1] https://youtu.be/xL58d1l-6tA https://youtu.be/xL58d1l-6tA [2] https://youtu.be/uFpK_r-jEXg https://youtu.be/uFpK_r-jEXg
- rhplus 2y agoJury’s out on their content direction. They cancel great shows before the first season has even had a chance to permeate. It’s clear they have little interest in long term artistic investments.
- autoexec 2y agoIt also shows that they have no interest in building a library of quality content. They could invest in shows that have proper endings but instead they pollute their library with unfinished works that are certain to either be avoided by netflix customers who know that the show was canceled before it concluded, or piss off any netflix customer who doesn't. 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.
- est 2y agoif anyone want a non-distributed but still very powerful counter service, I'd recommend Telegraf from Grafana or https://vector.dev/ https://vector.dev/
- fire_lake 2y agoI think the design could have been simpler with Kafka (which they touched on briefly): - 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.
- rshrin 2y agoYes, this is one of the approaches mentioned in the article and is indeed a valid approach. One thing to keep in mind is that we are already operating the TimeSeries service for a lot of other high ROI use cases within Netflix. There already exists a lot of automation to self-provision, configure, deploy and scale TimeSeries. There already exists automation to move data from Cassandra to S3/Iceberg. We somewhat get all that for free. The Counter service is really just the Rollup layer on top of it. The Rollup operational nuances are just to give it that extra edge when it comes to accuracy and reliability.
- fire_lake 2y agoDid you also consider AWS managed services? Like a direct write to Dynamo?
- rshrin 2y agoNot for this use case. Other use cases at Netflix use AWS Managed service when it makes sense from a use-case and cost perspective. In this case, using TimeSeries opens the door to a lot of other potential future use cases: 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.)
- bob1029 2y agoThis seems a bit overcooked to me. 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.
- Dylan16807 2y agoWell okay, that's some neat implementation stuff. 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.
- rshrin 2y agoUse cases fetching counts directly in the path of Netflix users/streaming, e.g. user-personalization, feature-gating > what features are shown when you load the home page, dictated by how many times these have been shown before for a given device. The article hints at this in the beginning. Also, there were some initial use cases related to interactive titles, details of which can't be publicly shared [although that is winding down now]
- Dylan16807 2y ago> user-personalization, feature-gating > what features are shown when you load the home page, dictated by how many times these have been shown before for a given device 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.
- rshrin 2y agoFor the Eventually Consistent counter, the low millisecond requirement is for reads and writes, not for the convergence of counts. For this category, the convergence is in the order of seconds (user-personalization, feature-gating fall in this category). For the "Best-effort" category, there are some use cases that run experiments in a single-region and need access to current counts at low latencies (they basically add increments and read the value back in the same call i.e. AddAndGet), but are willing to sacrifice "some degree" of accuracy for it. See the table in the 2nd section. There are multiple dimensions in terms of Read/Write Latency, Staleness, Global reads/writes etc. mapped to the two kinds of use cases. Maybe you are conflating a few things. Finally, there is the experimental type of Accurate counters that can get the current count with high degree of accuracy (but the latency there depends on a few things as explained in the article). The last type is more like what can be done using this approach, no current use case for it.