[MUSIC]. Okay, so let's talk a little bit about other large scale data processing systems besides MapReduce. And as a step to that, let's think about the design space of possibilities here. So, this is a breakdown that was proposed by Michael Isard who's developed a system called Dryad at Microsoft, which is a very nice system with the same sort of motivations as, as MapReduce, okay. And so, he divided this space into these three axes where he worried about, sort of, low latency, very interactive sort of speeds, quick turnaround time. Versus things that maximize kind of throughput, you know massive batch jobs operating on you know, thousands and thousands of computers at once. Versus another access here is sort of whether it's in a private data center or whether it's scaled out widely over the internet. and then maybe the third access is data parallel versus the shared memory, okay that we talked about. So the, the areas that we're mostly concerned with are going to be here which is what we're talking about currently. And then in a couple of segments we're going to talk about these Low-latency, smaller operations which you can think of as the no, the NoSQL systems. Okay. And then maybe where, where Michael placed older ar, relational databases, although he didn't label it the same way I have. Is down here in this quadrant where they're mostly in shared memories with the space and low latency. I, I say older databases to try to point out that not all relational databases operate in the space. Many are, most are in fact are data para, most parallel databases are. Are data-parallel, as well. And so this is, you know, MySQL and PostgreSQL, if you're familiar with those, are probably in this space. Shared memory, shared disk. Alright. And here, HPC means high performance computing, where you know, it's private data center and it's a big mass of shared memory. but there's, but, but it's a batch job submission system, right? You say, you're on your big compute, compute-intensive job and submit it to the machine. And it processes on it and eventually turns the results. And then, I'm not going to talk too much about this, but this notion of grid computing was really sort of focused on connecting up clusters of computers from different universities and letting them all sort of talk to each other. And so that's why he's pushed this up on this other axis. Increasingly, you're seeing these systems also being pushed up in the in this access. Internet scale, planetary, you know, distributed hash tables. With different kinds of layers for guaranteeing certain kinds of semantic. So the spanner system from Google fairly recently is, is a nice example of this. Okay. So, sort of wrap up what we talked about last time, large scale data processing you know, many tasks need to process big data and produce big data. And so you want to use hundreds of thousands of. The CPUs are, and hundreds of thousand or tens of thousands of computers to solve these problems, but this needs to be much easier. And so there are such things as parallel database. We talked about Databases and we sort of extolled the virtues of programming in that model, and they exist but they're often. Expensive. Well, they, they almost exclusively are expensive. And they're difficult to set up. And it's actually not totally clear that many of the parallel databases scale to really hundreds or thousands of, of machines. Ok. And so MapReduce. Came around at a time at a time as a, as a bit of response to this scenario. So it's more a light weight framework featuring, you know, automatic parallelization and distribution that we've been talking about, featuring I, a fault-tolerance. And I mentioned a couple of other things here that are probably less important but the I/O scheduling, the status and monitoring. So really sort of strip everything down to just parallel processing. Not all the features of parallel database is offered, just parallel processing, with the added benefit of fault-tolerance. Okay, and this really seemed to scratch a niche with people. I'll, I've, I'll argue here and I'll probably mention this again at some point that it's not totally clear to me that MapReduce would have been quite so popular had there been available a parallel open source. relational database product, but all the open source databases were all not, not parallel. In fact they were not even single threaded for query processing, you know only a single thread was working on an individual query at a time. But that's, that's, that's speculation. Well, actually, I have a little a bit of evidence for that that I'll lay out in a bit. Okay. So now I want to talk about, maybe, I guess I'm building up to that argument. I want to talk about parallel databases and how they work. And hopefully, show that there's some similarities, show where there are similarities and where there. There are differences. Okay. So we'll call that a key idea of. Relational databases was this notion of a relational algebra, where you could write sort of plans like this. And that this top-level language called SQL was the most common way of Producing a Relational Algebra plan, right? You wrote the query and sequel and it was automatically turned into a Relational Algebra plan by the system. OK. So this is kind of thrown out the window with MapReduce, arguably in favor of sort of flexibility in providing the program with more, more control. But let's go back to this model for a bit. Fine. So. Now we want to, we want to evaluate these queries, we want to do it in parallel now. And so there's two different terms that I want you to be familiar with, one is distributed query and one is parallel query or distributed query processing and parallel processing. And so they're both ways of sort of taking advantage of more computing resources for the same query but they're, behave a little differently. So the distributed query, what you're doing is taking a, a single large table and distributing it across a cluster, just like we talked about. And then you're breaking your query into individual pieces to operate on each of those partitions of the file. Okay. So this sounds like, well isn't that basically just the same thing as map reduce? It is, except for the fact that, all the results of those individual pieces are all sent back to the head node to a single server to sort of finish processing. So for example if you're doing a, a, a count, right, you want to count all the records that match some criteria. Well if you have a very large file that's split across several machines, these distributed query systems, Microsoft Sequel Server in particular is, is an example of this, is smart enough to break the query into a bunch a little pieces and run each of these of pieces in parallel. But as they start to string tuples out to be counted they'll send them all back to the head node. Actually, I guess that may not actually be true so maybe it's one of the thongs that you can just count things in parallel and add them. The map, but it's not hard to construct a query where, where you have this bottleneck of, of sending everything back to a single server. So it's essentially, you can think of it as having the map phase, but not really the reduce phase. Now, parallel query, every individual operator in the relation to algebra is implemented in parallel. So when you're doing joins, you're doing joins across. A bunch of nodes when you're doing groupings. You're doing grouping across a bunch of nodes. And we're seen how to implement relational join in map reduce and it's not too far off from how it's actually implemented inside databases. Okay, so if we how to implmenet join, that's usually the harder one. One, trust me that you can implement the other ones that way. Well, now we have a way to do parallel query processing with the relational algebra. You know, so why not do that? Well, the answer is that people do do that, and we'll come back to that in one second. Okay, so for a distributed query, I guess I was waving my hands a second ago trying to explain this, when, when it was all on the next slide. You can imagine constructing a view, and we talked about views, if you don't recall what that is, it's you know a, a named query that can be then accessed as a single table. So we say that the sales table is really the union of a bunch of smaller sales tables, one for each month, and in particular you could put each one of these sales table on a different disc or even a different server all together. Right? And then the, the user who is querying the sales table doesn't have to care about the fact this is actually distributed distributed table. They don't have to go gather up all the results from January, and then gather up all the results from February, and the gather up all the results from March and put them all together. That's done automatically by the system. However right, and so this, this is, this is the create table statement that we didn't talk about for constructing the sales table for sa, individual March, you know, okay. But again, however, when you process this stuff in parallel that works great but when you get the results you need to send them all back to the single node to for, to finish processing. And that's the limitation Distributed query. So it's great that you get some parallelism, but it can't do everything in parallel. And you, you can, you can see this when you run performance experiments. But a true parallel query example, would be, for example from a system called Teradata, which is a database company that many folks haven't heard of because they're selling very, very high-end databases to very, very high-end customers. And so they don't sort of need to have much word of mouth in the popular news media. But what's happening here is that as every. Individual rows inserted into the parallel database. It will be assigned to some particular server using a hash function. Okay. Fine. So, everything is automatically partitioned more or less randomly across the cluster and then whenever you're running queries on this, all of this, all the machines will access their, their data in parallel. Okay, and you can see this has a little bit of the flavor of how we did the relational join in in MapReduce. And that should be coming more clear in a second. Okay, so remember this is our query orders and line items, and this is the plan that we're going to do. We're going to select some orders and then join the orders with the items. Alright. So how this starts is these parallel processing units, units called amps in teradata terms, will each contain a piece of the data. A chunk of the data, and the chunk was to find sort of randomly by hashing. Okay. And so they all in parallel should begin to scan their, their individual chunk. And then they'll all in parallel apply the filtering conditions to throw out certain records that they don't want. And then they'll all in parallel hash on the so this isn't right. This should hash on the order. ID. Not the item ID. [BLANK_AUDIO]. So this is the join key, right. We're going to join on order. Here. and I probably had this wrong back here too. Yeah, this is wrong as well, this should be join on order ID. Doesn't quite make sense to call this item. Same thing here. Okay, so then they're all in parallel hash on the join, join attri/g. Tribute the order ID. And that will shuffle it, for lack of a better term to another set of amps. Perhaps the same set of amps, but typically another set of amps that will do the next step. And actually compete the join. And so this should look like a MapReduce job, right. You've got a map function that's scanning and selecting and, and then, and actually then hashing. And then we gotta re, reduce function coming up to actually produce the join. Okay. And for the other relation the same thing happens, you scan the items. And then hash on. Again, this is order, order, order. Right, so you're scanning the order, scanning and selecting on the orders and just scanning on the items. And then both are hashed on the appropriate joint attribute. And, lo and behold, all the items and all the orders that correspond to the same Order ID to the same join attribute, same join key end up on the same machine and you can actually process the join, just like the MapReduce example we saw. Alright. So then at the end of these two steps that I've shown you, amp four will have all the orders and all the line items where hash of order, goodness equals 1. And this AMP5 will have have all the orders and items where hash of order equals 2. And this one will have all the orders and line items where hash of order equals 3 and now it has enough enough to individually and in parallel finish the join and actually produce the result. And all these other the orders as well. Alright. So, fine, so the point is, is that the same machinery already exists in these parallel databases. And in Map Reduce you know, if you're interested in, in doing a join, you're sort of implementing this yourself. And so this observation was not lost on people that, you know, hey, it might be nice if there was sort of a standard way of doing join in MapReduce. And we didn't have to sort of rewrite it out ourselves every time and in fact, you know, there's, there's libraries on top of MapReduce that do this. And so there's a library called Pig from Yahoo, that encourages you to check out that Is recognizably relational algebra. Alright, there are operators called join, there are operators called group by. It does have a bit of a funny data model, where you're allowed to have kind of complicated nesting. As opposed to just straight tuples and straight relations. But the relation algebra is there, and in fact, this is sort of one of the points I want to make, is that. It, it, you know, it's important to sort of be able to modularize the concepts that come out of various communities and especially databases. They, it tends to be true that you know, it's kind of all or nothing. If you're interested in using databases, well then you have to take everything, you have, you have to take the whole package. You know, it's all or nothing. But increasingly what you're finding is that these concepts are leaking out into other systems. Which is why I'm really, emphasizing this relational algebra piece a lot. Is that you can use these concepts independently of buying in to a strict relational model. Okay. And certainly not a strict, a strict adherence to, to a particular information of it. Now, another system called HIVE is literally SQL on top of. [INAUDIBLE]. So it's goes one step even higher. Instead of just stopping the Relational Algebra level, it actually provides a sequel interface. Impala's a more recent system from Cloudera I should mention here. Cloudera by the way, is a company that has. align themselves pretty closely with the, the Hadoop stack, and so they have their own fork of the Hadoop system and a bunch of great tools for working with that ecosystem, and Impala is a new system that they produced that is, provides SQL over HDFS, and actually uses a lot of the code from the HIVE system. Okay. Cascading is another system that's maybe a little bit less common, but it's also very recognizably to be relational algebra. The Dryad system I mentioned has nothing to do with MapReduce directly except for a similar motivation, but it very obviously has relational algebra there. the Clustera system I, I mentioned, it's more a research project and is not clear to me that the code is available. But it's also very clearly relational algebra. So, you know, when you put your relational algebra goggles on, you start to see the world in this way, and it starts to come up everywhere, okay? So it's good to go back and understand those operations. All right.