> 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=3....