Turning SQLite into a Distributed Database
univalence.me
univalence.me
Their vision was to build the hardest parts of building a database, such as transactions, fault-tolerance, high-availability, elastic scaling, etc. This would free users to build higher-level (Layers) APIs [1] / libraries [2] on top.
The beauty of these layers is that you can basically remove doubt about the correctness of data once it leaves the layer. FoundationDB is one of the most (if not the) most tested [3] databases out there. I used it for over 4 years in high write / read production environments and never once did we second guess our decision.
I could see this project renamed to simply "fdb-sqlite-layer"
[1] https://github.com/FoundationDB/fdb-document-layer
That is very interesting and simple and valuable insight that seems to be missing from the wiki page. But also from the wiki page <https://en.wikipedia.org/wiki/FoundationDB>, this:
--
The design of FoundationDB results in several limitations:
Long transactions- FoundationDB does not support transactions running over five seconds.
Large transactions - Transaction size cannot exceed 10 MB of total written keys and values.
Large keys and values - Keys cannot exceed 10 kB in size. Values cannot exceed 100 kB in size.
--
Those (unless worked around) would be absolute blockers to several systems I've worked on.
Therefore, it's possible that the id is invalid (for the external storage) when referenced in the future. I think doing so only adds complexity as system grows.
It would be better to chunk your blob data to fit the DB, imho. It beats introducing external blob storage in the long run.
Depends! If the ID is a cryptographic hash, then as long as the blob is uploaded first, then the DB can't be inconsistent with the blob[1].
A Merkle Tree also allows "updates" by chunking the data into some convenient size, say 64 MB, and then building a new tree for each update and sticking that into the database.
[1] With the usual caveats that nobody is manually mucking about with the blob store, that it hasn't "lost" any blobs due to corruption, etc, etc...
Yeah, with those caveats. But how do you make sure they apply? If someone does manually muck about with the blob store, or it does lose blobs due to corruption, then your transaction is retroactively "un-atomicized" with no trace thereof in the actual DB.
But otherwise I agree.
It's basically trouble free unless you run below 10% free space on any instance, where things go bad.
Basically, fdb is great as long as you avoid this situation. If you do, woe unto you trying to being the cluster back online by adding nodes and hoping for rebalancing to fix things. It will, it is just very, very slow. I don't know if that's true in the current version.
[0] https://forums.foundationdb.org/t/brand-new-macos-installati...
How it scales as well.
https://blog.expensify.com/2018/01/08/scaling-sqlite-to-4m-q...
> basic specs: > 1TB of DDR4 RAM > 3TB of NVME SSD storage > 192 physical 2.7GHz cores (384 with hyperthreading)
Not to mention the risk of hardware failure.
That strikes me as the more interesting takeaway sentence.
It's not that much RAM if you think of it as 24 8-core machines.
It's still feels like a lot.
[1]: https://www.sec.gov/Archives/edgar/data/1476840/000162828021...
Why else would they benchmark to show they can achieve 4,000,000 QPS if all they need is 1 QPS?
Another approach using FUSE, making arbitrary SQLite-using applications leader-replica style distributed for HA: https://github.com/superfly/litefs (see also https://litestream.io/ for WAL-streaming backups, that's the foundation of this)
Like Aurora, some tweaks to the engine were required, but the core query engine is largely intact.
Has anyone seen postgresql on FoundationDB? Is there anything unique about SQLite that makes it better suited for this approach? One thing that comes to mind is how they were able to do block level locking instead of full db locking. That took some tweaking but probably significantly less than postgresql might require.
I wrote single machine MVCC in Java and I'm curious if there are other ways of implementing it. One way is event sourcing.
I use an integer to store the latest commit version and I only allow transactions to see versions less than the transaction's timestamp. This is the multiversion part.
The concurrency control part is enforced by checking if the read timestamp of the key is less than the reading transaction timestamp, if so someone got there before us and we abort and restart.
I am thinking how to build the distributed part.
I need a timestamp server the same way Google's Spanner needs TrueTime for monotonic timestamps but also some way of broadcasting read timestamps to detect conflicts between nodes. So I'm thinking of broadcasting timestamp events and using that to detect transactions that have dangerous dependencies.
The caveats for dqlite and rqlite always felt kind of awkward/risky to me -- in stark contrast to SQLite which is so stable/"built in" that you don't think about it's failure modes. Having to worry about what exactly I ran (ex. RANDOM()) was just a non-starter (IIRC rqlite has this problem but not dqlite? or the other way around -- one replicates at statement level the other at WAL level).
That said though, the biggest sticking point with all this SQLite goodness is how to make sure that certain libraries (any popular extension -- vsv, spatialite, libcephsqlite) were loaded for any application using SQLite -- there seem to be only a few options:
- calling load_extension[2] from code (this is somewhat frowned upon, but maybe it's fine)
- LD_PRELOAD (mvsqlite does this)
- Building your own SQLite and swapping out shared libs (mvqslite also does this, because statically compiled sqlite is a nuisance)
- Trapping/catching calls to dlopen (also basically requires LD_PRELOAD, but I guess you could go custom kernel or whatever)
This is probably the one big wart of SQLite -- it's a bit difficult to pull in new interesting extensions.
I also found this hack[3] which looks quite interesting for building something more general/reusable...
[EDIT] - Also while I'm here, I think FDB is probably one of the most under-rated massive-scale NoSQL databases right now. It gets nearly no press (to be fair because it went closed then open again), but it's casually a massive force behind Apple's services at scale.
[0]: https://docs.ceph.com/en/latest/rados/api/libcephsqlite/
[1]: https://github.com/rook/rook/issues/10689
[2]: https://www.sqlite.org/lang_corefunc.html#load_extension
It’s supper cool as it does change sqlite.
I love the use case of querying SQLite from a CDN with range requests, because it allows for real “serverless” querying. For example, this: https://github.com/psanford/sqlite3vfshttp
Author here. Actually I have a similar idea with mvSQLite. Provide a client-side-queryable API, but read-write instead of read-only. Security can be implemented with a "provenance"-style mechanism: the client proves they reached a page following a valid/allowed path, by presenting the path (along with necessary signatures) to the server. That way we can have "serverless" read-write transactions with table-level security.
Have you thought about encryption? Is it possible to do client-side symmetric encryption?
There are a couple of statements on that page I disagree with.
"mvsqlite is a distributed database, while dqlite and rqlite are replicated databases"
I disagree with this statement, and consider its definition of "distributed system" to be incorrect. rqlite[1] is a distributed database. A "distributed system" is simply a system that splits a problem over multiple machines, solving it in a way that is better, more efficient, possible etc than a single machine. rqlite uses distribution to provide fault-tolerance and high-availability. It uses distributed systems technology i.e. Raft, to make rqlite appear up-and-running, even in the face of node failures. That it replicates a full copy of the SQLite database to every node is correct, but that doesn't mean it's not a distributed system. Is Consul a distributed key-value store? etcd? By the definition quoted on that page they are not, but no one would actually agree with that.
rqlite is not just about replicating a SQLite database. I understand what the page is trying to say, but the point is that rqlite is distributed, just for fault tolerance and high-availability.[2]
"(+): mvsqlite runs on a production-grade distributed key-value store, FoundationDB, instead of implementing its own consensus subsystem."
rqlite doesn't implement its own consensus system either. It uses the same Raft consensus code that powers Hashicorp Consul. I don't see how this is any different, in principle, than using Foundation DB's consensus system.
[1] https://github.com/rqlite/rqlite
[2] https://github.com/rqlite/rqlite/blob/master/DOC/FAQ.md#rqli...
rqlite's readme seems to indicate that it uses Consul only for service discovery, and runs Raft internally?
I think the most important difference here is rqlite runs its "data plane" on a single consensus group/state machine while mvsqlite relies on FDB's distributed transaction system (well at the bottom there's a Paxos but it is only used for metadata coordination). Both approaches have advantages and disadvantages though.
I can make that clearer.
So projects like mvsqlite are solving a real need. I just disagree with some of its definitions.
IIRC rqlite replicates commands, dqlite replicates WAL frames, mvsqlite intercepts file system calls because it's a VFS implementation [0].
Many people can say, "wait, that doesn't matter to me at all, that happens at most once a year and all of my customers are also offline when it happens". Indeed, that's why people are doing just fine with postgres. But, some things are a little more mission critical than average, and so these technologies exist for them.
One thing that does tip the scale towards these distributed databases are operational concerns. "Apply security patch" or "migrate to new major version" look a lot to distributed systems like "tornado took out the data center" or "hard disk turned into many disconnected chunks of steel". So while you might not be super concerned about disasters, you might still be interested in keeping your compute nodes disposable or in doing regular maintenance without downtime. I don't think the operational balance is quite there yet, but it's definitely worth looking into every few months to see what the state of the art is.
The only way to a "free lunch" is via CRDTs.
edit: This is in no way a knock on FDB by the way. I think it's a really interesting piece of technology and I want to explore it more. Just that it doesn't have a way to do consensus faster than the speed of light :)
For example: you’re using Postgres. You send the request to FDB, it will ensure all Postgres transactions are consistent or tell Postgres to abort transactions.
Please excuse my ignorance, but what do all these Sqlite-but-make-it-X solutions offer over a simple, more established solution like (nobody ever got fired for choosing) Postgres?
I'm pretty into Litestream/LiteFS. Here's what I'm after:
1. Operational simplicity. Fly.io devs run small app servers in a bunch of places. They're usually read heavy. Running network database servers gets very complicated, very fast. You need a DB node in each place your app server lives.
2. Graceful failure. When you scatter app servers around the world, internet weather causes problems for an app. I want my DB to be ok when that happens. Reads should continue to work, if they can. And writes should fail in a way that makes it obvious what's happening.
3. Good for caching. Most fullstack apps use Postgres and then layer in Redis/Memcached for caching. This is yet another moving part. sqlite has amazing performance for cache workloads.
These all make an embedded DB interesting. If Postgres fails, your app server needs to know that Postgres failed and then also fail. Same if you add Redis. If your embedded DB fails, it's obvious to the app server that something is awry.
Another thing I'm after is drop in usage without touching code. We run a lot of peoples' apps and have limited ability to have them add libraries or write new code. Almost all of their frameworks speak sqlite, though.
Adding clustering to sqlite is not perfect. There are still networks to deal with and things will still break. All it's doing is shifting complexity and giving us different levers to use to keep things reliable.
Litestream/LiteFS are amazing projects. The FUSE-based approach is interesting (I'm implementing something similar in mvSQLite, thanks for the idea!)
> Graceful failure
mvSQLite is designed to continue to operate under degraded network (there is a fault-injection test specifically for checking this property: https://github.com/losfair/mvsqlite/blob/1dd1a80d2ff7263b07a...). Network errors and service unavailability are handled with idempotent retries and not exposed to the application.
> Good for caching
mvSQLite caches pages read and written, and does differential cache invalidation (only remotely modified pages are invalidated in the local page cache). The local cache is just a regular KV store with invalidation strategies, and can be moved onto the disk. So it essentially becomes a consistent local database snapshot.
For something larger, finding out writes to DB fail is not a big deal - apps have logs after all.
Another thing captured my attention is DB sizing/fault tolerance
1.keeping average sized DB of 100-300 GB locally, without much buffer pool/vfs cache with higher chances will slowdown the reads and joins. Something bigger like couple of TB should be even harder to do.
2. Syncing multiple concurrent writes can be even longer or on par - guessing here.
3. Both Redis and Postgres (and MySQL) can be (should be) kept behind LB. If node fails, clustering + LB will handle it. LBs of course can be present per location and provide local reads and forwarding writes to remote.
Regarding LB, probably simplest would be have different tcp ports for readwrite/readonly types of load. Reads can be split among multiple replicas be needed.
I have not used it for anything big, but from my experience small/average load (80k rps to Redis with 3Gbit of bandwidth) LB was not an issue. 500k rps may change the situation of course, not a silver bullet.
I can't speak for all, but keep in mind that SQLite does not require a server or sends requests over a network to get to the data. This means SQLite is a far simpler and cheaper deployment, which is totally acceptable if you don't have high reliability and horizontal scalability in mind. Sometimes it even outperforms production RDBMS.
Bolting on distributed access to SQLite adds horizontal scalability for (in the very least) cases where eventual consistency is more than enough to meet your requirements.