[MUSIC]. Okay, so we start with the same schematic that we were looking at when we were talking about Map Reduce and Scalability. Where we take a big data set and break it into chunks, and send those chunks to different machines. Okay, and here we are replicating this chunk to three different machines, which is the same thing we did for the Hadoop file system, for fault tolerance purposes. Where you know if this machine dies we still have two copies of the data to draw from, and we do this with every chunk, alright. But the two questions, the two requirements we need to speak to hear is we need to ensure high availability. So, that when something goes wrong the data still available, and we also want to support updates in this context which is different than what we were talking about before. So instead of just read performance or fault tolerance in the content of reads, we also want to make changes to this data now. And have them propagate to both other replicas, and in some cases to other consumers of that change, right. There might be other blocks of data that that are referred to as the same information, I'll give you an example on the next slide, okay. So, imagine a social networking application where people are updating their status and their friends get to, you friends get to see your status updates, okay. And so the right operation here is Sue updates her own status and the question we ask is, of her friends, what happens? Who sees the new one, who sees the old one? How do these, how does this status change propagate? And the answer to this question from a database perspective was, well look you know everyone must see the new change or no one does. Right, either the transaction commits and all copies of the data everywhere are synchronized simultaneously. And further anybody attempting to read the value in new media state is able to read only the old value, or is, or has to wait until the transaction commits, right. Which could be an arbitrarily, a pretty long time, deadlocks can happen which is why I said arbitrarily, fine. So, that's the answer given by databases, everything synchronous, everything must be updated, It's either all or nothing, okay. And the noSQL system just sort of make this observation, I said well look for really large applications, we simply can't afford to wait arbitrarily long for this to happen, right? I mean, you need status updates to be able to commit and respond, so the user can go on and do other things, right. They can't, sort of, look at a, at an hourglass while the synchronization is still ocurring. You know, and then further the observation is, well maybe it doesn't matter anyway. I mean is it really that important that, you know, here, if, if Sue's friend Joe, sees the new status while Kai still sees the old status, maybe who cares, right? As long as Kai eventually sees the new status, maybe that's good enough, okay. And so these observations suggested moving in a different area of the design space in sort of high scalability, high availability and, consistency, application consistency. and that motivated and, and those, that space of systems started to be associated with kind of anti-database, right, took a very different approach than databases data, and so in turn NoSQL came into play. It's actually unfortunate that the you know, the name that stuck was NoSQL, because it doesn't have a whole lot to do with SQL. It has more to do with the transaction processing side of the databases, which is not all that relevant for SQL, right? I mean the model of transactions, the sequences of reads and writes nothing to do with the query language. But, hey, that's what stuck. Now, I, I don't mean to that the term NoSQL only suggests these transaction models. It also sort of suggests a weaker data model and so on. And we'll talk a little bit about that. But I want this point to come across, because this is one of the key ideas. Okay, how did databases solve this problem, or why did they take so long? Well there's a, protocol called Two-Phase Commit that's fair, fairly standard in these situations for synchronous processing. And so, the motivation for why you want two-phase commit goes like this. If you want to have a bunch of replicas or other kinds of subordinates, anybody that wants to see the change, you know your, the, the, you make your status updated and your friends need to see it. The server's holding those different friends need to be told of the change, and so if you just go ahead and tell them, say look I made this change go ahead and update your internal state to reflect Sue's new status. Then you can have you know, some of them report back success, but one of them could fail. But now you're in trouble right, because this one has the old value, because it failed for some reason. Either you didn't hear back from the server at all or it said look something's wrong with my disk. I can't do it, so it responds with a failure, regardless. And these two have already successfully applied the, I'm going to put a check mark, have already successfully applied the transaction. And so now you're in an inconsistent state. Subordinate 3 has the old value, and these guys have the new values, okay? So, how you solve this problem, is two-phase commit. So the first phase here is, the coordinator sends a prepare to commit message, and the subordinates make sure, they take action to make sure that they can commit that transaction when asked no matter what. And so typically this means writing to a log the, the, the information related to the transaction. So, that even if the power goes out, when they wake backup they can pull it from the log, okay. And the subordinates reply with a yes, I'm ready to commit. And then in phase 2, if all subordinates say they're ready, then you'll go ahead and send the commit message. And if anyone failed, if, if rather instead anyone failed then you send back an abort message and need to be just clean up, okay. So this is fine. and here's the schematic of it and step one, they say prepare, these guys all write ahead to the log, and say I'm about to write, I'm going to commit this transaction. They response with yes, I'm ready to do so. The coordinator comes back with commit, and then finally all the work is done. And I'm not going to show the schematic for what happens in a failure, then essentially the coordinator needs to watch out for it and send back an abort if something had gone wrong, okay. Okay, so there's a couple of problems with this. One is there's some dependencies on the coordinator here, that if the coordinator fails at the wrong time, things can go kind of screwy. And a fully distributed protocol for ensuring mutual commitment of transactions, or other kinds of operations can be achieved. And one, one of the most successful and popular methods of doing this is an algorithm called Paxos, that we're not going to talk about in detail. But you're going to see that term, if you look in some of the reading for the NoSQL systems, okay. So think two-phase commit on a local cluster for a database, think Paxos for a distributed sort of peer-to-peer kind of protocol. And just briefly, what Paxos is essentially doing is it's a voting scheme. So, people sort of vote on, you know, the individual servers, well to, self determine whether or not they're supposed to commit the transaction or not. And a, the details can get a little bit subtle, but overall it's pretty simple given the nature, given the difficulty of the task involved. Okay, so fine, that's one problem. The other problem is just, with the Paxos sort of shared, is that, this can take awhile, right, if subordinates don't respond promptly, he might be waiting around. If things fail multiple times, and you just sort of abort and retry transactions if the application layer things can go slow. when there's, it doesn't necessarily scale when there's thousands or millions of the subordinates are needed to do this, you're kind of dead in the water. So, other protocols that I'm not going to talk about in too much detail include multi-version, but you will see in some of the papers mentioned, multi-version concurrency control. Where each write creates a new version of the data item, and the legality of the read is determined by checking the timestamp of the read transaction versus the current timestamp of the version that you're trying to read, okay? And if it's been updated since the time you're supposed to be reading it, then you know, prior to MVCC, all you could do is abort the ter, abort the read. And say, look, you, you're looking at dirty data, you're done. And but with multi-version concurrency control, you can actually keep multiple versions around, and redirect the read to the potentially to the prior version that is correct. Okay, and thereby avoid avoid aborting a certain transactions, fine. So, that mechanism still has the dependency on a coordinator role to administer the time stamp. A fully distributive scheme, where the decision to go forward of the transaction or to avoid a transaction is made. The revoting scheme among peers is Paxos. And Paxos is very successful and very widely applied. And you'll see it mentioned in some of the noSQL papers if you take the time to read them, and they're on the reading list. And so this relieves the dependency on having a central coordinator. but is still synchronous, and still has the potential for deadlock. And, can take some amount of time to reach consensus, to be able to know what's going on and what kinds of failures are happening. And so it's difficult to guarantee very high performance of very low latency response times, alright. So, then the term eventual consistency was originally defined, not so much in the context of its utility, and allowing systems to scale to very large levels. But just in this argument that the right, the only players in distributed systems that could make the appropriate decision about how to handle conflict where applications themselves. So, it's a version of this in to in argument that you may or may not come across in the context of networking. And so this was a paper in 1995 by Doug Terry, where this term was coined. And so he says, you know, we believe that applications must be aware that they may have read weekly consistent data, and that the right operations may conflict with those of other users and applications. And that applications must be involved in the detection and resolution conflicts, since these naturally depend on the semantics of the application. And so we'll make the argument in a few couple of segments, but I'm not, I'm not sure I totally agree with these assert, assertions. That it actually is better for the system to take care of this when it can. But what I wanted to do is let you know that this is where the term comes from, as oppose to the NoSQL system in the last 10 years or so, which it really would, would increase in popularity. Okay, so what does it mean? Well, what it means is, that in the absence of updates, all replicas will eventually converge towards identical copies, right? So, as long as things don't continuously change, as changes settle down, we'll all eventually see the same value, right? All your friends will see your status, alright? They won't be permanently stuck looking at an old one. But, you know, what the application sees in the meantime, what's one of your friends, which, which status one of your friends might be looking at, is really sensitive to the internal details of whatever application you're building and is difficult, and therefore is difficult to predict. Okay, and so for this reason it's, it's a little bit difficult to reason very precisely or formally about what eventual consistency means. Because it is so dependent on particular limitation details that are themselves difficult to formalize, okay. And in general, contrast as we've only been talking about relational databases and things like Paxos, where they guarantee strong consistency, but there maybe deadlocks. And so it's, you can prove that no system can be free of deadlocks and guarantee consistency. And so, relational databases like Paxos give up on this liveness property, meaning that they'll, they might allow deadlocks in favor of strong consistency. Now, they've, you can show that the cases where deadlocks can occur, can be made sort of rare, through different design decisions, but they can still happen. Okay, fine, so visually consistent models say, we can't afford the cost of waiting for these protocols to run and more over, they might not be necessary in certain application context. Alright, so where we are now is we're looking at this column. And I've already sort of marked this up a little bit, but what these words now mean, and we'll talk about these a little more when we talk about a few of these systems, is the scope of where strongly consistent transactions are supported. And so the scope here of a single record, means that I can update a multiple fields in one record, and either all the changes will occur, or none of them will occur, okay? by the way I filtered this list only include noSQL systems, so relational databases support this across arbitrary records, right. You can have, you could update a record over here and update a record over there and call that one transaction, and the system will only see both of those changes or neither of those changes. And that's what's not supported with these NoSQL systems. So, within interview record is supported, within some of these systems nothing is supported. you can't, there's no guarantees at all really. And what this EC means is eventually consistent, so its not really strongly consistent its not a transaction, but they do have eventually consisting guarantees at the record level. And that's what all these systems sort of guarantee, and then this system Megastore that's based on big table from Google. It's also a Google system. defines the notion of entity groups, and this is a set of related records for which transactions are strongly consistent for that group, okay. So, this is a little bit better than just one individual record as the one, only guarantee we can give, and it's a little bit less than any arbitrary record in the database. It's predefined into the groups that allow transactions. Okay, so this is sort of a compromise, fine. And then this most recent system form Google Spanner offers true strong consistency across all the records. And we'll talk about why they made that choice, in a little bit. Okay. So, another concept I want you to be familiar with is this so-called CAP theorem, from Eric Brewer in 2000, followed up later by Lynch in 2002. Where they define these three notions, consistency, availability and partitioning. And the way this is often described is you have to choose two of these. You know you can't get all three, you have to choose two or sacrifice per, performance. But I don't really like thinking of it that way. And Eric Brewer has also sort of described that maybe that's not the right way to think about it. And the reason is because it's not clear what it means to choose consistency and availability at the extent of partitioning. Okay, so what is partitioning? Partitioning means well, if you got a big distributive system with a hundreds of nodes involved, hundreds of servers all communicating with each other, and some segment of them lose communication with the other servers. Can those two segments still make forward progress in the application independently and sync up later? Or does everything have to stop and wait, or certain nodes have to stop and wait, in order to re, reestablish communication? So, for example, if you have a master node that controls everything, and you have some worker nodes that lose contact with the master. There are many designs at which you can't make any forward progress, until you reestablish connections with the master node. Right, you can, you can do no useful work, because you're waiting on communication, you're waiting on that last message from the, from the master to tell you what to do next. so in those cases you've given up availability, right, you go down. Those nodes are no longer accessible or, or, or doing useful work in the context of a network partitioning, okay. On the other hand, if you do say, well sure we're going to continue to do useful work even independently, then it's not difficult to show that you can arrive at an inconsistent state, right? Updates are coming in to this partition, and updates are coming to this partition of the network. And sometime down the road communications are reestablished. And you find out, oops, you know, your replica has one value, my replica has another. Which one's right? Well, we're going to have to sort it out, but meanwhile we've already sort of exposed these values to the application. So in some sense we're, we're demonstrably inconsistent, okay. It's the point is you can't get all three of these. All right. So, you either sacrifice availability or you sacrifice consistency by allowing things to continue working. So, conventional databases essentially assume that there is no partitioning. And again, this is a function of them only operating on tens of nodes at a time, all right? They didn't go to this thousand node scale, or planet wide distributed systems. and so you can kind of assume that there wasn't, there wasn't really a need to worry about. Well, what if queries are coming in to half my nodes and they can't talk to the other half of my nodes, and so on. They're all sitting there in a cluster that's in your data center. Not even in your data center, in your server room to some extent. And so that wasn't, really wasn't an issue that they were thinking about too much, okay. And the NoSQL systems do need to worry about this. They are very large, they are very distributed. There are different kinds of Byzantine failures happening all the time, because of, because of the sheer scale. And therefore, they choose to sacrifice consistency instead of availability, okay. And so, graphically you can look at this in this sort of triangle forming, you can put different systems on sort of an edge here where relational databases assume consistency and availability. but, but assume partitioning can never happen. While other systems need to tolerate partitioning, but give up on availability. and they are ensuring that certain kinds of transactions are going to be consistent, okay. And then other systems say, well, we're going to give up on consistency. but you can always use [UNKNOWN] work, okay, so fine. And really this is, the important thing here is the scope of the variation I put in my table's sort of critical here too. It's not sort of, nothing except for spanner over here, even tries to provide global transactions like relational databases do, okay. Alright. So fine, so I'll pick up here in the next segment.