[MUSIC]. So Rick Cattel wrote a nice paper in 2010 about scalable sequel and no sequel data stores where he had a taxonomy of these systems, and placed popular instances of these systems into that taxonomy. And so, this slide corresponds to his grouping, where each color is one of the groups that he defined. So he's already lumped them into key-value stores, document stores and extensible record stores. And so what he meant by that was, a document is you know, an example of this is an XML document or the JSON object that we looked at in that Twitter assignment, where you can now sort of arbitrary, arbitrary nesting. And it's also extensible, you can add new things to it whenever you want. Okay. So, no top down schema being enforced. Okay. An extensible record you can think of much as much like a database record, except that new attributes can be added. Okay. So, there is some sort of notion of a schema that's used for different various purposes, in particular there's sort of groups of attributes that are manipulated together these families. but you get single attributes in an individual row which are not to do with a relational database. And then finally a key value, I'm using the term object here, is a set of key-value pairs. And the difference here is that there is typically not a schema of any kind. So, you don't care what keys they are. They could be any keys at all. And there's no exposed nesting. And what I mean by that is you know, a value can be anything you want, so you might have some kind of complex object in the value, but its not going to be sort of visible to the system. Its just going to be a blob, a bla, a black box object that this doesn't know anything about. Okay. So sort of only one layer of nesting is, is aware of the system, unlike a document store or a document object that might have multiple layers nesting there are exposed to the system. Okay. And so his characterization of no-sequel features, you know, admitting that there's perhaps not a formal definition of no, well, there certainly isn't a formal definition of no-sequel. But the term tends to be applied in contexts of systems that have these features. So the sum-ability to scale simple operations through put to many, many servers. And by simple operation we mean key lookups, or maybe even attribute lookups, or reads and writes of just one or a few records, right? So, these sort of needle in a haystack kind of operations as opposed to these big analytic queries, like we've been talking about with MapRoduce and with databases. Okay. And, sequel criteria is the ability to replicate and partition data over many servers here. So you know, break a single large data set into multiple pieces, [UNKNOWN] automatically, which you automate just yourself. And here you might see the terms sharding and horizontal partitioning. the difference between the two, if there is any, is not particularly important, so you can think of them as synonyms. Whenever you see sharding, think horizontal partitioning of a of, of a database table. You'll see the term horizontal partitioning used more in the database community and sharding used more in the, the SQL community. Okay. And then simple APIs are no query language this corresponds to the first bullet, these are the simple operation, operations. And then critically a weaker concurrency model then, what I'm saying is ACID transactions now, and I might, we'll go, we'll talk a bit about ACID in the next slide, but I'm not going to go into a lot of detail here. There's, you know, 40 years of research on the, this topic, too much to cover in this course, especially when we're mostly focused on you know, reading data and analyzing data as opposed to a concurrency draw. But we will spend some time on some techniques of converge control in a minute okay. and then some of efficient use of distributing, he tells us efficient use of a distributed in X is and RAM for data storage. So, this is kind of minimizing latency, is their emphasis, as opposed to just throughput. All right. And then typically, they have this ability to add new attributes to data records in various ways, as we talked about on the previous slide, right. So the, the, the lack of a sche, no schema, is what you can I can think of here. Alright. So, this is no schema, no transactions. And so, we'll go into more detail there. No query language, we can go no SQL. Right. And high scale. Okay. So, ACID, and he, he talks about this term BASE that never quite caught on. I, I wouldn't typically use that, this term, and I'm not sure I recommend you do either. Certainly, ACID is much more permanent in the vernacular than, than this BASE is. it's, it's, okay. So ACID is an acronym standing for these four concepts: atomicity, consistency, isolation and durability. And just briefly, this is you know, the context here is when we're modifying records in, let's say a database. And we can be modifying lots of different records across different tables, anything we want. And the point is they're all lumped into one transaction. And so, each one of these refers to you know, that's the context for each one of these concepts. So, atomicity means that the entire transaction either needs to succeed or needs to fail. Right, you gotta learn to have partial transaction succeed. Consistency is the slipperiest one in my mind and this quote down here maybe captures that. And so, there's sort of any data written into the database must be valid according to all defined rules. And the question is well, where do these defined rules come from. Sometimes they're actually integrity constraints in the database. other times they're just sort of business logic rules, perhaps enforced by the application or just assumed by the application. So, it's a little bit difficult to say, prove a system is, achieves application level consistency but that's the goal. The point is, you, if, if there's only certain allowed states the database to have, you can't, you shouldn't have a system that allows transactions to put you into an invalid state, okay? Usually, you'd only go from working state to working state. Isolation means that while the transaction is occurring, other readers and writers can't sniff partially completed values. Okay. No, partially, you, you can't sniff values of data items before the transaction is complete. They only get final stage. And this one is the one most often relaxed in various ways. In part because it's very expensive and also because it's not usually all that critical. And then durability just means that if you report back that the transaction succeeded, it needs to have actually succeeded. Meaning that it needs to be written out to some kind of non-volatile storage, so that if the power goes out and the machine crashes, you don't say, hey, whoops that transaction that I accepted yesterday or committed yesterday. Well, you need to do that again because it didn't take. Alright, so that's not allowed. So fine. So these, these all make some sense with you know, a little bit of, a little bit of notion of consistency as I mentioned. And the pun here is that they're trying to sort of force an acronym on BASE, and this isn't Rick Cottell, this came out, else from elsewhere. but the idea is well, it's basically available. There's some notion of soft state and it eventually consistent. We talked about eventual consistency eh, eh, at least an overview in the previous segment. Fine, that's all I'm going to say about, that. So, something else I like about this paper is, he sort of says look, you know, the, the major impact systems here are these three. This Memcached or Memcache D, Dynamo from Amazon and BigTable from Google. And the reason he says these are the major impacts is, you can kind of trace the lineage and show that other systems are basically taking ideas from one of these three early systems. So memcache is very, very simple. And we'll talk a little bit about one particular technique that it made popular, in a minute. but it's essentially just, hey, look, let's just load everything into memory, scale it out across many, many machine. Right, and they'll be able to serve read requests without having to go sort of, query the data base. We'll just be able to do it directly from memory. And what's also made this very, very popular is you can kind of install it on top of your either scale, scale-out or non scale-out database. And it just sort of just works, right. It just makes things faster for, for read heavy workloads. And that was kind of a nice thing. So that's an older system, sort of around the scale of 2003, but it's still very widely used. And very, very popular, and there has been all sorts of extensions to it. And so, we'll talk about the probably the most basic version. Amazon's dynamo paper which has been somewhat more recently released as a cloud service called dynamo DB. what they did was, they didn't invent the concept of eventual consistency, but they did sort of show that if you relax the consistency notion, that will allow you to scale way, way out. Okay. And so data fetched can not be allowed, can not guaranteed to be up to date, but updates are guaranteed to be eventually propagated, everywhere they need to be. And we gave an example of why this was a good idea in the last segment. And then Google's BigTable that we'll spend some time on you know, demonstrated that record oriented storages could scale to 1000's and 1000's of machines, and that was something that data bases had not shown. Okay. So, let's talk about each one of these systems in turn. So memcached, as he says, main-memory caching service, no persistence, the basic version is no replication. Meaning there's not two copies, there's only one copy of every cached value. So, if something goes down, if that goes down then it's gone. That's okay, because it's sort of a cache, it's not assumed to be the golden copy of anything. That being said, there's been many extensions that provide various, these various features including membrain and membase. So it's a very mature system and still in wide use. And an important concept that they adopted, in this context was consistent hashing, so I want to explain a little bit what consistent hashing is. so that's one takeaway from, from this lecture, okay. So first for those of you without actually having too much of a background in programming. What is hashing, so what is regular hashing? Well, the problem we're looking at there in this co, hashing's a very, very general kind of, it's very fundamental to all programming. But in the context of what we're doing here, we're trying to assign data keys to a bunch of different servers. Okay. And the simplest way you might do this is sort of a round robin thing. Right? The first key goes to the first server, the second key goes to the second server, and you keep going untill you run out of servers. And you start back over by the first one. And that's implemented by this module. Okay. So, each of these data keys is placed somewhere on this, on one of these servers at various points. Fine, that's how hashing works. What's, what's wrong with that? Well, what happens if I want to add more servers to the mix, right. I want to scale out to, I want to double the number of servers. Well, every existing data key now needs to be reeval, it's place, it's location needs to be reevaluated by computing k mod 2 N instead of k mod N. Which means every single data item is going to be remapped at once. So every time you want to add a server, you end up having to move all the data that's already in the system and you're dead in the water. Okay. So what you want is some notion of consistent hashing, where consistent means when I play something somewhere and I add more servers, it's typically going to stay right where right where it is. And so there's a pretty good trick that's pretty simple to understand for doing this. Okay. So here's how it works. First key idea is you're going to map the server IDs into the same space as the key values themselves. Okay. So we apply a fa, a function that I'm going to leave sort of unspecified and map server 1 to some point on this circle and server 2 some point on this circle, and server 3 some place on this on this circle. And now, what that does is, divide this space up into three sections. Okay. And now, each key that comes around, I also map it into this circle. So, this gets key one and this gets key two, key three, key four, key five, key six, key seven and so on. Okay. And now this entire region one, server 1, server 2, server 3, this entire region is responsible for all of these data keys. Sorry, this server is responsible for all these data keys. And then this server is responsible for all the data keys in this region. And this server is responsible for all the data keys in this region. Okay. And so what's nice about this is, now when I add a new server, server ID equals 4. Well, let's say it comes around and it gets stuck right here. Well, that's a bad spot for, for my example actually. Let's say it comes around right here. Well, you just apply the same rule. It should be a po, it should be responsible for every key in this region, which means that these two guys need to be moved from server 3 to server 4. Right? But you only have to move that one section of data. And so, it splits at most sort of k over N data items. Alright. So this is a nice trick and there's all kinds of extension for supporting replicas we need to put data in obviously more than one place, well, just sort of hashed in two different places. So if you want to hash the same data under two different places compute, h of, lets say d as the data key, and then also put it to, you know, put it all three places, and you're done. Okay. So how do we serve request in this set up? Well, imagine the key space is divided across various servers in the same way we describe, and the request comes into the leader that may be elected among the servers or may be just assigned top down by the, by the system. Or could even be assigned randomly. And the naive way of, of doing this is well, this, this, this server would check to see whether it has the key being requested, and if it doesn't it would just forward the request on to the next guy. Okay. But this is no good because there could be many, many servers and this would encourage server latency every time you would want to do a read. So, a better way of doing this is for each server to memorize the locations of other servers in the ring. And which servers it memorizes is like this. So it knows where itself is, it knows, A plus 2, A plus 4, A plus 8, A plus 16 and so on. And what it does is, how it knows the key range being managed by each one of these servers. Okay. And so what it can do, is forward the request to the server that is closest to the key range it is looking for, okay. And so this takes a logarithmic number of hops away, you can imagine there are lots of servers here. And keeping all this information straight when new servers come in is still, each server only has to keep track of a algorithm of servers as well. So everything sort of ends up being algorithm to maintain this. Okay.