When we talk about fault tolerant distributed computing, using the state machine replication approach, it may seem obvious that a system of this kind should be capable of completely masking failures. In fact, however, this is not the case. Our ability to hide failures is really very limited.
When developers use state machine replication techniques (SMR), the usual approach is to replace components of systems or distributed services with groups of N members, and then use some sort of library that delivers the same inputs to each, in the same order. If the replicated component is deterministic, and if care is taken to initialize each component in a manner properly synchronized with respect to its peers, this is enough to guarantee that the copies will remain synchronized. Thus, we have an N-replica group that seemingly will tolerate N-1 faults.
Unfortunately, theory is one thing and reality can be quite a different matter. When people first began to experiment with SMR in the 1990's, developers quickly noticed that because software bugs are a major cause of failure, perfect replication will replicate many kinds of faults! Over time, a more nuanced approach emerged, in which the various replicas are proactively shut down and restarted in an uncoordinated way, so that on average there would still be N copies, but at any instant in time there might be N-1 copies, with one copy shutting down or rejoining. The trick is to transfer the current state of the group to the recovering member, and is solved using the virtual synchrony model, in which group membership advances through a series of epochs, reported via view upcall notifications, with state transfers performed during epoch transitions.
The benefit of this sort of staggered restart is to overcome so-called Heisenbugs. The term refers to bugs that are hard to pin down: they could cause non-deterministic behavior (in which case the replicas might diverge), or bugs that seem to shift around when the developer tries to isolate them.
A common form of Heisenbug involves situations where a thread damages a data structure, but the damage won't be noticed until much later, at which point any of a number of other threads could try to access the structure and crash. Thus the failure, when it occurs, is associated with logic remote from the true bug, and may jump about depending on scheduling order. If the developer realizes that the root cause is the earlier damage to the data structure, it generally isn't too hard to fix the problem. But if the developer is tricked into thinking the bug manifested in the code that triggered the crash, any attempts to modify that logic will probably just make things worse!
The reason that staggered restart overcomes Heisenbugs is that a restarting program will load its initial state from some form of checkpoint, hence we end up with N copies, each using different operations to reach the same coordinated state as the other N-1. If the data-structure corruption problem isn't a common thing, this joining process is unlikely to have corrupted the same data structure as did the others. With proactive restart, all N copies may be in equivalent yet rather different states. We can take this form of diversity even further by load-balancing read-requests across our N copies: each will see different read operations and this will be a further source of execution diversity, without mutating states in ways that can cause the N replicas to diverse.
With such steps, it isn't hard to build an ultra-resilient SMR service, that can remain alive even through extremely disruptive failure episodes. But can such a service "mask" failures?
The answer is yes and no.
On the "yes" side we find work by Robert Surton at Cornell, who created a clever little TCP fail-over solution called TCP-R. Using this protocol, a TCP connection can seamlessly roll from one machine (the failed server) to another. The cleverness arises because of the way that TCP itself handles acknowledgements: in Surton's approach, a service member accepts a TCP connection, reads the request (this assumes that the request size is smaller than the TCP window size, in bytes), replicates the received request using an SMR multicast, and only then allows TCP to acknowledge the bytes comprising the request.
If a failure disrupts the sequence, TCP-R allows a backup to take control over the TCP session and to reread the same bytes from the window. Thus the service is guaranteed to read the request at least once. A de-duplication step ensures that a request that happens to be read twice won't cause multiple state updates.
Replies back to the end user are handled in a similar way. The service member about to send the reply first notifies the other members, using a multicast, and only then sends the reply. If a failure occurs, one of the other members claims the TCP endpoint and finishes the interrupted send.
With TCP-R the end-user's experience is of a fully masked failure: the application sends its request, and the service definitely gets the request (unless all N members crash simultaneously, which will break the TCP session).
Lacking TCP-R, the situation is quite a bit more complex. In effect, the end-user would need to send the request, but also be prepared to re-send it if the endpoint fails without responding. For read-only requests, the service can just perform the request multiple times if it shows up multiple times, but updates are more complex. For these, the service would need logic to deduplicate requests: if the same request shows up twice, it should resend the original reply and not mutate the service state by performing the identical operation a second time. TCP-R masks the service failure from the perspective of the client, although the service itself still needs this form of deduplication logic.
On the "no" side of the coin, we need to consider the much broader range of situations that can arise in systems that use SMR for fault-tolerance. In particular, suppose that one SMR-replicated service somehow interacts with a second SMR-replicated service. Should all N members of the first replica group repeat the same request to the M members of the second group? Doing so is clearly the most general solution, but runs into the difficulty that bytes will be transferred N times to each member: a high cost if we are simply trying to mask a rare event!
Eric Cooper studied this question in his PhD thesis on a system called Circus, at Berkeley in the 1990's. Basically, he explored the range of options from sending one request from group A to group B, but reissuing the request if the sender in A failed or the receiver in B, all the way to the full N x M approach in which every member of A multicasts every request to every member of B, and the members of B thus receive N copies and must discard N-1 of them in the usual case. (TCP-R can be understood as an instance of the first approach, but with the client-side logic hidden under TCP itself, so that only the server has to be aware of the risk of redundancy, and so that it only arises when a genuine failure occurs.)
Cooper pointed out that even with high costs for redundancy, the N x M approach can often outperform any scheme that waits to sense a failure before retrying. His insight was that because detecting a failure can be slow (often 30s or more), proactively sending multiple copies of each request will generally produce extra load on the receivers, but with the benefit of ensuring a snappy response because at least some receiver will act promptly and send the desired reply with minimal delay.
Cooper's solution, Circus, is best understood as a design pattern: a methodology that the application developer is expected to implement. It involves a multicast from group A to group B, code in group B to remember recent responses to requests from A, and logic to de-duplicate the request stream, so that whether B receives a request 1 time or N times, it behaves identically and responds to A in the same manner.
In Derecho, we don't currently offer any special help for this heavily redundant approach, but all the needed functionality is available and the design pattern shouldn't be hard to instantiate. But in fact, many Derecho users are more likely to use a non-fault-tolerant approach when building a data processing pipeline. More specifically, while Derecho users would often use replication to make the processing elements of the pipeline fault-tolerant, they might decide not to pay the overhead of making the handoff of requests, stage by stage in the pipeline, ultra-reliable.
The reason for this compromise is that in the IoT settings where a "smart memory service" might be used, most sensors are rather redundant, sending photo after photo of the same car on the highway, or location after location for the cat in the apartment. The service receives this heavily duplicative input and will actually start by sorting out the good content and discarding the replicated data. Thus we suspect that most Derecho users will be more concerned with ensuring that the service itself is highly available, and less concerned with ensuring that every single message sent to the service is processed.
Indeed, in most IoT settings, freshness of data is more important that perfect processing of each and every data point. Thus, if some camera creates photo X of a vehicle on the highway, and then photo Y, and X is somehow lost because of a failure that occurs exactly as it is being sent, it often would make more sense to not fuss about X and just focus on processing request Y instead.
Microsoft has a system, Cosmos, in which a pipeline of processing is done on images and videos. It manages without fault-tolerant handoff between stages because failures are rare and, if some object is missing, there is always a simple recipe to create it again from scratch. Facebook apparently does this too. Both systems need to be extra careful with the original copies of photos and videos, but computed artifacts can always be regenerated. Thus, perfect fault tolerance isn't really needed!
Of course, one can easily imagine systems in which each piece of data sent to the service is individually of vital importance, and for those, an open question remains: is it really necessary to hand-code Cooper's Circus design pattern? Or could there be a really nice way to package the Circus concept, for example by using a higher level language to describe services and compiling them down to SMR replicas that talk to one-another redundantly?
I view this as an open research topic for Derecho, and one we may actually tackle in coming years. Until then, Derecho can certainly support a high quality of adaptation after crashes, but won't seamlessly hide crashes from the developer who works with the technology. But on the other hand, neither does any other technology of which I'm aware!
Showing posts with label availability. Show all posts
Showing posts with label availability. Show all posts
Monday, 31 July 2017
Friday, 17 March 2017
The CAP conjecture is dead. Now what?
CAP has been around now for something like 15 years. There are some circles in which the acronym is even more famous than FLP or SAT! But CAP was never a theorem (let's call it a "folk-theorem"), and by now we have more and more examples of systems that violate CAP.
Yet even if CAP is false, it was also a fantastic rule of thumb that was incredibly valuable to developers tasked with building high performance distributed systems on the cloud. If we start to teach students that CAP is false, what can replace it in this role?
A little historical context: CAP is short for "you can only have two from Consistency, Availabilty and Partition Tolerance", This assertion was initially put forward by Eric Brewer in a PODC keynote talk he gave in 2000. The justification he offered ran along the following lines. First, he observed that in today's increasingly globalized web systems, we invariably deploy services with high latency WAN links between them. These are balky, so services are forced to make a choice: either respond right now based on data locally available, or await restoration of the WAN link. A tradeoff that obviously favors availability over consistency, and already tells us that many web services will have to find a way to prevent responses that reflect stale data from being noticed by users, or if they are noticed, from causing problems. He used the term "partition tolerance" for this kind of fault-tolerance (namely, giving a response even when some link is down). Hence the P in CAP.
He suggested that we are favoring "A and P over C".
Then he observed that even in a single data center, if you want the highest possible levels of performance, you'll need to handle web requests on edge nodes that take their data from cache, without first talking a backend server first: again weaker consistency, but higher availability. So again, we see the A, as a kind of synonym for rapid response.
So he looked at all the different pairing: C and A over P, C and P over A, A and P over C. He concluded that in practice we always seem to find that by taking A and P we get the best scalability. But CAP itself asserted that although other mixes work, you always have to pick the two you like best, and you'll erode the third.
Although the definitions of A and P are a bit strange (P seems to have a bit of A mixed in, and also seems to sometimes mean fault-tolerance, since partitions never really arise within a data center), CAP is sort of catchy. But like many such acronyms, much depends on the details: if you try and give rigorous definitions, to turn CAP into a theorem (people have done so), you find that it only holds in some narrow situations.
The result is that as professors teaching cloud computing, we generally treat CAP as a clever acronym, but one that works mostly as a rule of thumb. In some sense, CAP is a useful kind of folklore.
CAP took hold back in 2000 because at that time, companies like eBay and Amazon were struggling with the high costs of ACID transactions in systems that the database SQL programming model. Scalable database performance poses issues that are much more nuanced than the ones Eric had in mind: there were puzzles of lock conflict, complexity, data pipelines with non-trivial asynchronous ordering requirements, etc. But the bottom line is that performance of the big systems at eBay and Amazon was erratic and often strayed outside the magic 100ms target for web-service and web-page responses. This is the point at which a human user feels that the system is "snappy" and everyone wants to be in the sub-100ms range.
So, the technology leaders at eBay began to tell their developers that it was absolutely fine to write SQL code as a starting point in web application design (a pragmatic decision: they people they were hiring mostly had extensive SQL experience from their database courses), but that once the SQL code was working, to weaken the transactions by turning them into a series of individual atomic actions. eBay began to talk about the resulting development methodology using a new acronym: BASE, by which they meant "Basically Available, Softstate systems with Eventual Consistency."
Notice that BASE doesn't abandon consistency. Instead, it points out to the developer that many web systems just don't need ACID guarantees to work correctly. ACID and SQL are great for creating a quick prototype that will be robust and easy to debug, but then you "optimize" it by taking away the ACID properties, without breaking the behavior in ways that violate the specification.
Amazon embraced BASE too. Around 2006, they decided to rebuild many of their core applications around a key-value technology called Dynamo, but SQL users found it hard to use, and by 2008 the adoption of Dynamo began to falter. To make the transition easier, Amazon layered in a NoSQL API for Dynamo, called Dynamo-DB: now SQL code could run on Dynamo, but with weaker guarantees than for a full SQL system (for example, NoSQL lacks join operations), and Dynamo-DB happened to be especially well-matched to BASE.
So you can see from these examples why CAP would be such convenient way to motivate developers to make the switch: it more or less tells them that if they optimize their code using BASE, it won't scale properly. Moreover, the eBay memos about BASE include step by step instructions to explain precisely how to go about doing it.
Today, fifteen years later, we've discovered more and more ways to implement scalable, high performance cloud services, edge caches that maintain coherence even at global scale, fully scalable transactional key-value storage systems, and the list goes on. Derecho is one example: it helps you build systems that are highly Available, Replicated and Consistent. Call this ARC.
The cool twist is that with ARC, you get lock-free strong consistency, right in the cloud edge! This is because we end up with very fast replication at the edge, and every replica is consistent. You do need to think hard about your update sources and patterns if you hope to avoid using locks, but for most applications that aspect is solvable because in the cloud, there is usually some sense in which any particular data item has just once real update source. So there is an update "pipeline" and once you have consistency in the system, you end up with a wide range of new options for building strongly consistent solutions.
An example I like very much was introduced by Marcos Aguilera in Sinfonia. In that system, you can make strongly consistent cache snapshots, and run your code on the snapshot. At massive scale you have as many of these snapshots as needed, and transactions mostly just run at the edge. But when you want to update the system state, the trick he suggested is this: run your transaction and keep track of what data it read and wrote (version numbers). Now instead of just responding to the client, generate a "minitransaction" that checks that these version numbers are still valid, and then does all the writes. Send this to the owner of the true database: it validates your updates and either commits by applying the updates, or rejects the transaction, which you can then retry.
Since so much of the heavy lifting is down when speculatively executing the transactions at the first step, the database owner has way less compute load imposed on it. Then you can start to think about systems with sharded data that has different owners for each shard, and limits the transactions to run within a single shard at a time (the trick is to design the sharing rule cleverly). Sinfonia and the follow-on systems had lots of these ideas.
ARC creates a world in which solutions like the Sinfonia one are easy to implement -- and the Sinfonia approach is definitely not the only such option.
I'm convinced that as we move the cloud towards the Internet of Things services, developers will need this model, because inconsistency in a system that controls smart cars or runs the power grid can wreak havoc and maybe even would be dangerous.
So now I want to claim that ARC is the best BASE story ever! Why? Well, remember that BASE is about basically available, soft-state systems with eventual consistency. I would argue that in examples like the Sinfonia one, we see all elements of the BASE story!
For example, a Sinfonia system is basically available because with enough consistent snapshots you can always do consistent read-only operations at any scale you like. The state is soft (the snapshots aren't the main copy of the system: they are replicas of it, asynchronously being updated as the main system evolves). And they are eventually consistent because when you do make an update, the mini-transaction validation step lets you push the update results into the main system, in a consistent way. It might take a few tries, but eventually, should succeed.
What about the eventual consistency aspect of BASE in the case of ARC? ARC is about direct management of complex large-scale strongly consistent replicated state. Well, if you use ARC at the edge, back-end servers tend to be doing updates in a slightly time-lagged way: batched updates and asynchronous pipelines improve performance. Thus they are eventually consistent, too: the edge queues an update, then responds to the end-user in a consistent way, and the back end catches up soon afterwards.
And ARC isn't the only such story: my colleague Lorenzo Alvisi has a way to combine BASE with ACID: he calls it SALT, and it works really well.
So, feeling burned by CAP? Why not check out ARC? And perhaps you would like a little SALT with that, for your databases?
Yet even if CAP is false, it was also a fantastic rule of thumb that was incredibly valuable to developers tasked with building high performance distributed systems on the cloud. If we start to teach students that CAP is false, what can replace it in this role?
A little historical context: CAP is short for "you can only have two from Consistency, Availabilty and Partition Tolerance", This assertion was initially put forward by Eric Brewer in a PODC keynote talk he gave in 2000. The justification he offered ran along the following lines. First, he observed that in today's increasingly globalized web systems, we invariably deploy services with high latency WAN links between them. These are balky, so services are forced to make a choice: either respond right now based on data locally available, or await restoration of the WAN link. A tradeoff that obviously favors availability over consistency, and already tells us that many web services will have to find a way to prevent responses that reflect stale data from being noticed by users, or if they are noticed, from causing problems. He used the term "partition tolerance" for this kind of fault-tolerance (namely, giving a response even when some link is down). Hence the P in CAP.
He suggested that we are favoring "A and P over C".
Then he observed that even in a single data center, if you want the highest possible levels of performance, you'll need to handle web requests on edge nodes that take their data from cache, without first talking a backend server first: again weaker consistency, but higher availability. So again, we see the A, as a kind of synonym for rapid response.
So he looked at all the different pairing: C and A over P, C and P over A, A and P over C. He concluded that in practice we always seem to find that by taking A and P we get the best scalability. But CAP itself asserted that although other mixes work, you always have to pick the two you like best, and you'll erode the third.
Although the definitions of A and P are a bit strange (P seems to have a bit of A mixed in, and also seems to sometimes mean fault-tolerance, since partitions never really arise within a data center), CAP is sort of catchy. But like many such acronyms, much depends on the details: if you try and give rigorous definitions, to turn CAP into a theorem (people have done so), you find that it only holds in some narrow situations.
The result is that as professors teaching cloud computing, we generally treat CAP as a clever acronym, but one that works mostly as a rule of thumb. In some sense, CAP is a useful kind of folklore.
CAP took hold back in 2000 because at that time, companies like eBay and Amazon were struggling with the high costs of ACID transactions in systems that the database SQL programming model. Scalable database performance poses issues that are much more nuanced than the ones Eric had in mind: there were puzzles of lock conflict, complexity, data pipelines with non-trivial asynchronous ordering requirements, etc. But the bottom line is that performance of the big systems at eBay and Amazon was erratic and often strayed outside the magic 100ms target for web-service and web-page responses. This is the point at which a human user feels that the system is "snappy" and everyone wants to be in the sub-100ms range.
So, the technology leaders at eBay began to tell their developers that it was absolutely fine to write SQL code as a starting point in web application design (a pragmatic decision: they people they were hiring mostly had extensive SQL experience from their database courses), but that once the SQL code was working, to weaken the transactions by turning them into a series of individual atomic actions. eBay began to talk about the resulting development methodology using a new acronym: BASE, by which they meant "Basically Available, Softstate systems with Eventual Consistency."
Notice that BASE doesn't abandon consistency. Instead, it points out to the developer that many web systems just don't need ACID guarantees to work correctly. ACID and SQL are great for creating a quick prototype that will be robust and easy to debug, but then you "optimize" it by taking away the ACID properties, without breaking the behavior in ways that violate the specification.
Amazon embraced BASE too. Around 2006, they decided to rebuild many of their core applications around a key-value technology called Dynamo, but SQL users found it hard to use, and by 2008 the adoption of Dynamo began to falter. To make the transition easier, Amazon layered in a NoSQL API for Dynamo, called Dynamo-DB: now SQL code could run on Dynamo, but with weaker guarantees than for a full SQL system (for example, NoSQL lacks join operations), and Dynamo-DB happened to be especially well-matched to BASE.
So you can see from these examples why CAP would be such convenient way to motivate developers to make the switch: it more or less tells them that if they optimize their code using BASE, it won't scale properly. Moreover, the eBay memos about BASE include step by step instructions to explain precisely how to go about doing it.
Today, fifteen years later, we've discovered more and more ways to implement scalable, high performance cloud services, edge caches that maintain coherence even at global scale, fully scalable transactional key-value storage systems, and the list goes on. Derecho is one example: it helps you build systems that are highly Available, Replicated and Consistent. Call this ARC.
The cool twist is that with ARC, you get lock-free strong consistency, right in the cloud edge! This is because we end up with very fast replication at the edge, and every replica is consistent. You do need to think hard about your update sources and patterns if you hope to avoid using locks, but for most applications that aspect is solvable because in the cloud, there is usually some sense in which any particular data item has just once real update source. So there is an update "pipeline" and once you have consistency in the system, you end up with a wide range of new options for building strongly consistent solutions.
An example I like very much was introduced by Marcos Aguilera in Sinfonia. In that system, you can make strongly consistent cache snapshots, and run your code on the snapshot. At massive scale you have as many of these snapshots as needed, and transactions mostly just run at the edge. But when you want to update the system state, the trick he suggested is this: run your transaction and keep track of what data it read and wrote (version numbers). Now instead of just responding to the client, generate a "minitransaction" that checks that these version numbers are still valid, and then does all the writes. Send this to the owner of the true database: it validates your updates and either commits by applying the updates, or rejects the transaction, which you can then retry.
Since so much of the heavy lifting is down when speculatively executing the transactions at the first step, the database owner has way less compute load imposed on it. Then you can start to think about systems with sharded data that has different owners for each shard, and limits the transactions to run within a single shard at a time (the trick is to design the sharing rule cleverly). Sinfonia and the follow-on systems had lots of these ideas.
ARC creates a world in which solutions like the Sinfonia one are easy to implement -- and the Sinfonia approach is definitely not the only such option.
I'm convinced that as we move the cloud towards the Internet of Things services, developers will need this model, because inconsistency in a system that controls smart cars or runs the power grid can wreak havoc and maybe even would be dangerous.
So now I want to claim that ARC is the best BASE story ever! Why? Well, remember that BASE is about basically available, soft-state systems with eventual consistency. I would argue that in examples like the Sinfonia one, we see all elements of the BASE story!
For example, a Sinfonia system is basically available because with enough consistent snapshots you can always do consistent read-only operations at any scale you like. The state is soft (the snapshots aren't the main copy of the system: they are replicas of it, asynchronously being updated as the main system evolves). And they are eventually consistent because when you do make an update, the mini-transaction validation step lets you push the update results into the main system, in a consistent way. It might take a few tries, but eventually, should succeed.
What about the eventual consistency aspect of BASE in the case of ARC? ARC is about direct management of complex large-scale strongly consistent replicated state. Well, if you use ARC at the edge, back-end servers tend to be doing updates in a slightly time-lagged way: batched updates and asynchronous pipelines improve performance. Thus they are eventually consistent, too: the edge queues an update, then responds to the end-user in a consistent way, and the back end catches up soon afterwards.
And ARC isn't the only such story: my colleague Lorenzo Alvisi has a way to combine BASE with ACID: he calls it SALT, and it works really well.
So, feeling burned by CAP? Why not check out ARC? And perhaps you would like a little SALT with that, for your databases?
Subscribe to:
Posts (Atom)