[MUSIC]. Okay, so just to wrap up this discussion of MapReduce versus databases, I want to go over some results from a paper in 2009 that's on the reading list, Where they directly compared Hadoop and a couple of different databases. And see if we can, maybe explain what some of these results tell us. Okay. So this was Emmy Pablo and some other folks at M, MIT and Brown, who did an experiment with this kind of a set up. So, the comparison was between three systems, Hadoop, Vertica, a which was a column-oriented database and DBMS-X. Which shall remain unnamed, although you might be able to figure it out and so, we haven't learned what a column oriented database is and what a row oriented database is. But we may have a guest lecture later that will describe that in more detail. But for right now, for our purposes, just think of these as two different kinds of relational database or two different relational database, with different techniques under the hoods, in the, under the hood. Okay. And so, there's two different facets to the analysis. One was sort of qualitative, about their discussion around the programming model and the ease of set up, and so on. And the other was quantitative, which was, performance experiments for particular types of queries, okay. So, the first task they considered was, what they call a Grep task and so this is a task to find a 3-byte pattern in a 100-byte record. And the data set was a very, very large set of 100-byte records. Okay. So, this was done in, with, this task was performed in the original MapReduce paper in 2004, which makes it a good candidate for a benchmark. And so, the data, the data set here is 10 billion records with, you know, totaling 1 terabyte spread across either 25, 50, or a 100 nodes. Okay, so you are just trying to find this record. So, this is much like this, you know genetic D sequence DNA search task that we described as a motivating example for sort of describing scalability. Okay. Fine. So what were the results? Just to load this data in, this is what the story sort of looked like. Hadoop and the system called Vertica that they are really the, the theme here is that they were the designers of the Vertica system. And so much that these results are going to show Vertica doing quite well, for, for a variety of reasons. So we're not going to talk about too much about those particular reasons, we're mostly going to be thinking about DBMS-X, which is a conventional relational database, and Hadoop. Okay. So here, loading is fast on Hadoop, while loading is slow on the db, on the relational database, and again it was sort of fast on Vertica as well. So, why is it faster on Hadoop? Well, there's not much to the loading right? You have to put it into this HTFS system, so, it needs to be partitioned, but that's about it. When you put things into a database, it's actually re-casting the data from it's raw form, into internal structures in the database and that takes time. Okay, and the process could be even worse. Because if you're building indexes over the data, you actually, you know every time you insert data into the index, you need to sort of maintain that data structure. Okay. And so load times are known to be bad. So, the take away here is remember that load times are typically bad in relational databases relative to Hadoop, because it has to do more work. Now, what is in the database you actually get some benefit from that and will see in a second these results. But actually we know, we know we can conform to a schema for example. Hadoop is just a pile of bits. We don't know anything at all, actually we run at out produced task on it, okay. And so, how much faster will [UNKNOWN] their experiments for the on 25 machines, you know, we're up here at 25,000. These are all seconds by the way. you know, 7500 seconds versus 25,000 and a little bit less as we go to more servers. Okay. Now, actually running the Grep task to find things, this is what we see. Again maybe ignoring Vertica for now, because I haven't explained to why, you know what the difference about Vertica that allows it to, to be so fast. But just think about a database from what we do understand. And Hadoop is, s, slower here and the primary reason is that it doesn't have access to a index to search. Okay. So again, no indexes available, Hadoop has to do. That's wrong. Okay, so Hadoop is slower than the database, even though both are doing a full scan of the data. The grep task here is not something amenable to any sort of indexing. You actually haven't touched any record, so there's no fundamental reason why the database should be slower or faster. But, partially because it gets a win out of the structured internal representation of the data and doesn't have to re, re-parse the raw data from disc like Hadoop does. And so, I said that there is no fundamental reason, there is a fundamental reason. Because it's already in,in a packed internal binary representation, which we paid for in the loading phase, but now we get the benefit from. Here in the query phase, even before we even talk about indexes. Okay. Now, a selection task we're not having to scan every record necessarily. You know, that is amenable to indexing as we discussed in the scalability segment. Well, the story is even you know, more extreme. [UNKNOWN] right? The Hadoop results are just way, way, way higher than both the database and in particular, the Vertica results. And so here, the reason is because you can build an index on the page rank attribute and zoom in directly to the records that your, you're interested in. Okay. Fine. So, those are sort of search and retrieval tasks not, arguably not exactly what Hadoop was designed for. Hadoop was designed more for analytically tasks. So, now consider these, so here the data set is 600,000 HTML documents. Which works out to be 60GB of data per node. Along with, another data set is 105, 155 million user visit records, and 18 million rankings records. So, this is kind of a web data processing task. Alright. And so, a simple aggregate task here is to add up all the adRevenue corresponding to a particular sub-domain. So, they do apply this function, SUBSTR, to the source IP, to pull out the first seven six characters. Right? The subnet mask of the first pretext of the, of the IP address. Group by that, and just add up the adRevenue. Okay, so this is, one thing is point, to point out is this actually very nicely and simply expressed as a, as a SQL query. You know, you don't necessarily have to write a bunch of Java code in MapReduce to express it in this particular case. Okay. And so here the results again you see a, a row oriented database. Beating Hadoop, Hadoop. And the reason here is maybe not quite so easy to explain, but essentially it's the, the internal representation in the database, a, again wins credit. There's no parsing that has to happen, okay. Okay, on this same schema there's a join task. which is defined, the sourceIP that generated the most adRevenue along with its average pageRank. And so, this is kind of a complicated thing involving multi-step, multiple passes over the data. You know, to sort of compute the average pageRank and then find the source IP that, Find the maximum adRevenue, find that corresponding source IP and then compute its average, pageRank. And so the implementations here are fairly complex SQL statement involving the use of temporary tables, and in MapReduce, it has to be three separate MapReduce jobs. Chained together, okay? And so for the complicated SQL, we won't go through this in too much detail but just notice that there's a join. And then there's a group buy. You know we looked at some complicated SQL, in which were going to show you how to break them down. And this is no different. So there's a, you know, you know, you know you see two tables which you, you should think yourself joined. And then you see a group buy. And so, those are really the two, tasks going on. And then the second step is to do a big sort, because you see the order by and just find the top most record. Okay. A join in a group, and so here are the results are also pretty imporessive and the reason is, again because of the Indexing, right? This join can be done very, very quickly because there's different kinds of way to do the join. The one we described from [UNKNOWN] is when you have no information abotu the scheme, all you've got are these two big relations and you have to scan them both in parallel. And shuffle them across the network on, with respect to the join key and then perform the join. But if one of them is indexed on that join attribute, you have other plans as available to you. And the database is automatically going to figure out the right one thanks to the magic of relational algebra. And so that's what's going on here. And so both Vertica and the relational database can do a lot better. Okay. Now, so that's fine, so that sort of paints the picture that maybe relational databases are, are great. And, boy, this, this, you know, MapReduce, framework is, is all wet, you know, and why would anyone use it. Well, we talked about fault-tolerance but a couple of other things, you know. There'e other ways to avoid sequential scans that you can actually implement directly in Hadoop. So, for example if you have a large relation and a small relation, one thing you could. And the relation is small in the sense that it, that it fits in, fits on a single node, it doesn't need to be partitioned anymore. Which happens a fair amount. You could actually broadcast that and make a copy of it, and send it to every machine in the cluster, or every, at least every machine that has a, has a copy of, of the other relation he's joined against, right? So, you're joining R and S, and S is small, and R is big. We'll just copy S to every partition of R and not you can do the join locally, without having to do this sort of [UNKNOWN] phase. Okay, and so I didn't get a chance to take advantage of that mechanism. Moreover, there are, especially in modern systems, this paper was in 2009, which is now a little bit old, or quite a bit old. There are ways to provide indexing capability in [UNKNOWN] stack and so. Sort of dead in the water when you aren't allowed to use indexing. So, the positive view, you know if you look sort of warmly on this work, you can think Great, you know, relational databases have all these benefits and G Hadoop can't really compete on somebody's even very basic queries. And another way of looking at this is, well these tricks that we already know work really well, like indexing, do indeed work. And so, all we gotta do is add those to Hadoop and we'll get the same kind of benefits. Okay. So, what's interesting here is to read about the response from Google when this paper came out, which was a discussion published in CACM. And one of their points was that, the largest known database installations were both at Ebay at the time. Which was a Greenplum on a, on about 100 nodes and a Teradata system on a, on also about 100 nodes. And the largest MapReduce in, installations at the time were way, way, way larger. Right? Nearly 4,000 at Yahoo and 600 plus at Facebook and again this is years, years ago. So, these numbers are much higher actually in both cases. But I think the overall point is still the same. The, the size of even perhaps typical Hadoop [UNKNOWN] is, pretty enormous, okay. To conclude the comparison, we said this a couple of times, but just to wrap it up one more time, what can MapReduce learn from databases? And, in the words of the authors of this paper, is that declarative languages are a good thing, schemas are important. And what can databases learn from MapReduce, is this query level fault-tolerance support for what I'm calling in situ data, Which is, you know, data as it lies, right, supporting without, don't require that the database sort of transforms and loaded before you can work with it. And then maybe embrace open-source, because if they, again, if there had been an open source parallel database available, you might not see the same popularity in, in MapReduce. Okay, other systems that are being considered in the same kind of frame work after the fact where Hadoop DB which became Hadapt, which I mentioned is now a start-up. And this is Hadoop filesystem but as Post, it's actually not Postgres anymore it's um, [UNKNOWN]. But the point being a relational database on individual nodes in order to get some of the indexing at some of the at lease local level, benefits of relational query optimization. And then Hive also came out since then. Fine, so that's the end of MapReduced com, [INAUDIBLE] of both, MapReduced itself and MapReduced compared to relational databases. And in the next segment we'll talk about no SQL systems that are solving a slightly different problem.