Designing a Distributed Cache with a Globally Aware Eviction Policy using Caffeine, InfintiSpan & RabbitMQ

Viewed 279

I’m using ScyllaDB / Cassandra as a global, distributed data store & Caffeine, Infinispan & Hazelcast as a local, in memory cache.

I have an application running on 1,000 nodes. When a node requests data from the global, distributed data store the data is cached locally using Caffeine / Infinispan. This way the application doesn’t have to request the same data over and over again - thereby reducing the load on the distributed data store.

When an individual node updates a given piece of data on the global data store it is relatively easy for that node to evict the corresponding data from its cache as it can simply signal to the local cache to evict/invalidate the data.

The problem is that any one of the 1,000 nodes can hold the same data and any node can update any piece of data at any time. If Node 539, for example, updates a specific piece of data and Node 877 holds a copy of that data in its local cache, I need Node 877 to evict the data from its local cache and retrieve the data in real time the next time it is needed. Of course, the same data can be cached on dozens of nodes and all of them would need to be made aware of the update made by Node 539 and evict the data accordingly.

What is the best way to design a system like this?

While I don’t want to reinvent the wheel I couldn’t find any existing solutions that is capable of achieving this so I devised my own plan:

My plan is to use a distributed messaging system such as RabbitMQ (and perhaps Kafka) where each of the 1,000 nodes subscribes to a topic which contains a list of data IDs that need to be evicted. Whenever a node updates a particular piece of data, it writes the “data ID” to the "eviction" topic. Every one of the 1,000 node subscribes to the "eviction" topic and evicts the data linked to the data ID in real time, if it holds that data in memory.

However, I have several concerns with this design.

First, it seems extremely inefficient. Every time any one of the 1,000 nodes updates a piece of data, the “data ID” will have to be propagated to all of the 1,000 nodes, since we don’t know which (if any) node(s) holds a specific piece of data. Additionally, it is likely that none of the nodes will cache most of the data being updated thereby making it even less efficient. It might even be more efficient just to read all the data in real time, every time.

Is there a more elegant/efficient design to achieve the desired goal?

I’m using Java.

Thanks

2 Answers

Basically you moved the problem from the DB space to the app space. Even if the DB is consistent, you need to make the local cache consistent (with which type of guarantee? there are many forms of consistency). It's not a simple problem to solve, it will need to survive failures, etc.

ScyllaDB caches really well, how about you'll increase its capability and move some ram from the local nodes to Scylla?

One alternative option is to always write to the DB first and have the nodes listen to the CDC stream and invalidate or repopulate their local cache

Btw, there is a nice local memcache project (in C++) - https://www.scylladb.com/2019/02/20/valustor-a-memcached-alternative-built-on-scylla/

What is the best way to design a system like this?

I don't think there is a best way to design a system like this without looking into the details of what you're trying to accomplish. So I'll try to provide several answers that cover a few different cases. I think the most crucial detail that is missing is the consequences of serving stale data from the local cache. As you note in your response in the comments to @SpaceTrucker above:

I'm trying to devise a system that works with data where the ramifications of holding onto outdated data is not good.

But how bad is it if you serve stale data? Can you tolerate a bit of staleness?

First, if you cannot tolerate any staleness, you can:

  1. Try to roll your own consistency solution at the local cache level. This will probably require some sort of consensus protocol or other really high complexity stuff that you probably do not want to write yourself. This is what @dor laor is noting above, saying that you are moving "the problem from the DB space to the app space."
  2. Use something like Redis as a shared remote cache and ditch the local cache. Obviously this will be a much higher latency solution compared to a local cache, but could be a good compromise solution, particularly if you have a need to reduce load on your database rather than merely improving latencies (though it could still provide better latencies than going directly to Cassandra, but that comparison is something that should probably be measured).
  3. Just go straight to your database. If you have problems with load, then just add more nodes.

I don't include your proposed solution here, because it leaves open the possibility of serving stale data. If, for instance, your RabbitMQ gets overloaded, you're going to stop processing evictions, and that means lots of stale data. If you can't tolerate serving stale data, then I think using an evictions notification queue isn't really a solution.

On the other hand, if you can tolerate some staleness, then the most elegant solution, it seems to me, would be to use a local cache with an appropriately sized ttl, and allow some stale data to get served. This could be combined with some other simple tools to achieve a decent functionality profile, for instance, if you can use sticky sessions to ensure that consumers always hit the same node and thus get the latest results from "their" cache instead of getting served from any of a number of caches that might give a different view of the state.

On top of all of that, I think your solution is good, but it seems like some eviction queue (or CDC stream as @dor laor suggests) is only really necessary in the rather niche case where you can tolerate some staleness in extreme conditions, but want it practically eliminated in normal circumstances.

Related