Let's build a distributed Postgres proof of concept
notes.eatonphil.com
notes.eatonphil.com
But to make it scalable and production grade, you have to run multiple Raft operations in parallel and manage them well. E.g., sharding them correctly to make the load balanced and adjusting dynamically.
Here's a pleasant article to read: https://en.pingcap.com/blog/how-tikv-reads-and-writes/
https://static.googleusercontent.com/media/research.google.c...
> The key enabler of these properties is a new TrueTime API and its implementation. The API directly exposes clock uncertainty, and the guarantees on Spanner’s timestamps depend on the bounds that the implementation provides. If the uncertainty is large, Spanner slows down to wait out that uncertainty. Google’s cluster-management software provides an implementation of the TrueTime API. This implementation keeps uncertainty small (generally less than 10ms) by using multiple modern clock references (GPS and atomic clocks).
Doesn't Google Spanner use Paxos instead of Raft? From the link:
> Our Paxos implementation supports long-lived leaders with time-based leader leases...
Another approach could be using PostgreSQL as a stateless scalable query execution engine, while stubbing out the storage layer with a remote replicated consistent storage solution.
This is also what I think Aurora does.
Also neon:
YugabyteDB splits tables into tablets (partitions), every tablet has its own raft. It’s a modified rocksdb under the hood.
OTOH, it’s cool to see how someone builds something like this from scratch and is willing to talk about it.
It looks like the apply method is responsible for writing the sql statements to the raft log as well as executing the sql queries. Is waiting for a quorum of writes to the raft log by the other members not needed? Or is this all just handled under the hood by the raft libraries being used.
Architecturally (in rqlite's case) a node looks something like this: https://docs.google.com/presentation/d/1Q8lQgCaODlecHa2hS-Oe...
It looks like Phil's post uses boltdb for the Postgres storage engine as well as for Raft log(via Hashicorpo's implementation Raft lib.)
Thanks for the link to the slides as well. I've seen rqlite mentioned a few times in the last few week and so it was on my short list of things to read up on.
https://github.com/otoolep/hraftd
This is kind of a reference use of Hashicorp's Raft.
It's like how TCP doesn't guarantee that your data will be written or processed or transmitted without corruption; only that it will be received without corruption. You will still get corruption after recv() and before send(), and you have to handle those cases, or your application will begin to introduce unknown data corruption. And when that happens you need the tooling and telemetry to address it after the fact, "distributedly".
I don't know for sure how this situation is supposed to be handled in a production system though.
So it's actually pretty simple for a given node to contact the Leader. If a node receives a request which must be performed on the Leader, and that node is not itself the Leader, it can do one of the following things:
1) reject the request with an error, but this isn't really a production-viable option.
2) reject the request with an error, but tell the client where the Leader can be found, so the client can retry the request, this time sending the request to the Leader.
3) transparently forward the request to the leader, wait for the Leader to execute the request, get the response, and return the response to the client. In this case the client doesn't even know the forwarding to the Leader happened.
rqlite supports mode 2 and 3, client can choose which behavior it wants, on a request-by-request basis. Option 3 is the default.
https://github.com/rqlite/rqlite/blob/master/DOC/DATA_API.md...
one idea would be, the db is completely synchronized on 3 redundant servers, so it's decentralized
each column of a table is stored on a different server, so it's distributed...?
^ semi off topic, but love the idea of distinguishing certain projects as 'just glue'. (Not a dis, glue matters). Especially interesting in the context of OSS tools whiteboxed by cloud vendors with 'proprietary glue'
Maybe this is a bit nit picky, but, still.
[0] https://datastation.multiprocess.io/blog/2022-02-08-the-worl...
For instance, a lot of Postgres tools use "pg_catalog" to do introspection on the datasource. Supporting pg_catalog is a bit of a pain -- this is the reason why Materialize doesn't work with Hasura.
CockroachDB also just barely doesn't work with Hasura OOTB because of a handful Postgres-specific functions and some metadata.
https://fly.io/blog/all-in-on-sqlite-litestream/
"The upcoming release of Litestream will let you live-replicate SQLite directly between databases, which means you can set up a write-leader database with distributed read replicas. Read replicas can catch writes and redirect them to the leader"
It's basically the same thing as monolith vs microservices but extending monolith to the data persistence layer. With horizontally scaling apps being the predominant architecture right now, I don't really see Litestream changing much.
If you're going to horizontally scale Litestream with a multiple writers you're going to end up introducing all the network and synchronization pieces Postgres architectures already have.
[1] https://github.com/jackc/pgproto3
[2] https://github.com/rqlite/rqlite
Disclaimer: I am the creator of rqlite.
Can you elaborate on this, in what sense does rqlite support distributed transactions?
My approach when learning new protocols like Raft or Paxos is to implement them in Pluscal (TLA+'s higher-level language) or P (https://github.com/p-org/P). I've found that helps separate the protocol-level concerns from the implementation-level concerns (sockets? wire format?) in a way that reduces the difficulty of learning the protocol.