Distributed SQLite for Go applications
github.com
github.com
What this project (and similar projects such as Rqlite) doesn't address is the hard problem -- sharding. Raft makes it quite trivial to write a master/slave system where every node is guaranteed to be identical, but it won't scale writes and won't distribute storage across your cluster.
You can build sharding on top of Raft, but in an SQL setting this would require distributed transactions, something that's difficult enough that relatively few projects have tackled it so far, the notable ones being CockroachDB and TiDB. (Google Spanner, being proprietary, is not relevant in this context.)
> It's fascinating how Raft has democratized distributed computing
First I thought you probably ment commoditized but then I realized how democratized actually applies perfectly well to consensus algorithms because by definition they apply decisions via voting by all participants :)But anyways, I agree very much with you that Raft has transformed the industry and we can't be thankful enough for it. But like you said, it's just one of the core building blocks of a distributed system that also shards data. Sharding is more use-case specific and so it can be valid for different databases to have different sharding logics.
FWIW, the CockroachDB people developed a modified version they call MultiRaft [1] that deals with the scaling challenges they have around shards.
If you don't have those requirements, dqlite or rqlite might get the job done with less moving parts and operational overhead.
It's pretty much the same argument of using SQLite vs (say) PostgreSQL, but translated into distributed systems. See:
Well, if you run it on a very fast and reliable local network without much load, not multi dc deployments over public internet, than sure, it can sort of work. Not without problems though once a fault occurs.
I'm not sure what this means. The Raft paper explicitly states that Raft is a fault-tolerant system -- and by definition this means it deals with faults. To quote the paper:
"Replicated state machines are used to solve a variety of fault tolerance problems in distributed systems."
Raft is a type of replicated state machine. To say that "Raft doesn't deal with faults" is not correct. Perhaps you mean that there is other work to be done to take the fault tolerance offered by Raft and build a functioning application, and a system that will stay up in the real world -- and deal with an even wider range of faults. That I definitely agree with.
As far as Google went with their paper; they were incredibly correct by saying that actually implementing a distributed system is hard because of the failure modes; rather than the success modes of the running system.
I never said that it did. That is what I mean by my statement that more work is required to build a real system in the real world. Much of the code of rqlite is about cluster management, built on a Raft substrate.
>how how does it re-sync the current, up-to-date state to the new node
I suggest you read the Raft paper, and study the Hashicorp Go implementation (https://github.com/hashicorp/raft) to see how this is done. It's all there.
>As far as Google went with their paper; they were incredibly correct by saying that actually implementing a distributed system is hard because of the failure modes; rather than the success modes of the running system.
I couldn't agree more.
- faulty node replacement: you can take any node off (or any node can crash) at any time. As long as there are enough nodes left to reach a quorum, your system will be available. If there are not enough nodes left, your system will be unavailable (but keep consistency)
- re-syncs, snapshots and the rest are all covered in the raft paper
- clients wanting to perform writes always talk to the node that is currently the leader, that node fails, clients will look for next leader (which will be eventually elected as long as there's a quorum of surviving nodes)
- overwhelming a system is a different concern, raft writes are serialized so the goal is usually not throughput (for that you might look at AP/AC storage solutions in the CAP spectrum, raft is CP). Designs that need high write throughput might still use raft as internal building block for coordination.
it is a core feature of the protocol.
> how does it make sure that when one of the nodes fails - that all of the clients do not overwhelm the rest of the system.
for write access (proposals), in a typical implementation (defined below [1]), one failed nodes actually slightly speed up the system with the cost of reduced reliability. Consider a typical setup with 3 nodes, the leader normally replicate state to 2 followers, the replication cost get cut in half once your cluster loses one member.
given that you can implement linearizable read by going through the write procedure, one can argue the above "speed up" can obviously be achieved on reads as well.
[1] independent replication to followers, for N followers, the same entry will be serialised and sent N times by the leader.
again - just talking about typical implementation and the described "speed up" comes with degraded reliability.
https://groups.google.com/forum/#!topic/raft-dev/BaS5Z2NkDxA
https://github.com/rqlite/rqlite/issues/266
dqlite is very interesting. If its patches to SQLite are accepted upstream, dqlite could become the storage engine for rqlite. But right now rqlite is deliberately built on vanilla-SQLite, so it can offer the same correctness guarantees as SQLite.
Yes, this is technically correct. However to meet its goals rqlite (and dqlites) does not require sharding. The point of rqlite is to provide fault-tolerance and high-availability. Again, you are technically correct that rqlite doesn't support sharding, but sharding simply isn't required for rqlite to do what it wants to do. Sharding is a solution to a different type of problem.
But ActorDB does seem to promote inter-actor transactions, and a pattern I've seen encouraged is where you use shared actors (especially with the key-value store functionality) to keep some kind of central lookup table that your app then can use to find actors. For example, if you're modeling HN with ActorDB, you'd have a global list of stories, but each story/comment tree could be a separate actor, and each user would be a separate actor. Posting a story only needs to access the story actor, but to post a comment transactionally, you have to do a transaction that inserts the comment and updates the user's comment history together.
You're right. What I meant was that a 'shard' in actordb isn't meant to be 'all data on one machine/node' - which is how many other databases view sharding, but is expected to be much more fine-grained. Thanks for the other explanations.
Another interesting implementation aspect is its use of LMDB for storing the SQLite pages.
Note for those who are interested in cool-stuff-with-SQLite, there's a similar project called RQLite (mentioned in the README):
https://github.com/rqlite/rqlite
The main differences are laid out in the README, but IMO it really all boils down to the fact that since it's single-writer quorum'd there's only one node actually doing writes, and they've just chosen to replicate right-before (right after?) the WAL commit (I think this is what they're calling frames, basically one chunk of the WAL content), not at the point of receiving a query. This is more like having a read secondary more than anything (I also assume writes are redirected to master).
It says in the readme that it should be expected to be slower than RQlite, I'd love to see some numbers.
I do really like this though -- can't wait to pour through the code and see if there's anything I can learn.
I think I do have something to offer -- mostly trying gossip and making use of aggressive automatic sharding -- but I'm also more interested in what I wanted to build on top of distributed SQLite as well. No need to get down in the weeds if someone's already done the hard work for me :)
By all means contribute to this project. Everything starts somewhere!
I've been looking at https://github.com/pingcap/raft-rs
You're correct: only the master accepts writes. During a failover clients can get an error and are expected to retry against the new master.
Loud and clear on the failure modes, the README was pretty clear about it
Things are still in flux, but yes, I plan to publish benchmarks before making a 1.0 release, as well as improving documentations and introduce some more abstractions to make it easier to use.
The reason it might be slower is some cases is that RQlite replicates statements and dqlite replicates WAL frames (where "frames" roughly means "disk pages"), which are typically bigger.
I didn't run benchmarks against RQlite yet, but I'd expect performance differences to be negligible for most use cases.
I see what you meant by it being slower, but that's only replication speed right? maybe that should be pointed out in the documentation when you find time to update it.
Does it seem like the replication patch is going to land in Sqlite? I sure would like it to...
when you perform a SQLite write, SQLite writes a new page on disk to the its write-ahead log. dqlite needs to replicate that page write across all nodes. A page is typically 4kb, so when you do a write on the leader node you need to transfer 4kb across the wire to all other follower nodes and wait for a quorum of them to come back and say "I got it" (which includes writing the page to disk). On the contrary rqlite just needs to transfer and store the SQL text, which is typically less than 4kb.
In practice it's still pretty performant, but I'll publish benchmarks later on.
I plan to submit the patch to SQLite upstream too, yes.
In principle rqlite could use dqlite.
Does that qualify as "high load"? That's just below 4 pages/second. Given how beefy server CPUs and SSDs are, I would consider something like 100+ requests/second "high load".
Sadly you cannot extrapolate at all the load of a server with their numbers of pages seen per day, depending on the use you can easily have peaks of more than 10x that.
But what kind of overhead postgress adds comparing to sqlite?
Advantages are clear:
- more advanced functionality is available when you need it.
- pgsql likely has better performance comparing to sqlite
does anyone know a resource to learn more about testing in this kind of situation? I hope there is an alternative to just try all cases in different real hardware setups (which is the mantra at my company)
Hardware does not matter that much, since SQLite is pretty hardened for handling any possible hardware failure (out of memory, out of disk space, memory/disk corruption) and propagate that to client code (including dqlite). Regarding network failures (hardware-related or not), raft guards you against them, and it's formally proved.
1) Distributed (e.g. n equal processes running on n machines) 2) Needs some shared state across nodes 3) Would like that state to have SQL semantics (relations, transactions, etc.) 4) Not to heavy on writes to this shared state (I might provide some ballpark numbers at some point) 5) Wants to avoid the operational overhead of an external storage system (mysql, postgresql) 6) Wants to be fault-tolerant and have transparent failovers
then dqlite might be a good choice.
Dqlite and rqlite are not in the same space as postgres and mysql, they're more comparable to tools like Zookeeper, etcd, or Consul[0] - low throughput CP data stores that are generally used for coordinating distributed systems. In this case, it saves users from having to manage an external data store to run lxd clustered.
[0] in fact, dqlite uses the same raft implementation as Consul, which is cool
I think the best explanation of why you'd use it over PG or MySQL is the same as why you'd use SQLite over them:
https://www.sqlite.org/whentouse.html
With dqlite you also get HA, fault-tolerance and transparent failover, if you're willing to trade it with a bit of performance.
edit: I'm sorry, I was thinking of BedrockDB.