LiteFS
fly.io
fly.io
What I'd like to have seen is how this compares to things like rqlite[1] or Cloudflare's D1[2] addressed directly in the article
That said, I think this is pretty good for things like read replica's. I know the sales pitch here is as a full database, and I don't disagree with it, and if I was starting from scratch today and could use this, I totally would give it a try and benchmark / test accordingly, however I can't speak to that use case directly.
What I find however and what I can speak to, is that most workloads already have database of some kind setup, typically not SQLite as their main database (MySQL or PostgreSQL seem most common). This is a great way to make very - insanely, really - fast read replica's across regions of your data. You can use an independent raft[3][4] implementation to do this on write. If your database supports it, you can even trigger a replication directly from a write to the database itself (I think Aurora has this ability, and I think - don't quote me! - PostgreSQL can do this natively via an extension to kick off a background job)
To that point, in my experience one thing SQLite is actually really good at is storing JSON blobs. I have successfully used it for replicating JSON representations of read only data in the past to great success, cutting down on read times significantly for APIs as the data is "pre-baked" and the lightweight nature of SQLite allows you to - if you wanted to naively do this - just spawn a new database for each customer and transform their data accordingly ahead of time. Its like AOT compilation for your data.
if you want to avoid some complexity with sharding (you can't always avoid it outright, but this can help cap its complexity) this approach helps enormously in my experience. Do try before you buy!
EDIT: Looks like its running LiteFS[5] not LiteStream[0]. This is my error of understanding.
[1]: https://github.com/rqlite/rqlite
[2]: https://blog.cloudflare.com/introducing-d1/
Litestream and LiteFS are by the same author and serve different purposes.
https://litestream.io/alternatives/ goes into this. It includes LiteFS and rqlite.
[Edited for clarification]
It is missing the D1 comparison though.
> What I'd like to have seen is how this compares to things like rqlite or Cloudflare's D1 addressed directly in the article
I think a post comparing the different options is a great idea. I'll try to summarize a bit here though. LiteFS aims to be an analogue to Postgres replication but with built-in failover. Postgres uses log shipping to copy database state from a primary to its replicas, as does LiteFS. LiteReplica is probably the closest thing to LiteFS although it uses a dual GPL/commercial license model. LiteFS uses Apache 2.
There are Raft-based tools like rqlite or dqlite. These have higher consistency guarantees, however, they tend to be more complex to set up -- especially in an ephemeral environment like Kubernetes. LiteFS has a more relaxed membership model for this reason. As for Cloudflare's D1, they haven't released many details and I assume it's closed source. It also requires using a custom JavaScript client instead of using a native SQLite client.
There are also several eventually consistent stores built on SQLite such as Mycelial[1]. These work great for applications with loose consistency needs. LiteFS still maintains serializable isolation within a transaction, although, it has looser guarantees across nodes than something like rqlite.
> Most workloads are already have database of some kind setup, typically not SQLite as their main database (MySQL or PostgreSQL seem most common)
Yes, that's absolutely true. I don't expect anyone to port their application from Postgres/MySQL to SQLite so they can use LiteFS. Databases and database tools have a long road to becoming mainstream and it follows the traditional adoption curve. I wrote a Go database library called BoltDB about 10 years ago and it had the same adoption concerns early on. Typically, folks try it out with toy applications and play around with it. Once they get more comfortable, then they create new applications on top of it. As more people use it, then late adopters get more comfortable and it further builds trust.
We're committed to LiteFS for the long term so we'll be making updates, fixing bugs, and we'll keep trying to build trust with the community.
Oh, I forgot to touch on this point. LiteFS uses some concepts to Litestream (e.g. log shipping), however, it doesn't use Litestream internally. It has much stricter requirements in terms of ensuring consistency since it's distributed so it performs an incremental checksum of the database on every transaction. It has additional benefits with its internal storage format called LTX. These storage files can be compacted together which will allow point-in-time restores that are nearly instant.
I think a strong - extremely strong - selling point is the point I made about "prebaked" data for your APIs, since the entire strength of these SQLite based systems reside in their fast read capacity (as mentioned elsewhere and in this article, its very fast for read heavy applications, which is most) you could take on an angle around that to get people "in the door" by showing a pathway of how this fits inside your existing data warehouse / storage model.
We found we liked the SQLite durability to do this. It was a bit smarter than just a plain cache (such as Redis) with better durability and (for our needs) comparable enough performance (I think in absolute terms, a tuned Redis instance will always be faster, but up to a certain point, speed isn't everything, especially when factoring cost).
We found it was cheaper - by a good margin - to do this over caching everything AOT in a redis cluster, and we could therefore much more cheaply go multi-region and have DB's sitting next to our customers that acted as a nearline cache.
The complexity - which is an area where this might help in the future, and why I'm mentioning it - is shipping changes back. What we ended up doing is setting up a write Redis cluster that clients write to, and we take those writes and trigger a propogation job back to the database. This allowed us to run a much slimmer redis cluster and made us feel more comfortable doing cache eviction since we could verify writes pretty easily. You could do this with memcache or whatever too.
Sounds convoluted, but it worked really well. It allowed us to keep our centralized database intact without having to spin up expensive instances to be multi-region or commit to ever growing Redis cluster(s). the SQLite flat file model + history of durability made things the perfect tradeoff for this use case. Of course, YMMV, however it was a novel solution that used "enterprise grade" parts all the way down, which made it an easy selling point.
You might find it worth exploring this more.
As far as the comparisons go, I think it'd be cool to see a deep dive, and run a test suite against each of the major SQLite as a distributed database model. For that, I don't think it has to be open source to do a reasonable comparison?
I think caches are an excellent use case for LiteFS early on. Sorry I didn't make that point in my previous reply. It's a good way to get benefits out of LiteFS without committing to it as your source of truth. Also related, Segment built a custom SQLite-based solution[1] for distributing out cached data that worked well for them.
[1]: https://segment.com/blog/separating-our-data-and-control-pla...
> We found it was cheaper - by a good margin - to do this over caching everything AOT in a redis cluster
Do you remember specifics of the cost difference? I can imagine it'd be pretty significant since you don't need to spin up servers with a bunch of RAM.
> Note: I apologize if this is overstepping, its hard to tell!
Not overstepping at all! It's great to hear folks' feedback.
With SQLite, we could just scale instances during peak demand (so you could read from different instances of the same data if there was a bottleneck) and scale back down again, without (!usually) losing the data.
It was a really complex - but fun - project all told, but its underpinnings were really simple.
How Segment approached the problem isn't dissimilar to how we did it, honestly.
The only thing that we (may) have done different is we had failover. If SQLite didn't respond our API layer could then talk directly to a database service to get the data. That was surprisingly complex to do.
Its entirely possible that even more robust setups than ours was when we did this would yield higher cost savings. We did this before we hit our next scale of customers, just to add a little more context
I mean, folks can do stuff like this on Fly with Redis backed by disk, too: https://fly.io/blog/last-mile-redis/
> LiteFS still maintains serializable isolation within a transaction, although, it has looser guarantees across nodes than something like rqlite.
Picking up a term from the consistency map here [0], what guarantees LiteFS makes across nodes?
However, during normal operation it'll function more like Snapshot Isolation. LiteFS does provide a transaction ID so requests could wait for a replica to catch up before issuing a transaction to get something closer to Serializable.
LiteFS aims to be the analogue to Postgres/MySQL replication and those generally work great for most applications.
We do have plans to run against the Tcl test suite[1] although most of that test suite is not applicable to LiteFS since it tests higher level constructs.
Since LiteFS acts on the raw pages, it really just functions similar to a VFS. Out of the 130K source lines of code in SQLite, only 4.8k are for the Unix VFS. As such, the testing coverage from the SQLite test suite mostly tests non-VFS code.
Ie i wonder if there's a way to can write your applications such that they have less/minimal contention, and then allow the databases to merge when back online? Of course, what happens when there inevitably _is_ contention? etc
Not sure that idea would have a benefit over many SQLite DBs with userland schemas mirroring CRDT principles though. But a boy can dream.
Regardless, very cool work being done here.
However, if you have a single writer and just need replicas to get updates when they periodically connect then yes, LiteFS could fit that use case.
This is misleading AFAICT. The article(s) is actually comparing remote RDBMS to local RDBMS, not Postgres to SQLite.
Postgres can also be served over a UNIX socket, removing the individual query overhead due to TCP roundtrip.
SQLite is a great technology, but keep in mind that you can also deploy Postgres right next to your app as well. If your app is something like a company backend that could evolve a lot and benefit from Postgres's advanced features, this may be the right choice.
No practical way to test this, of course.
In a number of cases, we were working with data that would really benefit from actual GIS tooling so PostGIS was kind of the natural thing to reach for. Many of these kiosks had slow, intermittent, or just straight up no internet connection available.
So we just deployed PostGIS directly on the kiosk hardware. Usually little embedded industrial machines with passive cooling working with something like a generation or two out-of-date Celeron, 8GB RAM, and a small SSD.
We'd load up shapefiles covering a bunch of features across over a quarter million square miles while simultaneously running Chromium in kiosk mode.
I know for a fact some of those had around 1k DAU. I mean, never more than one simultaneously but we did have them! I'm sure it would have handled more just fine if it weren't for the damn laws of physics limiting the number of people that can be physically accessing the hardware at the same time.
That said, we had the benefit of knowing that our user counts would never really increase and due to real-world limitations we'd never be working with a remote database because there was literally no internet available. In general I'd still say colocating your database is not an excuse to allow N+1 query patterns to slip in. They'll be fine for now and just come back to bite you in the ass when your app _does_ outgrow colocating the app/database.
Otherwise both are generally limited to the physical resources of the machine(memory, disk, etc). Generally speaking you can scale boxes much much farther than your data size for most applications.
Seems a bit unfair to call WAL mode emulation of "true" concurrent writes as I'm pretty sure a write-ahead-log (WAL) is exactly how other databases implement multiple concurrent writes. It's just always-on rather than being opt-in.
See [here](https://sqlite.org/wal.html) under "2.2. Concurrency" where it says:
> "However, since there is only one WAL file, there can only be one writer at a time."
SQLite is awesome, I'm a huge fan, but if you need to do lots and lots of writes, then SQLite is not your friend. Luckily that's a very rare application. There is a reason you never see SQLite being the end point for logs and other write-heavy applications.
Bump up the shared buffers, have a multithreaded or pooled driver for your app server so connections are not being opened and closed all the time.
It really 'flies' under such a scenario.
Largest app was a big analytics and drug tracking application for a small pharmaceutical company. Not sure of the size but it was very reporting heavy.
...come to think of it, they may not be the only network device maker that does this.
In my experience, I just cached the living hell out of the apps to avoid the I/O from the database. Providers in those days had huge issues with noisy neighbors, so I/O performance was quite poor.
I want to be clear though, I would never recommend anyone do this. For its time it was impressive but I suspect by now it's been re-engineered away from that model, if for no other reason than its a lot of eggs in one basket.
You definitely don't need an ORM to get to N+1 hell, a for-loop in vanilla PHP will do.
Because by default, no relationships are loaded and the way AR "loads" stuff makes is extremely easy to create N+M queries on 5 lines of ruby. Then, those people will tell you that you run out of database before you run out of ruby...
I'm, like, 1 for 10 maybe.
Where widgets might be "similar items", "people that bought x, bought y", "recent purchases of this widget", "best sellers in this category" and so on.
Even within the database you'll sometimes find examples: stored procedures looping over a cursor and running queries or other statements in each iteration. This doesn't have the same latency implications as N+1 requests over a network, but still does add significant time per step compared to a set based equivalent.
If you think about it, the database will sometimes do this itself when your database is not optimally structured & indexed for your queries (in SQL Server, look at "number of executions" in query plans).
Many moons ago in one of my first coding projects I not only used a separate query for each record, but a separate database connection too. In addition to this, I didn't use a loop to generate these queries, I copy and pasted the connection setup and query code for each record, leading to a 22k code file just for this single web page.
There were probably on the order of 1000 records (they represented premier league football players), and amazingly this code was fast enough (took a couple of seconds to load - which wasn't so unusual back in 2007) and worked well enough that it actually had ~100 real users (not bad a 14 year old's side project).
Joins are one possible workaround to the n+1 problem. However, they push the problem into the application to unpack, which detracts from having a nice high level abstraction in the database for fetching data. The many queries of n+1 would be ideal if you had an ideal database. You are querying on behalf of many different entities so many different queries is conceptually sound. But since we don't have ideal databases...
That's not to say that joins don't have their place. In the relational model, joins are essential. But the n+1 problem happens when you are not working in the relational model.
With LiteFS, we're aiming to easily run on low resource cloud hardware such as nodes with 256MB or less of RAM. I haven't tried it on a Raspberry Pi yet but I suspect it would run fine there as well.
However, "local SQLite vs. remote Postgres/MySQL" remains a false dichotomy when talking about network latency.
I'd pick a centralized network-reachable database with a strong relational schema for write-heavy applications, and a lighter in-process system for something that is mostly reads and where latency matters. It's not a false dichotomy — but more like a continuum that certainly includes both extremes on it.
It's not like the address-space separation is without benefits... heck, if it weren't, you could simply have embedded the whole application inside Postgres and achieved the same effect.
Litestream writes everything to S3 (or similar storage).
LiteFS lets different nodes copy replicated data directly to each other over a network, without involving S3.
In either case, the actual SQLite writes and reads all happen directly against local disk, without any network traffic. Replication happens after that.
That said, unless I've misunderstood the LifeFS use case, you're still going over the network to reach a node, and that node is still going through a FUSE filesystem. That would seem to create overhead comparable (potentially more significant) to talking to a Postgres database hosted on a remote node.
It just doesn't seem that obvious that there's a big performance win here. I'd be curious to see the profiling data behind this.
So your read queries should mostly be measured in microseconds.
You should check out the read latency for read-only requests over unix domain sockets with PostgreSQL. You tend to measure it in microseconds, and depending on circumstances it can be single-digit microseconds.
Regardless of whether your FUSE logic does nothing at all, It sure seems like there's intrinsic overhead to the FUSE model that is very similar to the intrinsic overhead of talking to another userspace database process... because you're talking to another userspace process (through the VFS layer).
When the application reads, those requests go from userspace to the FUSE driver & /dev/fuse, with the thread being put into a wait state; then the FUSE daemon needs to pull the request from /dev/fuse to service it; then your FUSE code does whatever minimal work it needs to do to process the read and passes it back through /dev/fuse and the FUSE driver, and from there back to your application. That gets you pretty much the same "block and context switch" overhead of an IPC call to Postgres (arguably more). FUSE uses splicing to minimize data copying (of course, unix domain sockets also minimize data copying), though looking at the LiteFS Go daemon, I'm not entirely sure there isn't a copy going on anyway. Memory copying issues aside, from a latency perspective, you're jumping through very similar hoops to talking to another user-space process... because that's how FUSE works.
There's a potential performance win if the data you're reading is already in the VFS cache, since that would bypass having to through the FUSE filesystem (and the /dev/fuse-to-userspace jump) entirely. The catch is, at that point you're bypassing SQLite transaction engine semantics entirely, giving you a dirty read that's really just a cached result from a previous read. That's not really a new trick, and you can get even better performance with client-side caching that can avoid a trip to kernel space.
I'm sure there's a win here somewhere, but I'm struggling to understand where.
That's a typical case for a lot of databases. Most of the "hot" data is in a small subset of the pages and many of those can live in the in-process page cache.
> The catch is, at that point you're bypassing SQLite transaction engine semantics entirely, giving you a dirty read that's really just a cached result from a previous read.
It's not bypassing the transaction engine semantics. For WAL mode, SQLite can determine when pages are updated by checking the SHM file and then reading updated pages from the WAL file. Pages in the cache don't need to flushed on every access or even between transactions to be valid.
> I'm sure there's a win here somewhere, but I'm struggling to understand where.
The main goal of LiteFS is to make it easy to globally replicate applications. Many apps run in a single region of the US (e.g. us-east-1) and that's fast for Americans but it's a 100ms round trip to Europe and a 250ms round trip to Asia. Sure, you can spin up a multi-region Postgres but it's not that easy and you'll likely deploy as separate database and application servers because Postgres is not very lightweight.
LiteFS aims to have a minimal footprint so it makes it possible to deploy many small instances since SQLite is built to run on low resource hardware.
As far as comparisons with Postgres over UNIX sockets, I agree that the performance of a single instance is probably comparable with a FUSE layer.
Yes, though most database engines end up managing their own cache and using direct IO, rather than the VFS cache.
> It's not bypassing the transaction engine semantics. For WAL mode, SQLite can determine when pages are updated by checking the SHM file and then reading updated pages from the WAL file. Pages in the cache don't need to flushed on every access or even between transactions to be valid.
That sounds a lot like at least the SHM & WAL checks wouldn't be cached by the VFS, but as I've been looking at the design more carefully, I'm starting to think I understand the idea here. Basically, the SHM & WAL get updated separately, so you might read a stale version, but since you aren't elected to be a writer, that just means you're looking at stale data, not creating an integrity problem.
> The main goal of LiteFS is to make it easy to globally replicate applications. Many apps run in a single region of the US (e.g. us-east-1) and that's fast for Americans but it's a 100ms round trip to Europe and a 250ms round trip to Asia. Sure, you can spin up a multi-region Postgres but it's not that easy and you'll likely deploy as separate database and application servers because Postgres is not very lightweight.
So, I get that multi-region Postgres can be tricky to set up, if you're doing multi-leader, but this seems about as complicated as a "single leader, many followers" set up, and given that what you're trying to do is shave off the hundreds of milliseconds from partially circumnavigating the earth at the speed of light, I'm not sure the perceived performance differences are significant (or really even measurable) compared to fluctuations in network latency of requests to the region-local node.
> LiteFS aims to have a minimal footprint so it makes it possible to deploy many small instances since SQLite is built to run on low resource hardware.
This part I'm getting and the objective a lot of sense to me (and certainly running local postgres instances on every node wouldn't be an obvious approach to me). I hadn't thought of FUSE + SQLite as a way to get there, so this is an interesting and surprising approach. I'm looking forward to how this plays out.
> As far as comparisons with Postgres over UNIX sockets, I agree that the performance of a single instance is probably comparable with a FUSE layer.
Interesting. I was thinking I was missing something. Thanks for all the insight.
1. kernel page cache fully removes the read overhead for in-memory pages
2. There's a FUSE_PASSTHROUGH mode that removes that overhead for all reads & writes. I haven't studied exactly what writes LiteFS needs to observe (just journal vs all data writes), but at least for reads it seems to pass them straight through. We could well submit a FUSE_PASSTHROUGH_READ patch to the kernel, and use that to remove all read overhead. The patch should be trivial, since the full FUSE_PASSTHROUGH mode is there already.
Disclaimer: I wrote the FUSE framework LiteFS uses, https://bazil.org/fuse
https://source.android.com/docs/core/storage/fuse-passthroug... https://lwn.net/Articles/674286/
Even if the reads are happening locally, if they're going through FUSE (even a "thin" pass-through that does nothing), that means they're getting routed from kernel space to a FUSE daemon, which means you're still doing IPC to another process through a kernel...
The gains in latency with sqlite won't matter as soon as throughput starts to dominate.
BeckrockDB is a MySQL compatible interface wrapped around SQLite and as such, is intended for remote RDBMS use.
Expensify has some interesting scaling performance data for their use of it.
https://blog.expensify.com/2018/01/08/scaling-sqlite-to-4m-q...
https://sqlite.org/np1queryprob.html#the_need_for_over_200_s... talks about this in the context of Fossil, which uses hundreds of queries per page and loads extremely fast.
It turns out the N+1 thing really is only an anti-pattern if you're dealing with significant overhead per query. If you don't need to worry about that you can write code that's much easier to write and maintain just by putting a few SQL queries in a loop!
Related: the N+1 problem is notorious in GraphQL world as one of the reasons building a high performance GraphQL API is really difficult.
With SQLite you don't have to worry about that! I built https://datasette.io/plugins/datasette-graphql on SQLite and was delighted at how well it can handle deeply nested queries.
But it doesn't change the fact that standalone database server processes are designed to support specific queries at lower frequencies. This is one of the main points of The SQL language is to load precisely the data that is needed in a single statement.
Relying on this is a design pattern would only scale in specific use cases and would hit hard walls in changing scenarios
I see bad judgement calls coming from small numbers, often due to failing to do the cost x frequency math properly in your head. If you're looking at an operation that takes 3ms or 3μs that is called a million times per minute, or twenty thousand times per request, you don't have enough significant figures there and people make mistakes. 3ms x 12345 = ~40s, not 37035ms. Better if you use a higher resolution clock and find out it's actually 3.259 ms, leading to a total of ~40.23s
Point is, when we are doing billions of operations per second, lots of small problems that hide in clock jitter can become really big problems, and a 5-10% error repeated for half a dozen concerns can lead to serious miscalculations in capacity planning and revenue.
I segment sqlite files (databases) that have the same schema into the same folder. I haven't really had a case where migrations was really a concern, but I could see it happening soon.
Seems like in my deployment, I'm going to need an approach to loop over dbs to apply this change... I currently have a step of app deployment that attempts to apply migrations... but it is more simplistic because the primary RDBMS (postgresql) just appears to the application as a single entity which is the normative use-case for db-migrate-runners.
I wonder how this would present itself in real-world usage? Would the replicas go out-of-date for 10-30s but continue serving read-only traffic, or could there be some element of cluster downtime caused by this?
Yes, a large migration could end up sending a lot of data out to replicas. One way to mitigate this is to shard your data into separate SQLite databases so you're only migrating a small subset at a time. Of course, that doesn't work for every application.
> Would the replicas go out-of-date for 10-30s but continue serving read-only traffic, or could there be some element of cluster downtime caused by this?
Once WAL support is in LiteFS, it will be able to replicate this out without causing any read down time. SQLite is still a single-writer database though so writes would be blocked during a migration.
How can you ensure that a client that just performed a forwarded write will be able to read that back on their local replica on subsequent reads?
A couple years ago someone posted a solution to that here. I'm not sure if it works for SQLite, but it worked for Postgres. The basics of it were that each replica was aware of the latest transaction ID it had seen. On a normal read you'd deal with the usual set of eventually consistent issues. But on a read-after-write, you would select a replica that was ahead of the write transaction.
Ultimately that's a very odd flavor of sharding. What I don't recall is how they propagated that shard information efficiently, since with sharding and consistent hashing the whole idea is that the lookup algorithm is deterministic, and therefore can be run anywhere and get the same result. WAL lag information is dynamic, so... Raft?
100 edits a minute from distinct sessions is 1000 sessions pinned at any moment. If they read anything in that interval it comes from the primary. The only question is what's the frequency and interval of reads after a write.
(The client needs to make sure to consume that in a single pread/read syscall, or it could observe a sheared state.)
Disclaimer: I wrote the FUSE framework LiteFS uses, https://bazil.org/fuse
https://github.com/superfly/litefs/blob/52e269d4b04070690ce2... https://github.com/superfly/litefs/blob/52e269d4b04070690ce2... https://github.com/superfly/litefs/blob/a5cf33d1a3a91873d4ad...
- udp support
- container2vm overhaul
- private networks aka 6pn
- some key flyctl (cli) commands like flyctl ssh, flyctl proxy
- metrics
- litefs
- perhaps, the imminent overhaul of the orchestration layer (?)
- the upcoming authz layer
[1] https://community.fly.io/u/thomas/activity/topics | https://fly.io/blog/author/thomas/
> - litefs
That glory goes to https://news.ycombinator.com/user?id=benbjohnson
Of course, I claim no inside knowledge, so I may very well be mistaken (:
This is just one of many examples: https://www.phoronix.com/news/MTI1MzM
In contrary, in some cases FUSE is even faster than doing a regular kernel mount(). There is an experimental research distribution called distri that is exclusively relying on fuse mounting and figured out that FUSE was faster for doing mounts. https://michael.stapelberg.ch/posts/tags/distri/
The first one is as it's name suggests.. a user space implementation of ZFS.
ZoL (now unified with OpenZFS) is implemented as a kernel module and as such does _not_ run in user space. It performs significantly better as a result.
FUSE is still slow, which is why there's ongoing effort to replace things like NTFS-3G (the default NTFS implementation in most linux distros) with an in-kernel implementation: https://news.ycombinator.com/item?id=28418674
Edit: Also I don't want to imply that FUSE is near in-kernel filesystems, but it is certainly performing much better than 12 years ago.
Disclaimer: I wrote the FUSE framework LiteFS uses, https://bazil.org/fuse -- and I also have some pending performance-related work to finish, there...
Also, what happens if the Consul instance goes down?
If my application nodes can't be ephemeral then this seems like it would be harder to operate than Postgres or MySQL in practice. If it completely abstracts that away somehow then I suppose that'd be pretty cool.
Currently finding it hard to get on board with the idea that adding a distributed system here actually makes things simpler.
Yes, each node has a full copy of the database locally.
> If so, is there another copy somewhere else (e.g. S3) in case all nodes go down?
S3 replication support is coming[1]. Probably in the next month or so. Until then, it's recommended that you run a persistent volume with your nodes.
> What happens if the Consul instance goes down?
If Consul goes down then you'll lose write availability but your nodes will still be able to perform read queries.
> If my application nodes can't be ephemeral then this seems like it would be harder to operate than Postgres or MySQL in practice.
Support for pure ephemeral nodes is coming. If you're comparing LiteFS to a single node Postgres/MySQL then you're right, it's harder to operate. However, distributing out Postgres/MySQL to regions around the world and handling automatic failovers is likely more difficult to operate than LiteFS.
That's a tricky choice of words, since it looks like you lose Consistency, while retaining Availability and Partition Tolerance. If Consul is down everyone reads stale data, but no writes. Right?
Of course, it's harder for Consul to go down than it is for your database to go down, so the Venn Diagram of "Consul unhappy, Database Happy" is fairly heavily populated with "user error". Which is why 'when in doubt, use Consul' is not terrible advice.
> To improve availability, it uses leases to determine the primary node in your cluster. By default, it uses Hashicorp's Consul.
Having a satellite office become leader of a cluster is one of the classic blunders in distributed computing.
There are variants of Raft where you can have quorum members that won't nominate themselves for election, but out of the box this is a bad plan.
If you have a Dallas, Chennai, Chicago, and Cleveland office and Dallas goes dark (ie, the tunnel gets fucked up for the fifth time this year), you want Chicago to become the leader, Cleveland if you're desperate. But if Chennai gets elected then everyone has a bad time, including Dallas when it comes back online.
We don’t support any kind of tiering for candidates in different regions. That’s not a bad idea though.
A centralized database handles consistency, and vends data closures to distributed applications for in-process querying (and those closures reconcile via something like CRDT back to the core db).
Does this exist?
See also: materialize.com, readyset.io, aws elastic-views (forever in preview).
For instance, we run kubernetes on multiple VPS providers, including public clouds with serverless onramp/offramps deployed on edge location. Anything under 15 minutes are processed by serverless. Anything longer is offloaded to one of the VPS containers available in every part of the world.
I have some more feedbacks ready if you are interested, its a neat idea but not exactly as seamless and easy as the idea proposed since public clouds already offer a way to do this.
Are the readable replicas supposed to be long-lived (as in, I don't know, hours)? Or does consul happily converge even with ephemeral instances coming and going every few minutes (thinking of something like Cloud Run and the like, not sure if Fly works the same way)? And do they need to make a copy of the entire DB when they "boot" or do they stream pages in on demand?
Right now, yes, they should be relatively long-lived (e.g. hours). Each node keeps a full copy of the database so it's recommended to use persistent volumes with it so it doesn't need to re-snapshot on boot.
We do have plans to make this work with very short-lived requests for platforms like Lambda or Vercel. That'll use transactionally-aware paging on-demand. You could even run it in a browser, although I'm not sure why you would want to. :)
How does that compare/contrast with what my grandparent is alluding to ?
[1] https://github.com/benbjohnson/litestream/issues/140
[2] https://www.rsync.net/resources/notes/2021-q3-rsync.net_tech...
LiteFS will do the same although I'm only targeting S3 initially. That seemed to be the bulk of what people used.
As with the litestream component, it would be wonderful to have an SFTP transport for LiteFS.
When you say 'S3' I assume you mean "S3 compatible API" so that would be fairly open and portable but SFTP would be even more so.
Thanks again.
CouchDB had this same issue with its database per user model and eventually consistent writes.
Traditionally, it's complicated to replicate your database to different regions of the world using something like Postgres. However, LiteFS aims to make it as simple as just spinning up more server instances around the world. It connects them automatically and ships changes in a transactionally safe way between them.
what makes it complicated?
$ curl -v https://fly.io/blog/introducing-litefs/
* Trying 2a09:8280:1::a:791:443...
* Connected to fly.io (2a09:8280:1::a:791) port 443 (#0)
* ALPN, offering h2
* ALPN, offering http/1.1
* successfully set certificate verify locations:
* CAfile: /etc/ssl/cert.pem
* CApath: none
* (304) (OUT), TLS handshake, Client hello (1):
curl: (35) error:02FFF036:system library:func(4095):Connection reset by peer
Requesting via ipv4 works $ curl -4v https://fly.io/blog/introducing-litefs/
* Trying 37.16.18.81:443...
* Connected to fly.io (37.16.18.81) port 443 (#0)
* ALPN, offering h2
* ALPN, offering http/1.1
* successfully set certificate verify locations:
* CAfile: /etc/ssl/cert.pem
* CApath: none
* (304) (OUT), TLS handshake, Client hello (1):
* (304) (IN), TLS handshake, Server hello (2):
* (304) (IN), TLS handshake, Unknown (8):
* (304) (IN), TLS handshake, Certificate (11):
* (304) (IN), TLS handshake, CERT verify (15):
* (304) (IN), TLS handshake, Finished (20):
* (304) (OUT), TLS handshake, Finished (20):
* SSL connection using TLSv1.3 / AEAD-CHACHA20-POLY1305-SHA256
* ALPN, server accepted to use h2
* Server certificate:
* subject: CN=fly.io
* start date: Jul 25 11:20:01 2022 GMT
* expire date: Oct 23 11:20:00 2022 GMT
* subjectAltName: host "fly.io" matched cert's "fly.io"
* issuer: C=US; O=Let's Encrypt; CN=R3
* SSL certificate verify ok.
* Using HTTP2, server supports multiplexing
* Connection state changed (HTTP/2 confirmed)
* Copying HTTP/2 data in stream buffer to connection buffer after upgrade: len=0
* Using Stream ID: 1 (easy handle 0x135011c00)
> GET /blog/introducing-litefs/ HTTP/2
> Host: fly.io
> user-agent: curl/7.79.1
> accept: */*
>
* Connection state changed (MAX_CONCURRENT_STREAMS == 32)!
< HTTP/2 200
< accept-ranges: bytes
< cache-control: max-age=0, private, must-revalidate
< content-type: text/html
< date: Wed, 21 Sep 2022 16:50:16 GMT
< etag: "632b20f0-1bdc1"
< fly-request-id: 01GDGFA3RPZPRDV9M3AQ3159ZK-fra
< last-modified: Wed, 21 Sep 2022 14:34:24 GMT
< server: Fly/51ee4ef9 (2022-09-20)
< via: 1.1 fly.io, 2 fly.io
<
<!doctype html> ...I don't have a mac, but if the man page[1] is right, something like this should work to see how big of a packet you can successfully send and receive:
ping -6 -D -G 1500 -g 1400 fly.io
(you may need to run as root to set packet sizes). You should get back a list of replies with say 1408 bytes, then 1409, etc. The last number you get back is effectively the largest IPv6 payload you can receive (on this path, it could be different for other paths), and if you add the IPv6 header length of 40, that's your effective path MTU.Use tcpdump to see what TCP MSS is being sent on your outgoing SYN packets, for IPv4, the MSS is MTU - 40 (20 for IPV4 header, 20 for TCP header), for IPv6, the MSS should be MTU - 60 (40 for IPv6 header, 20 for TCP header). If your TCP MSS is higher than the observed path MTU to fly.io, that's likely the immediate cause of your problem.
If you're using a router, make sure it knows the proper MTU for the IPv6 connection, and enable MSS clamping on IPv6, if possible --- or make sure the router advertisement daemon shares the correct MTU.
Hope this gets you started.
Just for reference, the ping command is a little different
sudo ping6 -D -G 1500,1400 fly.io
I set an MSS to 1492 which pfsense (my router) translates to an MSS clamp of 1492-60 for IPv6 and 1492-40 for IPv4. This is a German Deutsche Telekom Fiber connection. Now everything works fine, I can request fly.io (and also discovered that https://ipv6-test.com was not working before and now does with the MSS clamping)Does MSS clamping have any disadvantages? Are there any alternatives in my case?
The only downside to MSS clamping is the computational expense of inspecting and modifying the packets. On a residential connection, where you're running pfsense already, it's probably not even noticeable; but your ISP wouldn't be able to do clamping for you, because large scale routers don't have the processing budget to inspect packets at that level. I've seen some MSS clamping implementations that only clamp packets going out to the internet, and not the return packets... that can lead to problems sending large packets (which isn't always very noticeable, actually; a lot of basic browsing doesn't send packets large enough to hit this, unless you go to a site that sets huge cookies or do some real uploading)
The alternative would be to run a 1492 MTU on your LAN, but that has the marginal negative of reducing your maximum packet size for LAN to LAN transfers.
* Trying 2a09:8280:1::a:791:443...
* Connected to fly.io (2a09:8280:1::a:791) port 443 (#0) ALPN,
* offering h2 ALPN, offering http/1.1 CAfile:
* /etc/ssl/certs/ca-certificates.crt CApath: /etc/ssl/certs TLSv1.0
* (OUT), TLS header, Certificate Status (22): TLSv1.3 (OUT), TLS
* handshake, Client hello (1): TLSv1.2 (IN), TLS header, Certificate
* Status (22): TLSv1.3 (IN), TLS handshake, Server hello (2): TLSv1.2
* (IN), TLS header, Finished (20): TLSv1.2 (IN), TLS header,
* Supplemental data (23): TLSv1.3 (IN), TLS handshake, Encrypted
* Extensions (8): TLSv1.2 (IN), TLS header, Supplemental data (23):
* TLSv1.3 (IN), TLS handshake, Certificate (11): TLSv1.2 (IN), TLS a
* header, Supplemental data (23): TLSv1.3 (IN), TLS handshake, CERT
* verify (15): TLSv1.2 (IN), TLS header, Supplemental data (23):
* TLSv1.3 (IN), TLS handshake, Finished (20): TLSv1.2 (OUT), TLS
* header, Finished (20): TLSv1.3 (OUT), TLS change cipher, Change
* cipher spec (1): TLSv1.2 (OUT), TLS header, Supplemental data (23):
* TLSv1.3 (OUT), TLS handshake, Finished (20): SSL connection using
* TLSv1.3 / TLS_AES_256_GCM_SHA384 ALPN, server accepted to use h2
* Server certificate:
...Same on my macbook
This level of database replication isn't going to put individual rows in different geographies.
What exactly are we talking about here? A WebSQL thats actually synced to a proper RDBMS? Synced across devices? I'm not clear about an end to end use case.
Edit: Honestly, this line from the LiteFS docs[0] needs to be added to the top of the article:
> LiteFS is a distributed file system that transparently replicates SQLite databases. This lets you run your application like it's running against a local on-disk SQLite database but behind the scenes the database is replicated to all the nodes in your cluster. This lets you run your database right next to your application on the edge.
I had no idea what was being talked about otherwise.