We have solved these issues in a few ways, mainly:
- working with the relevant product teams to implement appropriate rate limiting or improving data modeling.
- introducing our own query layer, written in Rust that sits in front of Cassandra that uses a form of micro-caching called read coalescing, and also other forms of query throttling/load shedding to reduce work the database must do for hot keys/pathological patterns of access. We expose a GRPC interface from this - and this lets us centralize control of the client driver and tune it appropriately, while also getting to leverage the ever growing open source grpc traffic routing solutions (envoy, etc...)
and ultimately,
- switching to ScyllaDB, a C++ rewrite of Cassandra which is of course void of any garbage collection issues, and features faster overall performance and lower latencies.
Scylla, however, is not without its own set of issues - and somewhat strict hardware requirements[0] thanks to the seastar engine it is built on top of. Their team however has been delightful to work with, and our platform is markedly more stable in current year than it was in years past thanks to the above factors.
Operationally, however, Scylla and Cassandra are quite easy to run, the trickiest part is repairs. Common operations such as cluster expansion, or replacement of node are so common an operation that they are at this point mundane. Be wary however about read/write amplification issues inherent to LSMT databases, choosing the correct compaction strategy and tuning it appropriately can be quite key. Additionally tombstones can be quite bad for performance.
In current day we offer a new more generic solution that sits on top of scylla (it would work with Cassandra too) that provides a simple interface to query KKV based data, without having to worry too much about problems like large partitions, hot keys, or tombstones! With a design like this, the underlying cluster thus far has been issue free and very easy to operate.
[0]: https://discord.com/blog/how-discord-supercharges-network-di...
If you had fully explored using those and found them insufficient then it would have been an interesting case study in them "not working". Supposedly this is a nearly "solved problem" but I am very curious to hear real world use cases and whether it truly works or not.
- Use a better a GC like Azul's C4.
- Leverage monitoring to shard off hot spots.
- If possible, get away from CPU- & memory-inefficient JVM and dynamic languages.
- Steer clear from ScyllaDB and its radioactive license. And, there's little customer demand for it. Astra DB for a commercial distro or YugabyteDB is from the folks behind Cassandra.
- Optimize persistence with multiple choices of backends that abstract away the underlying datastore: the eases the pain of going through "divorce" from any particular technical solution. Writes vs. Reads ratios, geographic replication, object size, data structure (K/V, tuplestore, time series, files, etc.)