5 ms·
From the FAQ <http://hyperdex.org/faq/> http://hyperdex.org/faq/>: "So, the CAP Theorem says that you can only have one of C, A, and P. Which are you sacrifici
by ericmoritz 15y ago
From the FAQ <http://hyperdex.org/faq/> http://hyperdex.org/faq/>:
"So, the CAP Theorem says that you can only have one of C, A, and P. Which are you sacrificing?
HyperDex is designed to operate within a single datacenter. The CAP Theorem holds only for asynchronous environments, and well-administered datacenters enable us to sidestep this tradeoff entirely."
I'd like to see how they pull that off when a node goes down. I guess in "well-administered" data centers, nodes don't go down.
Sounds like they're sacrificing "A" to me because they're doing synchronous replication.
- rescrv 15y agoHyperDex uses value-dependent chaining, which offers fault tolerance properties similar to those provided by chain replication (http://www.cs.cornell.edu/home/rvr/papers/osdi04.pdf http://www.cs.cornell.edu/home/rvr/papers/osdi04.pdf). A single node failure will be recovered from quickly without issue. Multiple concurrent failures are handled the same as the single failure case, so long as our failure assumptions are not violated (e.g., every node in the datacenter fails simultaneously).
- tlb 15y agoThat can work if the server process is killed so that the master is immediately notified. But what about other failure modes? For example, if the disk has soft errors and writes start taking several seconds to complete, the system can't decide in a small amount of time that the node is dead.
- rescrv 15y agoIn this case,the node should report its own failure. It must be able to do so, otherwise it will be deemed faulty by the entity it reports to.
- snewman 15y agoIt's a bit of a middle ground. Yes, the replication is synchronous, which impacts availability. However, the master can remove a failed replica from the chain fairly quickly. In principle, with proper tuning, a node failure would merely cause a brief hiccup. This would feel more like a period of increased latency than a full-blown outage. So there really needn't be much sacrifice of availability. However, there's also a sacrifice of partition tolerance. If the master is unable to communicate with any replica, the system can't serve requests. Also, the master is implemented as a collection of Paxos nodes; if these nodes are partitioned from one another, the entire system would grind to a halt. Since this is intended for intra-datacenter use, one could argue that a full network partition might be unlikely. (Depending on what sort of data center you hang out in.) But in CAP terms, it's possible, of course. (I base all this on the value-dependent chaining paper cited below.)
- rescrv 15y agoWhat you said was right on. I just wanted to add a few things. The coordinator is only involved for recovering from failures, so the cluster can still serve requests until server (non-coordinator) nodes start failing too. I would also add that if there is a intra-datacenter partition so severe as to violate HyperDex's failure assumptions, it will likely impact applications built on top of HyperDex as well. It would be necessary to survive such failures with an inter-datacenter system (which could be built on top of HyperDex).
- ypcx 15y agoWell first Thank you for providing another data storage possibility, and a great one at that. I'd just like to ask - the benchmarks where HyperDex beats even Redis - these are strictly clustered benchmarks - is that true? Or is the way HyperDex stores data so efficient, that it beats Redis even on a single core / single thread? Thanks!
- gizzlon 15y agoWithout more context graphs like that is pretty useless.. For all we know they just invented those numbers. I'm not saying they did, but you get my point.. Those numbers seem way to low for running on the same machine, and if not shouldn't the network be the bottleneck and show similar results for both? I'm sure there's a reasonable explanation, just as I'm sure they picket benchmarks that makes themselves look good.
- nicktelford 15y agoThis sounds like a CP system to me. There's nothing wrong with that btw, I don't know why people are so reluctant to admit this. AP systems have some useful properties, but they're also (typically) more difficult to reason about. The "hiccups" you describe are periods of unavailability. The increased latency is caused by an element of the system waiting for the data to become available again, a totally valid strategy for coping with transient failures/partitions. Your argument about intra-datacenter partitions being unlikely are true, but they do happen. You also make a good point about such partitions also affecting client applications. Both of these are indicative of CP systems and, like I said: there's nothing wrong with that. Personally, I think both AP and CP distributed systems are equally interesting. What I consider a red flag is attempting to rationalize how a system "beats CAP".
- tlb 15y agoWhen a node dies, the master reconfigures all the servers and clients with a new topology excluding the failed node. "Operations which are interrupted by reconfiguration exhibit at-most-once semantics." So while the system is reconfiguring after a node failure, updates can be lost. Time windows of "at most once semantics" mean the system has none of C, A, or P. Which doesn't mean it's not a good database for many purposes.
- rescrv 15y agoThe only scenario in which the operation has "at most once semantics" is when the node the client is directly communicating with fails. No other failure is visible to the clients. Furthermore, every failure scenario provides the following guarantees: * If the result of an operation is visible by one client, it is visible by all clients, always and immediately * Updates to the same key are always applied in the same order on all servers. The presence of "at most once semantics" do not harm our consistency guarantee. In the database world, this would be equivalent to a client sending the final "commit" message, and then losing internet connectivity. In such a scenario, the operation may or may not happen, but the client will not know one way or the other. Edit: Formatting of the list
- tlb 15y agoThe only scenario. In essentially every case of a node failing in an active database, clients will be communicating with it. When talking about durability, you have to assume things are failing in operation.
- rescrv 15y agoHyperDex utilizes value-dependent chains for replication. Updates move forward in the chains, while acknowledgements flow in reverse. To issue a PUT or a GET, the client contacts the head of the chain responsible for the object it is modifying/accessing. If other nodes in the chain fail, the chain will transparently recover. If the point leader fails (the head of the chain), then the client does not know if the operation completed. This is analogous to a database library opening a socket and sending "BEGIN; INSERT INTO data ("x", "y", "z"); COMMIT" and then the client losing connection (or crashing entirely). There is always some point at which the server may complete, and then the client may immediately crash before receiving notification that the operation is complete. Even if this happens, however, HyperDex's GET and PUT operations are linearizable.
- reitblatt 15y ago"I'd like to see how they pull that off when a node goes down. I guess in "well-administered" data centers, nodes don't go down." They offer f-fault tolerance. They can have f nodes go down in a single "zone" and keep chugging as long as no more nodes in the same zone go down before the master reconfigures. Note that the f faults are per-zone, not per-system, so in fact many more than f nodes can be down in a single system without a problem. But, more importantly, you seem to be confusing partition tolerance and fault tolerance. CAP is about partition tolerance: offering "CA" in the presence of arbitrary partitions. They offer a specific form of fault tolerance: "CA" in the presence of any failure or partition that affects less than f nodes.
- lobster_johnson 15y ago> "So, the CAP Theorem says that you can only have one of C, A, and P. Which are you sacrificing? Which is wrong. You can only have two of the three at the same time: CA, CP or AP. If they get something as fundamental as this wrong, you have to wonder about the rest of the project.
- mononcqc 15y agoTechnically you can have at most two of the three of the same time. It's well possible for some stores to provide only one or no property at all!
- lobster_johnson 15y agoWell, yes. But that's not what the FAQ says; it says you can have at most one. The distinction is significant.
- reitblatt 15y agoOr it could just be a typo on the brand new webpage of a brand new project. Which it obviously is if you look at their explanation of CAP.
- rescrv 15y agoJust a simple typo. I've fixed it. Thanks for pointing this out.
- xijhing 15y agoPACELC[1] helps distinguish just the kind of CAP sacrifice this makes. [1] http://dbmsmusings.blogspot.com/2010/04/problems-with-cap-and-yahoos-little.html http://dbmsmusings.blogspot.com/2010/04/problems-with-cap-an...