[MUSIC]. Okay, so Dynamo from Amazon in [SOUND] 2007, which is a paper and again a few years later, it's been released as a cloud service called DynamoDB. Okay, so here we're looking you know, scale, 2000 of nodes. you can look up things by primary index, and basically nothing else, just a key value store, just like nimcash, okay? Right. So ,what are some of the tricks it let to? Well, so some key features are that it has. It's sort of some, so talking to, in terms of DynamoDB, which is the implementation you can now go and use, and pay for. One of the neat things here, is that it offers a service level agreement on performance. And so, at the 99th percentile, you know, they promise to respond within 300 milliseconds for 99.9% of its request. And the reason they do this on the 99th percentile as opposed to some sort of notion of the average the mean or the medians, or the mean and the median, is that would artificially penalize the people who are using it heavily, right. They would get a disproportionate number of failed requests. It would be easy to satisfy the average by only focusing only on the lightweight users for example, okay. So Dynamo, the system is a distributive hash table, that's what DHT stands for. And each key is stored at, or sorry, each value is stored at locations, multiple locations for replication purposes, and its up for replication factor of N. And so it at location K, K plus 1 all the way up to K plus N minus 1. And they achieve eventual consistency through vector clots, which I'll describe in the next couple of slides. And so reconciliation of potential conflicts when things are being read and written. Happens at read time, which is another maybe interesting feature of Dynamo. Okay, so, rights never fail and they site in the paper, that the reason for this is poor customer experience, alright. So, if you're sort of typing into your Google Doc, well, [LAUGH] [INAUDIBLE]. if you're billing application, a web application where you say you update your status on some social networking site, and it comes back with an error message and says, Sorry couldn't commit, you know, somebody else was editing the same. Or you're editing the same status from somewhere else. Their claim is that that's more disruptive than getting the, you know, the wrong read, which seems reasonable to me, okay. All right, and so, conflict resolution for many applications may be the most recent write is the one that wins. Or you can actually have the application controlled. In some cases, you may even sort of go back and ask the user to resolve the conflict manually, okay. Okay, so the goal with Vector Clocks is to detect conflicts in a concurrent read-write scenario. But not to necessarily do anything about them automatically. Okay, so, In this scheme every data item is associated with a list of server timestamp pairs that indicates its version history. And so, in this example, [SOUND] sum value D was read by a client and D1 was written back at the server called SX. And so what SX does is append this fact to this vector clock. So, at timestamp one, server SX, you know, created a change. Then some other client reads D1 and writes back D2. And, you know, you might want to, append to the vector clock, both values. But, you know, this change descends from, D2 descends from D1. It was handled by the same server, and so, you can garbage collect this part of the vector clock right? So, it's the same server with a higher timestamp, means that the old version where the older timestamp is not needed anymore. Okay, and since there were no other conflicts to work on, okay. But now, independently two different clients read D2, and write back different values. One writes back D4, and one writes back D3, and these two requests were handled by different servers, SY and SZ. And so these facts get recorded in the vector clocks since they're different. Now, the contexts here will, this, we call this sort of vector clock contacts. The contacts will reflect this fact when the next read comes in, it will see that, oh wait, there's a conflict, because there's the same timestamp but two different servers. And you can either ask the client what to do or you apply some of your instinct where the later one run, because these might not be timestamps like integers. They could be sort of the actual clock time stamps, in which case you make an arbitrary decision, just pick it and go, okay. So, that's how vector clocks work. So, in the example, just to run through what we just saw. A client writes D1 to server SX and creates this value. Another client writes D, reads D1 and writes back D2, also handled by SX, and D1 was garbage collected. Then separate clients read D2 and write back D3 and D4, and two different servers, SY and SZ. And then, another client reads D3 and D4, and notice, it then finds that there's a con, the system reports that there's a conflict to be handled. Okay, so let's practice with these. Here's, [SOUND] two different vector clocks and you, figure out whether there's a con, whether they represent a conflict or not. So, in this case we have server SX with a timestamp of 3, and on this data server Sx with a timestamp of 3. And then each one is a different server with different timestamps. So, is there a conflict? Well, yeah, there is, because on one version path, SY made a series of changes, and on another version path, SZ made a series of changes. And they didn't talk to each other, because they don't reflect each others' changes. So, yes, there is. And on this one it have the same server at a later timestamp. So, is there a conflict here? Well no, because they weren't handled by different servers. So, really just this one subsumes that one, and we're okay. So, on this one we have server SX with 3, server SX with 3, server SY with 6, server SY with 6 so they agree so far. And this was an extra change of S, of server SZ, with the timestamp of 2. And so, no, there's no conflict here, because they agree wherever there's, on, on, this is just an extra change on top of this one. So this guy wins, okay. In this next one, server SX with timestamp of 3, server SX with a timestamp of 3. Server SY with a time stamp of 10, and this has servers with a timestamp of 6, and then some later change at SZ. So, is there a conflict here? Well there is because this one is later in time on SY, right, 10 then this one is. But then, if that was all there was. If it wasn't for this guy, if it wasn't for this guy here. I'd be okay, we would just pick this one, because it's later. But, because this one's here, we now have some changes at SZ. And some changes SY they were both forward in time from the latest point that we've agreed on. And so we don't know how to resolve that. And so yes there is a conflict. And them similarly here, SX in timestamp with 3 and SX in timestamp with 3. Well here we have SY and SY 10 is the same as the last one. But here instead of 6, it's 20. And so it's later than this 10, and then further we have a change at SZ. And so is there a conflict here? well, no, because this one is strictly later than that one is. On all the servers that they share it has later timestamps. So it's strictly subsumitive. And so no, there is no confluence. Ok, so those are, in the last segment we talked about consistent hashing, in this segment we talked about vector clocks. These are two little gadgets to be familiar with. Because they come up time and again in these no sequel systems, and in other systems in general, okay. All right, so Dynamo also talks about a way to parameterize the level of consistency. And this comes up, occasionally in papers, so I just want to make sure you're exposed to it. So, the idea here is that you have two parameters, R and W. And R is the minimum number of nodes that need to participate in a successful read. Okay and W is the minimum number of nodes that are needed to participate in a successful write. So, this is, sort of, how many replicas you write to, and how many replicas need to respond from a pool to know that you, sort of, have all the information, right? Because if everybody is updating everything all the time, there might be that you have 20 different servers, they all have a different version of the data you're trying to read. And so the question is how many of these you need to sample before you feel like you have the right one? So for a replication factor of N, if R plus W, is greater than N, then you can claim consistency, but where, often you want to set R plus W less than N, in order to achieve lower latency. So, you don't want to have to actually contact, you know too many servers in order to satisfy some read request or write request, okay. So, if you see that notation this is what it means. But I'm not going to describe too much more about it, because I think this sort of, you know, the formula falls over a little bit under a little bit of scrutiny. But that's what they're talking about when you see this, discussion on say, blog posts, okay.