5 ms·
I'm a beginner in this topic and I find this topic interesting. I really want there to be an easy-to-deploy consistency solution. If I have a distributed micro
by samsquire 3y ago
I'm a beginner in this topic and I find this topic interesting. I really want there to be an easy-to-deploy consistency solution.
If I have a distributed microservice architecture and I want to keep multiple datastores in synchronization or "consistent" what's the industry best practice?
A few days ago I was trying to solve the inconsistency problem with "settled timestamps" which is a kind of multiversioning idea except that timestamps elapsed with the absence of reported error represent a valid save/commit. Kind of like two phase commit with the second phase being time. The idea is that we watch the clocks of other servers and if they don't update then we know we cannot trust their settled timestamps. (My intent was to allow scaling consistency across many servers, because we don't need to wait for response for every update, we only need to wait for the next timestamp interval)
Here's my Multithreaded multiprocessing Python code to test indeterminancy. 10 threads all send eachother random updates. They also broadcast their own timestamp and the timestamps of their own perspective of the timestamps every other server.
https://replit.com/@Chronological/InconsistencySimulation#main.py https://replit.com/@Chronological/InconsistencySimulation#ma...
(click Run and watch the output, you'll have to wait 10 seconds)
A read in this simulation is the MIN of all timestamps of all servers reported timestamps.
10 seconds into the simulation, we ask every thread for its own perspective of what the counter value is. Sometimes they will all report the same value, a lot of the time they shall be split brained.
I am aware that wall clock timestamps are not suitable for ordering in a distributed system and that logical or vector clocks should be used for ordering.
If you can get the simulation to all report the same number at any point in time, then that would be great :-)
Ordering in distributed systems is significant, as the eventual consistency of the simulation means that some values can arrive late but affect the value, meaning it is not linearizable. Bloomlang tries to solve this.
I'm specifically interested in scaling WITH consistency but I think this is quite difficult.
- rawgabbit 3y agoFor distributed systems, the main idea is a centralized write ordering journal that is replayed by individual nodes. Multiple systems write sequentially to the central journal. The journal is simply taking requests like a key value store. The journal is replicated to all the nodes. The nodes read from the journal and performs the complex logic requested.
- nine_k 3y agoIf business software practices is not enough of a proof, look at any massive online games. They all use one central server as a source of truth about the game world, and broadcast that state to the clients. Anything a client reports that diverges from the central server view is either corrected, rejected, or becomes a reason to disconnect the client for cheating attempts. If you need strict order, that order should happen in strictly one place. (The universe itself does not support strict order at a distance, as Special Relativity shows.)
- bawolff 3y agoBitcoin: Am I joke to you? Everyone: yes.
- rusk 3y agoBitcoin is an interesting exception to the above. Issue is latency.
- rcxdude 3y agoThere's no performance gain from bitcoin's approach though. The distributed consensus is there for trustless operation, not performance. It's many orders of magnitude slower (and a few more orders of magnitude more power inefficient!) than just having one server handle it.
- acuozzo 3y agoIs your inter-machine messaging asynchronous and is it possible for one of your machines to crash?
- layer8 3y agoI’d suggest to re-evaluate if you really, really need (a) distributed data stores and (b) synchronous consistency. Things become much simpler if you can forego one of them.
- infogulch 3y agoI suspect TigerBeetle DB will be the industry benchmark for consistent, high throughput, fault tolerant, distributed databases in 5 years.
- theptip 3y agoLook at the Raft protocol. In general you’d rather not integrate at the protocol layer; it’s standard to use a consistent store like etcd that implements Raft for the coordination/data that needs to be serializable. (Kubernetes uses etcd so it scales pretty well for a strongly-consistent k/v store.) You said “multiple datastores” so I’m assuming you have heterogenous data and something like CockroachDB isn’t an option. > I'm a beginner in this topic and I find this topic interesting. Not trying to gatekeepe but rolling your own is dangerous. See https://aphyr.com/ https://aphyr.com/ for the gold standard in testing (great educational material). You can use Jepsen to test your distributed systems. But better to just use datastores that Kyle has shown are solid.
- samsquire 3y agoThanks for your reply. I've experimented with a toy Raft implementation but I haven't Jepsen tested that and it's incomplete I did write a Jepsen test for a different eventually consistent protocol which understandably fails the linearizability test because I'm still learning - eventually consistent is not linearizable. https://GitHub.com/samsquire/eventually-consistent-mesh https://GitHub.com/samsquire/eventually-consistent-mesh I want to have my cake and eat it too. Scalability and consistency.
- deleted 3y ago[deleted]
- ris 3y ago> If I have a distributed microservice architecture and I want to keep multiple datastores in synchronization or "consistent" what's the industry best practice? Not to use a distributed microservice architecture.
- mhuffman 3y agoThis is the answer. But if they are going to, I would recommend a book called "Designing Data-Intensive Applications" by Martin Kleppmann. It is a book with the most clear communication on this topic and others that I have ever read. Even if you have a degree in CS, I think it is eye-opening to see how very complex ideas can be communicated vs the same ideas in textbooks.
- matthewsinclair 3y agoSo much this. I was part of an engineering team that built a payments switch from scratch in the early oughts. We built it in Java, on commodity hardware and OS, on top of our own replicated, stateful, distributed computing platform. This was a bonkers thing to do then, pre-cloud. It’s probably still a bonkers thing to do _today_. Anyway, we did a lot of work with 2PC and other consensus mechanisms and came to the conclusion they 2PC wasn’t up to scratch for what we needed (it’s actually provably less than ideal). We ended up building (again, from scratch) one of the earliest (that I know of) implementations of the Virtual Synchrony protocol. VS has some robust maths behind it that you can use to make some stronger consistency claims than you can with 2PC. These are important when you’re dealing with interbank settlements for payments switching. If we started again today I’d say that we might use something like Raft, but I’ve been away from the space for ~20 years now so I’m not entirely sure. However, one thing I do know is this: distributed consensus is Very Hard™ to get right. If the answer to your question involves any kind of multiple master implementation using distributed consensus I can unequivocally guarantee that (other than in a few very specific circumstances that you almost certainly do not have) you’re asking the wrong question.