4 ms·
Consistency is only an issue if you switch database servers mid session. Just "stick" to a node by issuing all subsequent requests to the same server and there
by quinthar 10y ago
Consistency is only an issue if you switch database servers mid session. Just "stick" to a node by issuing all subsequent requests to the same server and there are no consistency issues.
- thesmallestcat 10y agoAnd then laugh maniacally when the partition heals. What's a "session"? Can you explain commitCount more, and how a client is supposed to use it intelligently without deadlocking. Also there's a lot of master/slave talk and I'm wondering how leaders are elected during a partition, and what recovery looks like.
- quinthar 10y agoHeh, there's a lot of detail to be captured I agree. But in short: - All nodes connect to all other nodes - Each node has a priority; a Paxos algorithm is used to identify the highest priority node, which "stands up" to be the master - All nodes respond to read queries using the local database. - So long as you always talk to the same node, there are no consistency issues: you are guaranteed that each request will answered with a database at least as fresh as the last. - However, if you switch databases (eg, a node goes down and you are forced to go to a different one that might not be as fresh), it will wait until it is as fresh as the node you were using - Write requests are escalated to the master, which processes them with a two-phase commit distributed transaction - By default, the master waits for a "quorum" of the cluster to approve the transaction before committing - However, this limits write capacity to 1/median(rtt) of the slaves - Accordingly, we have a "selective synchronization" feature where you can optionally designate some transactions as requiring "less consistency". For example, you can specify that the master should commit when at least one other node approves. - For the highest performance, you can designate a transaction as "asynchronous" and the master will commit immediately - For example, when someone is reimbursed, we absolutely want to get quorum approval of the transaction: we can't risk losing that if the master crashes. But if someone just adds a report comment, we can risk losing it upon master crash. - When the master does go down (either gracefully or not), the next highest priority node steps up seamlessly - Any command escalated to the master but not processed will be re-escalated to the new master - When a higher priority node comes up, it synchronizes to get up to date with everything it missed, and then becomes master Anyway, lots of details here and I agree they haven't been well captured in the site yet. Coming soon!