IceFireDB: Distributed disk storage database based on Raft and Redis protocol
github.com
github.com
SET: 253232.12 requests per second
GET: 2130875.50 requests per second
The 10:1 throughput ratio for GET vs SET is interesting. Redis being in-memory, the rates there are pretty close to the same for read/write.Is a 10:1 ratio typical for a storage backed distributed kv store?
Edit: Looks like CockroachDb has roughly a 3:1 ratio, similar for YugabyteDB:
https://www.cockroachlabs.com/docs/stable/performance.html
https://forum.yugabyte.com/t/large-cluster-perf-1-25-nodes/5...
Also ~3:1 for etcd:
'Typical' is a matter of what guarantees you want to give.
Waiting to report success until a majority have committed allows you to make guarantees with a straight face.. "it will probably be committed in the near future" is not the same thing.
But since you brought it up, FPaxos can have both lower latency and higher throughput than MultiPaxos/Raft[0] ;)
The reason is that in Raft if a node acknowledges to the leader that it wrote something to the log it must not later accept a different write in the same log position.
This mean if for some reason server rebooted with dirty buffered writes that could not be flushed in time. it’s supposed to forgot everything it know and rejoin the cluster using a brand new node id.
I haven't checked the code though so I might be off.
In a single-node system, the best way to increase your write throughput is to batch requests over small chunks of time. Ultimately, the amount of writes you can perform per unit time is either bounded by the underlying I/O sequential throughput, or the business constraints regarding maximum allowable request latency. In the most trivial case, you are writing a buffer containing the entire day's work to disk in 1 shot while everyone sleeps. Imagine how fast that could be.
A distributed system has all of the same properties, but then you have to put this over a denominator that additionally factors in the number of nodes and the latency between all participants. A single node is always going to give you the most throughput when talking about 1 serial narrative of events wherein any degree of contention is expected.
Things that can make a difference: Databases have subtly different definitions of "durability", so they aren't always doing semantically equivalent operations. Write throughput sometimes scales with the number of clients and it is not possible to saturate the server with a single client due to limitations of the client protocol, so single client benchmarks are misleading. Some databases allow read and write operations to be pipelined; in these implementations it is possible for write performance to sometimes exceed read performance.
For open source databases in particular, read and write throughput is significantly throttled by poor storage engine performance, so the ratio of read/write performance is almost arbitrary. That 3:1 ratio isn't a good heuristic because the absolute values in these cases could be much higher. A more optimal design would offer integer factor throughput improvements for both reading and writing, but it is difficult to estimate what the ratio "should" be on a given server absent a database engine that can really drive the hardware.
One sees a lot of 3:1 in practice due to the replication factor. If you have 3 copies of the data and the client can read from any node, you get 3x the read performance as having to have a quorum write on two out of three nodes.
To the GP, for a rough swag of what is possible out of given hardware, a combination of FIO and ACT (measures IO latency under a fixed load) is a good start.
Source: several years dealing with vault and consul.
What demons did you encounter with vault/consul?
It’s a good idea in principle but it’s not got a good ROI
Most of the problems I've had with Vault have been around it's Terraform provider which they've improved enough that it's not an issue anymore.
I think the only thing about Raft that folks don't realize is how disk hungry it gets, if you want fast write performance you gotta make dang sure all those fsyncs can keep up. Our largest Consul cluster today runs on storage-heavy boxes as it does ~500Mb/s of writes pretty much 24/7.
I use Consul w/ Vault today instead of the internal storage for Vault just cause Consul has really nice monitoring around some stuff that Vault doesn't (path-based stuff for the most part), I think the internal storage is a really good option for 90% of use-cases.
We did this for 30 years fine before someone invented this stack on deployments much larger then the average consul or vault deployment these days.
I had something running 15,000 dynamic rps on Apache about 15 years ago.
People are blinded from simplicity by complexity. Eventually complexity owns you. You can only own simplicity.
At the end of the day this is one way to solve a problem that doesn’t need to be solved that someone has convinced you is a problem.
What does this mean? This doesn't mean anything?
Are you saying to push the DB to the client?
It's less than a few hundred lines of Go that just wraps two other databases (syndtr/goleveldb and ledisdb/ledisdb) with a third library (tidwall/uhaha) that provides a Raft API.
Here's the LSET code:
https://github.com/gitsrc/IceFireDB/blob/main/lists.go#L232
func cmdLSET(m uhaha.Machine, args []string) (interface{}, error) {
if len(args) != 4 {
return nil, rafthub.ErrWrongNumArgs
}
index, err := ledis.StrInt64([]byte(args[2]), nil)
if err != nil {
return nil, err
}
if err := ldb.LSet([]byte(args[1]), int32(index), []byte(args[3])); err != nil {
return nil, err
}
return redcon.SimpleString("OK"), nil
}
So what "IceFireDB" is:1. tidwall/uhaha - Raft server (m uhaha.Machine, rafthub)
2. tidwall/redcon - Read/write redis protocol (redcon.SimpleString)
3. ledisdb/ledisdb - Redis-compatible with disk persistence via leveldb (ldb.LSet)
4. syndtr/goleveldb/leveldb - Provides snapshots, other scattered references throughout code
It also includes this seemingly random file below, which seems to implement some string slice overloads using unsafe.Pointer:
https://github.com/siddontang/go/blob/master/hack/hack.go
// no copy to change slice to string
// use your own risk
func String(b []byte) (s string) {
pbytes := (*reflect.SliceHeader)(unsafe.Pointer(&b))
pstring := (*reflect.StringHeader)(unsafe.Pointer(&s))
pstring.Data = pbytes.Data
pstring.Len = pbytes.Len
return
}
// no copy to change string to slice
// use your own risk
func Slice(s string) (b []byte) {
pbytes := (*reflect.SliceHeader)(unsafe.Pointer(&b))
pstring := (*reflect.StringHeader)(unsafe.Pointer(&s))
pbytes.Data = pstring.Data
pbytes.Len = pstring.Len
pbytes.Cap = pstring.Len
return
}Aerospike and ScyllaDB both support a subset of the Redis API and run as a durable cluster. Both have been tested with Jepsen.
Only if you want to be a drop-in replacement and take advantage of existing compatible libraries. With the recent tea around the official Elasticsearch Python library, it becomes a more interesting question.
Tendis is a high-performance distributed storage system which is fully compatible with the Redis protocol.
Makes me wonder if there is any spec for the Redis commands. I.e., in the same way that SQL defines an interface, but leaves the details up to individual implementations, is there a "Redis" interface that leaves the details up to the implementation?
I'm thinking of something similar to ISO or RFC.
The only real pitfall was what part of the CONFIG stuff I needed to implement to make popular redis client libs talk to me and/or use the newer protocol features.
The rest was pretty straight forward, just read the docs for a command, implement the stuff, run the test suite, fix any bugs, repeat.
As far as I know there is no RFC let alone an ISO standard.
A hosted disk based redis protocol compliant capable of sub TB size datasets would be a dream for me.
There are already several community solutions for Redis persistence - this one provides different guarantees.
The name implies the goal is to make it easy to mix "hot" (from memory) and "cold" (from disk) data. The author suggests this.
I might be off (and probably am) but if I remember correctly Redis persistence is more for disaster recovery - you can create snapshots and recover them or replay a log file. That's very different in terms of performance guarantees from persisting the data itself to disk and reading from it.
I was under the impression that's what tools (like this one) and stuff like Ardb try to solve.
It's just that Redis is mostly an in-memory database and if the process is terminated and restarted (for all sorts of reasons) the data can be restored from disk.
So what IceFireDB might be good for is data which would not fit easily into the memory of one node.
Again, it's really not clear to me.
It's unclear.