7 ms·
Turning SQLite into a Distributed Database
- didip 4y agoI am surprised that SQLite can be plugged into FDB just like that.
- conradev 4y agoI really like the idea of using FUSE as SQLite backend, rather than injecting into the process 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 https://github.com/psanford/sqlite3vfshttp
- losfair 4y ago> I love the use case of querying SQLite from a CDN with range requests 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.
- conradev 4y agoInteresting! I had not thought about access controls when all of the pages are mixed together, but that makes sense. Have you thought about encryption? Is it possible to do client-side symmetric encryption?
- nik736 4y agoWhat makes this special compared to rqlite or dqlite? Edit: https://github.com/losfair/mvsqlite/wiki/Comparison-with-dqlite-and-rqlite https://github.com/losfair/mvsqlite/wiki/Comparison-with-dql...
- hardwaresofton 4y agoThe approach is different -- both rqlite and dqlite sit "on top" of SQLite and are much more like replication and coordination layers. IIRC rqlite replicates commands, dqlite replicates WAL frames, mvsqlite intercepts file system calls because it's a VFS implementation [0]. [0]: https://www.sqlite.org/vfs.html https://www.sqlite.org/vfs.html
- otoolep 4y agorqlite author here. 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 https://github.com/rqlite/rqlite [2] https://github.com/rqlite/rqlite/blob/master/DOC/FAQ.md#rqlite-is-distributed-does-that-mean-it-can-increase-sqlite-performance https://github.com/rqlite/rqlite/blob/master/DOC/FAQ.md#rqli...
- ComodoHacker 4y agoLooks like in their view "fully distributed" == "write scalability", which is, of course, very limited view.
- baq 4y agoSounds good in theory and if it really works as a drop in replacement it’s amazing tech, but numbers, please! I can make a beefy Postgres server handle 5-10-20k tps relatively easily, what can I expect from this?
- losfair 4y agoI'm working on the benchmarks. There shouldn't be surprising results with read scalability, as it is basically guaranteed by FoundationDB. For writes it is indeed more complicated and a benchmark would help here.
- tmikaeld 4y agoPlease post the results on HN when you have them, would love to read it
- jrockway 4y agoI don't think distributed databases are aiming primarily at increasing write throughput over a single node. The question is, how many transactions can you write per second from your California datacenter to your New York postgres leader when someone typos the BGP rule controlling that database's IP address? Probably zero. That's the performance number these databases are trying to increase. 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.
- 4y ago
- samsquire 4y agoThis is interesting..I would like to learn more how they implemented MVCC. 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.
- losfair 4y agoFoundationDB handles the hard part! It provides monotonic "versionstamps", externally consistent transactions, along with other useful features. I recommend FDB's architecture docs: https://apple.github.io/foundationdb/architecture.html https://apple.github.io/foundationdb/architecture.html
- hrgiger 4y agoThere are plenty of leader election/replication libraries, shouldnt be problem if you pick one and give a shot if you are not confident about your self implementation or just benchmark and see how it goes.
- hardwaresofton 4y agoSo I just fell down the rabbit hole of figuring out how to use SQLite with Ceph (turns out a thing called libcephsqlite[0][1] exists) -- awesome to see this new take on distributed SQLite. 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/ https://docs.ceph.com/en/latest/rados/api/libcephsqlite/ [1]: https://github.com/rook/rook/issues/10689 https://github.com/rook/rook/issues/10689 [2]: https://www.sqlite.org/lang_corefunc.html#load_extension https://www.sqlite.org/lang_corefunc.html#load_extension [3]: https://github.com/cventers/sqlite3-preload https://github.com/cventers/sqlite3-preload
- mike_hock 4y agoWait, so the one method that isn't a massive hack is the one that's "frowned upon"?
- hardwaresofton 4y agoAh sorry I didn’t state that clearly — calling the function From SQL (rather than C/FFI) is frowned upon I think
- alberth 4y agoDon’t forget BedrockDB (built on SQLite) that’s used in production at Expensify. How it scales as well. https://bedrockdb.com/ https://bedrockdb.com/ https://blog.expensify.com/2018/01/08/scaling-sqlite-to-4m-qps-on-a-single-server/ https://blog.expensify.com/2018/01/08/scaling-sqlite-to-4m-q...
- edf13 4y agoThis seems like and incredibly expensive way of achieving their query goals… > 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.
- samatman 4y ago> To be clear, the above specs would be pointless for most databases, as almost nothing scales to handle this kind of hardware well — and almost nobody tries. 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.
- Macha 4y agoFor a database? At expensify scale? It's still feels like a lot.
- robertlagrant 4y agoWhat is expensify scale?
- Macha 4y agoTen million users [1] averaging maybe averaging 40 or so interactions per year (filling out an expense report isn't a common task for a lot of those users)? Napkin math of 1.2 qps. Even if you 10x that for backing workflows, and double it because users are more active than expencted, that's still only 30qps. [1]: https://www.sec.gov/Archives/edgar/data/1476840/000162828021020115/expensifys-1.htm https://www.sec.gov/Archives/edgar/data/1476840/000162828021...
- endisneigh 4y agoI’m interested in how FoundationDB can make anything consistent. 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.
- lifeisstillgood 4y agoI am fascinated by how the initial authors of all these packages went from "that would neat if I had 3 months free" to "yeah, doing it now"
- solarkraft 4y agoI've seen a lot of Sqlite hype here in the past weeks and months. 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?
- deleted 4y ago[deleted]
- otoolep 4y agoCan only speak for rqlite: https://github.com/rqlite/rqlite/blob/master/DOC/FAQ.md#why-would-i-use-this-versus-some-other-distributed-database https://github.com/rqlite/rqlite/blob/master/DOC/FAQ.md#why-...
- arinlen 4y ago> Please excuse my ignorance, but what do all these Sqlite-but-make-it-X solutions offer over a simple (...) 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.
- mrkurt 4y agoYou need to consider some of what sqlite is good at. Querying local sqlite is very fast. Much lower latency than querying postgres. It's also very reliable – a local file will always work better than a network service. 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.
- monstrado 4y agoThis is exactly what the engineers behind FoundationDB (FDB) wanted when they open sourced. For those who don't know, FDB provides a transactional (and distributed) ordered key-value store with a somewhat simple but very powerful API. 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 https://github.com/FoundationDB/fdb-document-layer [2] https://github.com/FoundationDB/fdb-record-layer https://github.com/FoundationDB/fdb-record-layer [3] https://www.youtube.com/watch?v=OJb8A6h9jQQ https://www.youtube.com/watch?v=OJb8A6h9jQQ
- zasdffaa 4y ago> 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. 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 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.
- monstrado 4y agoThis project (mvSQLite) appears to have found a way around the 5s transaction limit as well as the size, so that's really promising. That being said, I believe the new RedWood storage engine in FDB 7.0+ is making inroads in eliminating some of these limitations, and this project should also benefit from that new storage engine...(prefix compression is a big one).
- EGreg 4y agoHow does this compare to, say, CockroachDB?
- Joel_Mckay 4y agoI suspect because most of CockroachDB is written in Go, that it performs much better. Of course, this assumes no sequential keys or auto-increment indices (global state interdependence chokes distributed dbs to a crawl).
- medv 4y agoAlso https://github.com/rqlite/rqlite https://github.com/rqlite/rqlite It’s supper cool as it does change sqlite.
- yencabulator 4y agoRelevant reading about FoundationDB building a SQL database on top of a distributed key-value store: https://www.voltactivedata.com/blog/2015/04/foundationdbs-lesson-fast-key-value-store-not-enough/ https://www.voltactivedata.com/blog/2015/04/foundationdbs-le... (That one replaced SQLite's btree, this one puts pages of the btree as values in the key-value store.) Another approach using FUSE, making arbitrary SQLite-using applications leader-replica style distributed for HA: https://github.com/superfly/litefs https://github.com/superfly/litefs (see also https://litestream.io/ https://litestream.io/ for WAL-streaming backups, that's the foundation of this)
- formerly_proven 4y agoSpeaking of "turning SQLite into things it wasn't really meant for", does anyone know of a compressing time-series layer for SQLite?
- stereosteve 4y agoWhile not sqlite… duckdb is inspired by the design of sqlite https://duckdb.org/ https://duckdb.org/
- abujazar 4y agoBut why?
- phamilton 4y agoIf I understand this correctly, it's similar in design to AWS Aurora or GCP AlloyDB.The underlying storage provides the distributed primitives and the DB itself just reads and writes to it. 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.