12 ms·
Cache made consistent: Meta’s cache invalidation solution
- bklaasen 4y ago> Phil Karlton famously said, “There are only two hard things in computer science: cache invalidation and naming things.” ...and off-by-one errors.
- hinkley 4y agoAnd assuming that certain problems are different from each other and should be counted separately. (I'm starting to think that 'naming things' and 'cache invalidation' are the same thing, since you need to invalidate the cache anytime the description of the thing changes)
- uvdn7 4y agoI am the author of the blog post. I believe the methodology described should be applicable to most if not all invalidation-based caches. I am serious when I say that cache invalidation might no longer be a hard thing in computer science. AMA!
- politician 4y agoI read the article. It seems to be suggesting that cache invalidation might no longer be a hard thing in computer science because of the insights uncovered by your tracing and observability solution. IOW, now that you can measure and audit the system, you're able to find and fix bugs in the system. Is that the correct take-away?
- uvdn7 4y agoYep. You're exactly right. And the approach is generic and I think it should work for everyone. The idea is that with all these observability capabilities, debugging cache inconsistencies is getting very actionable and close to how we debug an error with message telling us exactly where an exception happened.
- deleted 4y ago[deleted]
- robmccoll 4y agoSaying a problem isn't hard to solve because you have tools to analyze the correctness of your solution seems like a stretch.
- uvdn7 4y agoBut that's not what I am saying though ... > you have tools to analyze the correctness of your solution That's half of it. Cache invalidation is hard not only because of the complexity of cache coherence protocols, some of which is not very complicated. But cache invalidation does introduce countless races that manage to introduce cache inconsistencies in ways that are just hard to imagine ahead of time (in my experience). IMO, that's the harder part of cache invalidation – when cache inconsistencies happen, answering the "why" question is much harder than having a cache invalidation protocol (you can have TLA+ for one if you will). And answering the "why my cache is inconsistent" is the problem we solved, which I think is the harder part of the cache invalidation problem.
- latchkey 4y ago> answering the "why my cache is inconsistent" is the problem we solved That should be the title and focus of your post. Instead, it feels like grandiose claims about solving cache invalidation itself.
- uvdn7 4y ago> Instead, it feels like grandiose claims about solving cache invalidation itself. That is definitely something I worried about. But at the same time, I do think we solved the harder part of the cache invalidation problem. TAO and Memcache serves quadrillions queries a day. Based on our experience, answering the question of "why my cache is inconsistent" is the hardest thing about cache invalidation. Cache invalidation protocols can be complicated, but some are pretty managable. You can also verify it using TLA+ if you will. But once the rubber hits the road, some cache entries tend to be inconsistent. I definitely worry about people taking this the wrong way, but at the same time, I stand by the claim of "cache invalidation might no longer be a hard thing in computer science".
- jitl 4y agoHow do you handle caching derived data assembled from multiple individual records? For example, how would you maintain a cache for a query like getFriendsOfFriends(userId)? My context is working on Notion’s caches for page data. Our pages are made out of small units called “blocks” that store pointers to their child content. To serve all the blocks needed to render a page, we need to do a recursive traversal of the block tree. We have an inconsistent cache of “page chunks” right now, how would you think about making a consistent cache?
- uvdn7 4y agoThis is fantastic question! We face something very similar (if not identical). There are two parts to solve this problem 1. monitor and measure how consistent the cache is 2. figure out why they are inconsistent I will focus on #1 in this comment. You can build something very similar to Polaris (mentioned in the blog) that - tails your database's binlog so it knows when e.g. "friendship" data is mutated - it can then perform the computation to figure out which cache entries "should have been" updated. E.g. if Alice just friended Bob, then Alice's friends-of-friends and Bob's friends-of-friends should reflect the change. And your monitoring service will "observe" that and alert on anomalies
- ahahahahah 4y agoSo if you just implement the cache invalidation logic in your monitoring system, you can tell if you got the cache invalidation logic correct in your caching system. That sounds really helpful!
- uvdn7 4y agoThat's not the case though. Polaris acts as a client and only monitors client observable effect, and assumes no knowledge of the server internals. I am trying to help.
- ahahahahah 4y agoThis continues to point out how you just completely don't understand the quote. You very clearly seem to think that the "hard" part of cache invalidation is how to implement invalidating it when you know exactly what needs invalidation. The "hard" part is actually in knowing what needs invalidation. Your grandiose claims make you, your team and org, and your company look bad.
- continuational 4y agoWith the risk of stating the obvious - the hard part of cache invalidation is to know when to invalidate the cache.
- uvdn7 4y agoWe invalidate cache upon mutations. When else would you do it?
- continuational 4y agoPlease see https://news.ycombinator.com/item?id=31672541 https://news.ycombinator.com/item?id=31672541
- rajesh-s 4y agoOff topic but what tool did you use to create those cache hierarchy diagrams?
- simonw 4y agoIs your argument here that cache invalidation may no longer be hard because you can implement systems like Polaris which continually test your cache implementation to try and catch invalidation problems so you can then go and fix them? EDIT: Already answered in this reply: https://news.ycombinator.com/item?id=31671794 https://news.ycombinator.com/item?id=31671794
- uvdn7 4y agoPolaris is actually fairly simple to build. The harder question to answer is "why cache is inconsistent" and how you debug. The second half of the post talks about consistency tracing, which tracks all cache data state mutations. Distributed systems are state machines, with consistency tracing keeping track of all the state transitions, debugging cache inconsistencies in an invalidation-based cache is very actionable and managable based on our experience.
- ArrayBoundCheck 4y agoI didn't understand the tracing part. Is tracing 100% inside of polaris? If not does it start at the database? The first cache? Does the invalidation need to be predictable before you can use tracing? What kind of parameters do you have so you don't have too much logging or too little?
- uvdn7 4y agoTracing is not in polaris. It's a separate library that runs in every cache host. Its main job is logging every cache state mutation for traced writes. So when polaris detects cache inconsistencies, we know what happened and why cache is inconsistent. It starts all the way from client initiated write (where we mint a unique id, and plumb it all the way through). > What kind of parameters do you have so you don't have too much logging or too little? This is the key question! If you go to the second half of the post, it talks about an insight about how we managed to log only when and where cache inconsistencies _can_ be introduced. There's only a small window after mutation/write where cache can become inconsistent due to invalidation. So it's very cheap to trace and provide the information we need.
- nixpulvis 4y agoCorrect me if I'm wrong, but don't we have general solutions that ensure correctness if you allow for enough time? Wouldn't performance be a critical aspect of the claim that this is "no longer a hard problem"? How about the CAP theorem?
- uvdn7 4y agoThe "hard problem" defined here is more about the engineering side. The analogy I have is Paxos the protocol and Google's Paxos Made Live paper. We do have many cache coherency protocols that are provably correct (using TLA+ if you will); but making them actually consistent in production is a completely different story; and a hard problem for an invalidation-based cache.
- gnomeduck 4y agoCan you speak more on why ‘making them actually consistent in production is a completely different story’? Curious because I’ve been learning TLA+ recently and interested to more know about cases where an algorithm has been proven but actual an implementation of it fails.
- uvdn7 4y agoLet try to put this as concisely as possible. When you put an algorithm in code, when rubber hits the road, you have to make certain assumptions of how things work, eg how events are ordered, how fsync works, what kind of failure scenarios you are expecting, etc. More likely than not, the reality will be a little different. No to measure just innocent bugs in the code. My favorite example is Paxos. Its algorithm fits on a single slide. But it’s notoriously hard to make it actually work correctly in production.
- ntoskrnl 4y agoCorrect me if I'm wrong, but the title seems a little clickbaity. "Cache invalidation might no longer be hard" --> "We built a tracing library and manually fixed the bugs it uncovered"
- uvdn7 4y agoI am the author; so I am obviously biased here. I am serious when I say cache invalidation might no longer be a hard thing in computer science. In the post, I explained why cache invalidation is hard; and how we solve/manage its unique challenge and complexity. By my definition, we are solving the cache invalidation problem. The analogy I have is Paxos. It's notoriously hard to implement Paxos correctly. Google published a paper on Paxos Made Live just on how they managed the implementation and productionization complexity.
- continuational 4y agoPlease be careful with such bold claims. You don't really address the hard part of cache invalidation, which is to figure out when to do it.
- uvdn7 4y ago> which is to figure out when to do it. Can you elaborate?
- continuational 4y agoSure - you typically cache the result of some expensive query. The hard part of cache invalidation is to detect when an update somewhere in your system is going to affect the result of that query, such that you need to trigger an invalidation.
- uvdn7 4y agoGood point! We have memcache which is a look-aside cache that serves this type of workload. What you described can be solved by adding one level of abstraction and letting reads and writes all go through it. Now on your read path, you can construct arbitrary complex sql queries or whatnot, but it must take some kind of input to filter on. Those become part of the "keys". The invariant is that as long as the "context/filter" you encode covers all the mutations which would impact your cache data, you should be good. Based on our experience, it has been fairly managable.
- VWWHFSfQ 4y agoAs naive as it is, I always kinda liked MySQL's query cache invalidation strategy: blow away the whole cache anytime a mutating statement arrives at the server. Simple. Effective. But it obviously doesn't work for everything.
- Trufa 4y agoCause who needs scaling anyway!? Joking aside, strategy seems logical for smaller stuff.
- pyrolistical 4y agoI don’t get it. If you already have the version in the cache, then when a client does GET X, the cache can reply X@v4. As long as the client submits transactions with “I calculated this based off of X@v4” then the database can reject if there is a newer version of X. The client can then replay the transactions against the cache and by then there will be a X@v5. With this scheme you can track the transaction replay rate instead of having to build a new system to read-back the cache for consistency. To get the traceability on which cache node is stale, the client can forward a header from the cache node that identifies it. Then using the same transaction rejection rate ops has visibility on which cache node isn’t being updated. No cache invalidation needed. Always just cache forever X@version, it’s just that the cache allows an unstable query for GET X.
- jitl 4y agoTracking all data dependencies used for all writes in a large system seems rather challenging. Plus what about read workloads, like “pull message from queue and send it as a push notification to user Y”? I guess it’s fine if the push is re-delivered due to stale cache?
- dragontamer 4y agoVersions of cache-data aren't numbers, they're vectors-of-numbers (at a minimum). Here's a bunch of versions of variable "X": X@Alice-V1, X@Alice-V2, X@Bob-V1, X@Alice-V2/Bob-V2, X@Carol-V1 Which version of "X" should you read? EDIT: Alice, Bob, and Carol are three different servers. What has happened here is that Alice and Bob are closer together, so their updates are faster between each other (synchronizing the writes). Carol is slower for some reason (bad connection?), and is going to update off of old data. In this case, the two valid caches are X@Alice-V2/Bob-V2, and X@Carol-V1. The other cached data is old and can be discarded. Things get complicated when you have not only reads (such as in your simplified question), but also writes that are cached (such as what happens in practice).
- londons_explore 4y ago> In other words, 99.99999999 percent of cache writes are consistent within five minutes. Those are the kind of rates that lead to very hard to find bugs. You won't properly write code in downstream systems to properly handle that failure case, and suddenly someone will get paged at 3am because it has just happened years after the code was deployed and nobody really understands how this corner case suddenly happened! When I'm building systems, every happy codepath should either happen frequently or never. Having a cache whose design allows it to be stale but only very rarely seems like a bad idea. If I was building something on top of this cache I would have a layer on top to force staleness for 1 in 1000 requests just to force the application built on top to be designed to handle that properly.
- uvdn7 4y agoI agree with everything you said. > Having a cache whose design allows it to be stale but only very rarely seems like a bad idea. And that's exactly why we first brought this number from 6 9's to more than 10 9's; and we are not done yet. Also the cache is not "designed" to be inconsistent. It's like Paxos is designed to solve the consensus problem, but when implemented, most Paxos implementations do not work properly.
- dang 4y agoThe submitted title ("Cache invalidation might no longer be a hard thing in Computer Science") broke the site guidelines badly. They say: "Please use the original title, unless it is misleading or linkbait; don't editorialize." - https://news.ycombinator.com/newsguidelines.html https://news.ycombinator.com/newsguidelines.html Editorializing the original title to make it more linkbait is going the wrong way down a one-way street. Please don't. (If you're the author, then please either use the original title or, if it's linkbait, change the HN title to make it less so.)
- uvdn7 4y agoAccepted. Just to be clear, it’s not a link bait as far as I understand it for two reasons 1. I am dead serious about making cache invalidation a simpler problem 2. “Cache invalidation might no longer be a hard problem in Computer Science” is the subtitle of the corresponding systems @scale talk I just gave. I respect the guidelines and I am OK with the rename. I just want to clarify that I don’t think the old title is misleading.
- PKop 4y agoIt's still massively click-bait until such time as the claim isn't pure speculation. Also, casually optimistic claims about complex and hard problems are the essence of click-bait. There isn't much unique about your instance of it
- uvdn7 4y ago> until such time as the claim isn't pure speculation Which I don't think it is. > casually optimistic claims about complex and hard problems are the essence of click-bait A group of people at Facebook worked on this for many years. TAO and Memcache serves quadrillions of queries a day. I wanted to make this claim – cache invalidation might no longer be a hard problem in Computer Science – years ago; but I didn't because I don't want to jump the gun. We baked the solution for years. There's nothing casual about this.
- 4y ago
- Trufa 4y agoThis seems to be a pure click bait title, cache invalidation will always be hard since it's basically a balanced give and take issue which simple logic dictates can't avoid.
- uvdn7 4y agoCan you elaborate? I went into details in the post about what I believe is the root cause of why cache invalidation is hard, which seems different than what you are saying.
- deleted 4y ago[deleted]
- judofyr 4y ago> At a high level, Polaris interacts with a stateful service as a client and assumes no knowledge of the service internals. This allows it to be generic. We have dozens of Polaris integrations at Meta. “Cache should eventually be consistent with the database” is a typical client-observable invariant that Polaris monitors, especially in the presence of asynchronous cache invalidation. In this case, Polaris pretends to be a cache server and receives cache invalidation events. I'm a bit confused here. Polaris behaves as a regular cache server which receives invalidation events, but isn't the typical bug related to cache invalidation that a service forgets to contact the cache server for invalidation? So this will only catch cases where (1) you remember to contact Polaris, but you forgot to contact other cache servers [which Polaris happens to know about], OR (2) you're not handling errors during invalidation requests to the cache server [and the request to Polaris was successful]? Or are you cache servers "smart" and might have internal logic which "ignores" an invalidation? What am I missing? EDIT: Reading through "A real bug we found and fixed this year" and I'm still a bit confused. It seems like a very contrived bug directed directly to how you deal with versioning (e.g. you allow the latest version to be present with stale metadata, or something?). My main concern with cache invalidation is what to invalidate at what time.
- uvdn7 4y ago> isn't the typical bug related to cache invalidation that a service forgets to contact the cache server for invalidation? That's not the case based on our experience at FB. At-least-once delivery is a solved problem basically. But you are absolutely right that if there's an issue in the invalidation delivery, it's possible that polaris won't receive the event as well. Polaris actually supports a separate event stream (all the way from client initiated writes) to cover this case.
- benlivengood 4y agoIt helps that TAO is a write through cache so clients can't really forget to invalidate, correct? If someone were to directly write to MySQL shards there would be stale data ~indefinitely. I'm assuming the versioning is what ensures proper write, invalidation, and fetch ordering so that e.g. slow mysql writes/replication don't cause remote TAO clusters to receive an invalidation message and then read a stale value from a replica?
- rossmohax 4y agoHow both of these can be true? > Data in cache is not durable, which means that sometimes version information that is important for conflict resolution can get evicted. > We also added a special flag to the query that Polaris sends to the cache server. So, in the reply, Polaris would know whether the target cache server has seen and processed the cache invalidation event. To make special flag work, cache server need to track not only current version state, but also past versions. If it tracks past versions, then conflicts can be resolved at the cache server level, but whole premise of article is that cache servers can't resolve conflicts by themselves.
- uvdn7 4y agoGreat question! I will try to answer this one without too much implementation details. First of all, cache items can be evicted and are not durable, which is just a fact. But that doesn't mean we can't track progress (in terms of "time" or "logical time"/"order" in distributed system terms). The social graph is sharded (not surprisingly). We can keep progress of each cache host per shard, which is just a map kept in memory, which doesn't get evicted. I hope this answers your question.
- isodev 4y ago
- gtirloni 4y agoYour comment does not contribute anything useful to this technical subject.
- yodsanklai 4y agoTotally off-topic, but... https://www.facebook.com/help/152637448140583 https://www.facebook.com/help/152637448140583 > No, we don't sell your information. Instead, based on the information we have, advertisers and other partners pay us to show you personalized ads on the Facebook family of apps and technologies. A lot can be said about Meta but I find it counterproductive to make up facts. It's not correct that they "sell personal details".
- anonymoushn 4y agoSure, and American prisons don't sell slaves to American farms. They lease slaves to American farms.
- deleted 4y ago[deleted]
- deleted 4y ago[deleted]
- AtNightWeCode 4y agoSo basically, classic timestamp versioning with some consistency checking. Might work. Cache invalidation is hard. The only problem I ever faced while working in the biz that I could not solve even at theory level was cache related. What I thought the article would tackle is the dependencies between stored objects. Some solve it with super complicated graphs others by just flushing the entire cache. At Meta, why not just use flat files or a store that act as one, pubsub the changes, listen on the changes and update aggregated views as needed and store them. Then just short-term cache everything on the edge.
- uvdn7 4y ago> What I thought the article would tackle is the dependencies between stored objects. Like tracking dependencies, so it knows what to invalidate on mutations?
- AtNightWeCode 4y agoYes. The Boss changes his name from Mark. It will directly hit a lot of caches. It will need an update of all bosses that materialized this. Then all the staff that materialized the boss of the boss and so on. This is simple example because it is a tree. This can be circular dependencies as well.
- gwbas1c 4y ago> Based on consistency tracing, we know the following happened... I'm a little confused: - Is the database updated notification going out while the transaction is still in-progress? Doesn't it make more sense to delay the notification until after the transaction commits? - If a cache tries to update itself while there's an open transaction against the data it's trying to read, shouldn't the database block the operation until the write is complete? And then shouldn't another notification after the transaction commits trigger the cache to update?
- uvdn7 4y ago> Is the database updated notification going out while the transaction is still in-progress? No, that's not what happened. > If a cache tries to update itself while there's an open transaction against the data it's trying to read, shouldn't the database block the operation until the write is complete? Yes. Generally speaking, having a read transaction here would be better. There are unrelated reasons why it's done this way. The point of the example is that it's really intricate and we can still identify the bug regardless.
- gwbas1c 4y agoAhh, dirty reads when you need data integrity are dangerous. (They're fine for things like a UI, where the consumer of state is essentially transient.) I'm not too familiar with your database, but I wonder if you could have a "read transaction" with such a short timeout that you could return a "try again soon" error to the cache? (The only time I dealt with dirty reads like this was Microsoft SQL, I suspect you're using something with very different transaction semantics or guarantees.)
- benlivengood 4y agoDo you guarantee referential integrity in TAO? Last I heard the answer is no; complex queries are not supported and clients would have to do their own work with internal versions to achieve a consistent view on the whole graph (if such a consistent view even exists). But it seems to work fine since global referential integrity doesn't seem to be a big deal; there aren't a lot of cases where dependency chains are long enough or actions quick enough that it matters (e.g. a group admin adds a new user, makes them an admin, who adds an admin, and then everyone tries to remove the admin who added them [always locally appearing as if there is >=1 remaining admin], but causes the group to ultimately have no admins when transactions settle). Run into any fun issues like that? Contrasting that with Spanner where cache coherency is solved by reading from the recent past at consistent snapshots or by issuing global transactions that must complete without conflict within a brief truetime window. I am guessing the cost of Spanner underneath the social graph would be a bit too much for whatever benefits might be gained, but curious if anyone looked into using something similar.
- uvdn7 4y agoFor the fun issues you described, read-modify-write (either using optimistic concurrency control or pessimistic) can work, if I understand your question correctly. Spanner is awesome and one of my favorite systems/papers. I think it would be very computational and power expensive to run the social graph workload on a spanner-like system. Do you have an estimate of how many megawatts (if not more) are needed to support one quadrillion queries a day on Spanner?
- benlivengood 4y agoI wish I had more solid data. Cloud Spanner claims about 10K read QPS and 2K write on a "node" which costs $1/hour. The Spanner paper reports about 10K read QPS per core. As I understand the cloud deployment it's 3 replicas at that price, so $0.33/hour should buy me about 8 cores and so I'm not sure what the disparity is (maybe markup? CloudSQL is 2X the cost of compute), but I'll go with 10K QPS/8 cores at the low end. Anyway, something like 500-1000W for ~100 cores so something between 100 QPS/W and 10 QPS/W using the 1000W and high and low performance estimates, or 10^15 / 86400 = 11.6e12 QPS for somewhere between 116MW and a a little over a GW. Sounds comparable(?) to MySQL+TAO on the low end to ridiculously expensive at the high end. EDIT: I have no clue how efficient MySQL+TAO really are but figure that at least tens of thousands of machines go into it.
- moralestapia 4y ago??? We seem to have different concepts with regards to the cache invalidation problem. The one I know has to do with "Is there a change in the source data? How should I check this? How often?" and such things. Yours seems to be "In my distributed system, I cannot guarantee the order of the messages that get exchanged, and thus, different nodes may end up having different information at times". IMO you have a consensus/distributed data problem, not a cache invalidation problem (the inconsistency between your cache nodes is a consequence of the underlying design of your system).
- uvdn7 4y agoI think https://news.ycombinator.com/item?id=31672541 https://news.ycombinator.com/item?id=31672541 might address your question.
- moralestapia 4y agoNo, mine was not a question, it was a statement. You misnamed the problem you are trying to solve.
- uvdn7 4y agoOK. I am here to learn. Say we go with your definition of cache invalidation problem. Do you think that's THE hard part of cache invalidation? Say, you are just caching simple k/v data (no joins, nothing). Are you claiming cache invalidation is simple in that case? Also the reason why I didn't mention the cache invalidation dependencies is that _I believe_ it's a solved problem (see the link above, we do operate memcache at scale). I am happy to discuss if and why it wouldn't work for your case.
- moralestapia 4y ago>Are you claiming cache invalidation is simple in that case? No. >I believe it's a solved problem It's not. Suppose you have source-of-truth A (doesn't really matter if it's a key-value store or whatever, it could even be a function for all intents) and a few clients B1, B2, B3, ... that rely on the data from A. You have to keep them in sync. When should B* check if A has changed? Every time they need it? Every minute? Every hour? Every day? This is the cache invalidation problem; which IMO is not even a problem but a tradeoff, but whatever. Epilogue: With all due respect, you have an unbelievable career, IBM, Autodesk, Google, Box and now Meta. None of those companies would give me five minutes of their time because I am self-taught, yet here we are :)
- irrational 4y agoPhil Karlton famously said, “There are only two hard things in computer science: cache invalidation and naming things and off-by-one errors.”
- deleted 4y ago[deleted]
- eckesicle 4y agoand exactly once delivery and exactly once delivery.
- spread_love 4y agoThis joke keeps getting extended so much it will be "there are only 2 hard problems: [list of 15 things]" in a couple years
- irrational 4y agoI’m not sure. I’ve seen the “and off-by-one error” version of the joke for 15+ years. I’ve never seen anything added to it.
- Izkata 4y agoconcurrency Next up:
- LeonB 4y agoNot quite.
- lucideer 4y agoThis would have been a really good article in itself, without the frankly unbelievable hubris of the author. TL;DR this is a blogpost about some very interesting problems with distributed cache consistency and observability. The article doesn't address cache invalidation and - based on their replies here - the author doesn't seem to understand the cache invalidation problem at all.
- uvdn7 4y agoI am glad that you liked the content. On the definition of cache invalidation, and specifically why it's hard. Can you send me a link to any definition of it? This is what's in wikipedia and I think it's reasonable. > Cache invalidation is a process in a computer system whereby entries in a cache are replaced or removed. And I think I am describing that process, and what's hard about it. Some comments here explicitly talk about dependencies, which I can see why it's hard. My point is that even without dependencies, cache invalidation remains a hard problem. Now about dependency tracking, some of my thoughts are captured here https://news.ycombinator.com/item?id=31674933 https://news.ycombinator.com/item?id=31674933.
- lucideer 4y agoNumerous replies to your comments here have pointed out why cache invalidation is hard. You've responded to them by saying "Good point!" (always with an exclamation point), and then proceeded to demonstrate in your response that you didn't understand their point. This comment describes the "hard" part of cache invalidation best: https://news.ycombinator.com/item?id=31674251 https://news.ycombinator.com/item?id=31674251 - your response to them makes very little sense. First, you acknowledge that they're correct, though you use obtuse language to say so: you seem to like using the very abstract term "dependency" to represent the very simple concept of detecting updates. In your second paragraph you then go off on an unrelated tangent by saying: > Now getting back to the “what’s really hard about cache invalidation” part. No. You're not "getting back" to that - you're changing the subject back to the topic of your article, which is unrelated to cache invalidation. > Say you just have a k/v store, no joins, nothing. Cache invalidation is about invalidation - it's not about your store architecture. It's not about the implementation of logic that processes the clearing/overwrite of values, it's about the when and nothing else. > just doing TTL, would be simpler. Yes. It is simpler. If you want to avoid solving the hard problem, you can use a TTL. Now... how long should it be?
- orf 4y ago> Take TAO, for example. It serves more than one quadrillion queries a day That’s 11,574,074,074 requests per second.
- throwdbaaway 4y agoBigger than the human population queries per second. 1 request from human would likely create multiple queries internally.
- pdevr 4y agoFrom the article: >>Polaris pretends to be a cache server De facto, it is a cache server, with the associated performance decrease. Isn't the main purpose of caching to improve performance? >>We deploy Polaris as a separate service so that it will scale independently from the production service and its workload. Assuming the scaling is horizontal, so then, to synchronize among the service instances, what do you do? Create another meta-Polaris service? Not rhetoric or sarcasm - hoping for an open discussion.
- deleted 4y ago[deleted]
- uvdn7 4y ago> De facto, it is a cache server, with the associated performance decrease. Isn't the main purpose of caching to improve performance? Polaris only receives cache invalidation event, and doesn't serve any client queries. > so then, to synchronize among the service instances It doesn't synchronize among the service instances. It's actually basically stateless. It pretends to be a cache server only for receiving invalidation events, and acts as a client, and assumes no knowledge of the cache internals.
- pdevr 4y agoThanks for the clarification. However, in that case, won't the querying of all the cache copies/replicas turn out to be the bottleneck at high volumes? Because you are going to have to check all the existing copies/replicas somewhere, right?
- uvdn7 4y agoYep. Polaris checks can be sampled, but still covers pretty much all types of workload and scenarios/races over time.
- Hnrobert42 4y agoYou can’t even visit the FB blog while using NordVPN. Amazing.
- rgbrenner 4y agoI see a couple of monitoring/reporting systems.. but no caching solution. These are tools to catch bugs in the solution you're using. Good work, but not a solution for cache invalidation. And regarding those tools: it doesnt sound like Polaris would handle a network partition well. If a cache invalidation triggers it to check the other caches for consistency... that assumes Polaris will receive that invalidation message. Imagine a scenario of 5 cache servers, and a Polaris server. On one side of the split is Polaris and 2 servers, and the other has 3 cache servers... It's possible for the 3 cache servers to receive an update that is not received by the polaris+2 network... And polaris would not just be unaware of the inconsistency, but it also wouldn't know to check later for the inconsistency when the network partition is resolved. I also feel like the consistency tracing is assuming that only one fill and invalidate is occurring at a time (when in practice, there may be multiple fills and invalidates occurring in parallel on the same data)... and that those calls will arrive in order. If they arrive out of order, it doesnt sound like it would catch that.. and I think you're relying on Polaris to catch this case, but high latency cannot be differentiated from a network partition except the length of the delay... so these two types of errors can be seen together.. in which case, you'd have a cache error that neither tool would detect. I would like to hear why Im wrong. I understand this is being used in prod, but network partitions don't occur everyday... and Im not convinced this has seen enough flaky networks to work out the bugs.
- uvdn7 4y ago> but not a solution for cache invalidation Yeah I have learned about some people have different perceptions of what cache invalidation means (I have a narrower definition, a subset of what some people think as cache invalidation). I will not die on this hill. I am happy to rename/rephrase/anything to be helpful. > it doesnt sound like Polaris would handle a network partition well ... Good question. And you are right. My thought on this is that over time, polaris would still catch those unique issues happened at network partition or whatnot. Polaris doesn't promise have 100% coverage all the time. However, over time, if there are flaws in the system, Polaris should help surface it. And once it did (even with just one inconsistent sample), it should be very actionable, and we can find out why it happened, and fix the underlying cause. > I also feel like the consistency tracing is assuming that only one fill and invalidate is occurring at a time (when in practice, there may be multiple fills and invalidates occurring in parallel on the same data)... and that those calls will arrive in order. If they arrive out of order, it doesnt sound like it would catch that.. It doesn't make that assumption actually. Distributed systems are state machines. Consistency tracing essentially logs (it doesn't do much detection) state transitions. So when polaris detects an anomaly, we have all the information to help us diagnose. And you are right that fills and invalidations can happen in parallel and it's fine. E.g. if fill happens before invalidation, we can always track state mutations caused by the invalidation. If invalidation is the culprit (the earlier fill can't be because it happened earlier), we would have a log for it.
- uvdn7 4y agoI understand some people have a different definition of cache invalidation. I am using the following definition from Wikipedia > Cache invalidation is a process in a computer system whereby entries in a cache are replaced or removed. It's the smallest unit of function that "cache invalidation" must perform. Some people define cache invalidation as a problem of figuring out "when/who" to invalidate. That's not the definition I am using (and my definition is narrower in that sense). The confusion is not intended, as I believe (as I explained in the post) that even with the narrower definition, cache invalidation is still insanely hard. If it helps, let's essentially break cache invalidation into a few parts 1. knowing when/who to invalidate 2. actually processing the invalidate My argument is that #1 can be very managable with simpler data models (as we did). #2 can't be avoided; and #2 is very hard. And the post is about how we believe we have a systemic approach for managing #2. For #1, say, you are caching a result of joining two tables with two ids that you are filtering on. It's still very managable to track the dependency and know when to invalidate. It can easily grow out of hand (join 10 tables with 100 lines of SQL). Then solving the "when/who" to invalidate problem is essentially equivalent to doing "joins" on the write/invalidation path. First of all, it's unbounded. The number of cache entries you need to invalidate can be unbounded (not bounded by the number of indices, but a function of data in the database instead). My argument is that why do this to begin with? I acknowledge this is hard. But why do it? On the other hand, you can have simpler data models (e.g. TAO), fetching and stitching everything together on the read path scales fairly well. It's essentially doing "joins" on the read path. But it's all hitting caches, so it's fast still. For some complicated queries, you can cache secondary indices (which is easier to figure out the "when/who" question, just as how DB figures out which index entry to update on transactions) to make your read-path join faster. The write amplification / cache invalidation fanout is bounded. You don't do "join" on writes/invalidations. Let's discuss the "when/who" problem specifically, if say we just have to solve it. E.g. in its most generic form, a cache can store arbitrary materialization from any data source. Now when updating the data source, in order to keep caches consistent, you essentially need to transact (cross system transaction) on both the data source and cache(s). Usually cache has more number of replicas, I am not sure running this type of transactions is practical at scale. What happens if we don't transact on both systems (the data source, and cache)? Well, now whenever the asynchronous update pipeline performs the computation, it's done against a moving data source (not a snapshot of when the write was committed). Now let's say the data source is Spanner, which provides point-in-time snapshots. On Spanner commit you can get a commit time (TrueTime) back. Now using that commit time, to read the data and compute cache update asynchronously can be done. On the other hand, it's very easy to feel like #2 is an easy problem, which probably explains why people think #1 is what Phil Karlton was referring to. The analogy I like to use is Paxos. The protocol fits on a single slide. It's easier to feel like you have Paxos work; but it's very hard to have Paxos actually work. Here's the recording of my talk at systems @scale which I hope people find it helpful (no matter how you define cache invalidation – the broader vs. narrower definition) – https://www.facebook.com/watch/live/?ref=watch_permalink&v=384770970336225 https://www.facebook.com/watch/live/?ref=watch_permalink&v=3....
- shrimpx 4y agoMy problem has never been invalidation as stated in the article, but how to manage interdependent cache entries, so that when one is invalidated, the dependent ones must also be invalidated. For example get_user(123) and get_users_for_org(456). Suppose user with id 123 is part of the org with id 456. When the user is deleted, you have to invalidate the get_users_for_org(456) entry. I haven’t seen any convincing “design pattern” for managing such dependencies.
- MaxMoney 4y agoIn the past I've used a timestamp key for the org. Create a key like f"org_{org.pk}" that has a timestamp. Then append that timestamp to the `get_users_for_org`. Now all you have to do is invalidate the org timestamp to generate a rolling key (requires more memory).
- uvdn7 4y agoYep. I respect this definition, and I tried to clarify things a bit here https://news.ycombinator.com/item?id=31676102 https://news.ycombinator.com/item?id=31676102. I am sorry for the confusion. We have internal abstractions that deals with the problem you described (especially in simpler forms, the example you have is fairly managable actually). EDIT: I can see why it gets complicated if your data model and query is very complicated.
- somedudetbh 4y agoBased on the wide disagreement over whether or not this solves cache invalidation, this post provides strong evidence that the hardest bit is in fact naming things.
- deleted 4y ago[deleted]
- deleted 4y ago[deleted]
- c3534l 4y agoI've heard this sort of thing is almost as hard as naming things.
- hoseja 4y agoYeah the nines are very trendy and impressive but wouldn't it be more comprehensible to say "one in a million" and "one in ten billion"?
- FranksTV 4y agoThe nines are just a way to say 1e-5 without having to use notation.
- DeathArrow 4y ago>Distributed systems are essentially state machines. If every state transition is performed correctly, we will have a distributed system that works as expected. I imagine huge distributed systems as Markov chains, where transitions are probabilistic rather than deterministic.
- uvdn7 4y agoInteresting… what are your thoughts on Lamport’s TLA+?
- DeathArrow 4y agoSo you measure with Polaris and have found you have a data consistency of up to 10 nines in a period of 5 minutes. But aren't inconsistencies building over timer? What if you measure for 1 hour instead of 5 minutes? Still 10 nines?
- uvdn7 4y agoThis is a very good question. The number of nines actually go up with increasing timescale windows because Polaris checks anomalies at write-time/invalidation-time. I guess the behavior of what you are describing is cache inconsistencies at read-time. E.g. if I introduced an inconsistent cache entry at time T, it can be exposed to many subsequent reads (hence the "building up over time" as you mentioned). This is an important metric as well – we actually measure it as well. The key difference between read-time cache consistency measurement and write-time cache consistency measurement is about "purpose". Write-time cache consistency measurement is more actionable, as it captures the moment (or very close) of when cache becomes inconsistent. If one wants to debug something, you want to get close to when the anomaly happens. Read-time cache consistency measurement is more about measuring the negative impact of cache inconsistencies (which are client facing).
- mbj111 4y agoTweet thread from veteran engineer - balance between the desired system properties and the costs of coordination (https://twitter.com/MarcJBrooker/status/1534944325997453312 https://twitter.com/MarcJBrooker/status/1534944325997453312)
- uvdn7 4y agoMarc put it extremely well. I agree with every single word of his thread. I should have applied a narrower and more specific definition of cache invalidation in the blog post. I apologize for any confusions it caused.
- jerDev 4y agoCould the diagram on this page be any worse?