5 ms·
Scaling does not imply availability. Sharding user data certainly is a good strategy to achieve scalability (i.e., support a large number of users, distributed
by jonburs 16y ago
Scaling does not imply availability. Sharding user data certainly is a good strategy to achieve scalability (i.e., support a large number of users, distributed across a fleet of commodity hardware), but it doesn't a priori make that data highly available.
Let's imagine my personal netflix data (queue, ratings, etc.) live on a lightly loaded (since things are well shareded) DB instance. What happens if the hardware hosting that DB fails? Ok, let's replicate to a slave, and put into place automated failover mechanisms to promote the slave if the master goes down (since with many shards it's really hard to effect a failover manually in a timely fashion). Should the master and slave be in the same datacenter? That after is all is the failure scenario Netflix originally wanted to address -- what happens if your single datacenter goes down (which eventually will happen, no matter how much redundancy you think it may have)?
Fine then, separate the master and slave in different DCs (with fat and fast pipes so 2PC or synchronous replication can keep things consistent without impacting latency too much). Now where does the auto-promotion logic live -- the DC with the slave or the master? What happens if my client can reach both DCs, but there are connectivity issues on the links connecting the DCs? Network hardware can fail like any other kind, backhoes accidentally sever buried optics, and routing protocols do take time to converge.
P is now rearing its ugly and inevitable head, and you're faced with the C/A decision -- can I see and alter my queue (leading to inconsistency) or is it unavailable? Is this scenario pathological?
- lsd5you 16y agoSorry, but this argument is the worst kind of wrong - it is half true. 100.00000% consistency/availability is of course impossible in any circumstances. 'It is only relevant for a few dozen companies' ... because of size/volume. Your argument would make it relevant to all companies, unless netflix and other large companies have to live to a higher standard - which i do not think is being argued. The CAP theory really is about bottlenecks, high volumes of possibly conflicting data to different nodes cannot be synchronised with guarantees. It is possible to have distributed transactions on lightly loaded servers. These could complete in a timely fashion (<10 secs) 99.9% of the time with, theoretically a geometric drop off for the probability of longer delays. The answer to the disconnected data centre is to have 3 with majority rule. Then the loss of a single link can no longer break consistency. What is more, the article talks about the data being internally inconsistent (dangling references... etc.) this seems to be totally unnecessary for something like private user data. For the rest (recommendations), as discussed other strategies can be employed that at least guarantee internal consistency of a node.