1
23:59:59,500 --> 00:00:05,853
[MUSIC]. 

2
00:00:05,853 --> 00:00:10,473
Okay, so just to wrap up this discussion 
of MapReduce versus databases, I want to 

3
00:00:10,473 --> 00:00:16,659
go over some results from a paper in 2009 
that's on the reading list, 

4
00:00:16,659 --> 00:00:21,390
Where they directly compared Hadoop and a 
couple of different databases. 

5
00:00:21,390 --> 00:00:25,580
And see if we can, maybe explain what 
some of these results tell us. 

6
00:00:25,580 --> 00:00:30,326
Okay. 
So this was Emmy Pablo and some other 

7
00:00:30,326 --> 00:00:40,790
folks at M, MIT and Brown, who did an 
experiment with this kind of a set up. 

8
00:00:40,790 --> 00:00:44,475
So, the comparison was between three 
systems, Hadoop, Vertica, a which was a 

9
00:00:44,475 --> 00:00:49,342
column-oriented database and DBMS-X. 
Which shall remain unnamed, although you 

10
00:00:49,342 --> 00:00:53,177
might be able to figure it out and so, we 
haven't learned what a column oriented 

11
00:00:53,177 --> 00:00:57,374
database is and what a row oriented 
database is. 

12
00:00:57,374 --> 00:01:01,689
But we may have a guest lecture later 
that will describe that in more detail. 

13
00:01:01,689 --> 00:01:05,374
But for right now, for our purposes, just 
think of these as two different kinds of 

14
00:01:05,374 --> 00:01:08,894
relational database or two different 
relational database, with different 

15
00:01:08,894 --> 00:01:13,170
techniques under the hoods, in the, under 
the hood. 

16
00:01:13,170 --> 00:01:15,339
Okay. 
And so, there's two different facets to 

17
00:01:15,339 --> 00:01:18,861
the analysis. 
One was sort of qualitative, about their 

18
00:01:18,861 --> 00:01:23,910
discussion around the programming model 
and the ease of set up, and so on. 

19
00:01:23,910 --> 00:01:27,082
And the other was quantitative, which 
was, performance experiments for 

20
00:01:27,082 --> 00:01:31,779
particular types of queries, okay. 
So, the first task they considered was, 

21
00:01:31,779 --> 00:01:36,378
what they call a Grep task and so this is 
a task to find a 3-byte pattern in a 

22
00:01:36,378 --> 00:01:41,911
100-byte record. 
And the data set was a very, very large 

23
00:01:41,911 --> 00:01:46,910
set of 100-byte records. 
Okay. 

24
00:01:46,910 --> 00:01:50,222
So, this was done in, with, this task was 
performed in the original MapReduce paper 

25
00:01:50,222 --> 00:01:53,740
in 2004, which makes it a good candidate 
for a benchmark. 

26
00:01:55,510 --> 00:01:59,932
And so, the data, the data set here is 10 
billion records with, you know, totaling 

27
00:01:59,932 --> 00:02:04,620
1 terabyte spread across either 25, 50, 
or a 100 nodes. 

28
00:02:04,620 --> 00:02:06,100
Okay, so you are just trying to find this 
record. 

29
00:02:06,100 --> 00:02:12,015
So, this is much like this, you know 
genetic D sequence DNA search task that 

30
00:02:12,015 --> 00:02:20,430
we described as a motivating example for 
sort of describing scalability. 

31
00:02:20,430 --> 00:02:20,620
Okay. 
Fine. 

32
00:02:20,620 --> 00:02:23,207
So what were the results? 
Just to load this data in, this is what 

33
00:02:23,207 --> 00:02:29,556
the story sort of looked like. 
Hadoop and the system called Vertica that 

34
00:02:29,556 --> 00:02:33,576
they are really the, the theme here is 
that they were the designers of the 

35
00:02:33,576 --> 00:02:38,001
Vertica system. 
And so much that these results are 

36
00:02:38,001 --> 00:02:43,290
going to show Vertica doing quite well, 
for, for a variety of reasons. 

37
00:02:43,290 --> 00:02:46,618
So we're not going to talk about too much 
about those particular reasons, we're 

38
00:02:46,618 --> 00:02:49,894
mostly going to be thinking about DBMS-X, 
which is a conventional relational 

39
00:02:49,894 --> 00:02:52,870
database, and Hadoop. 
Okay. 

40
00:02:55,700 --> 00:02:59,060
So here, loading is fast on Hadoop, while 
loading is slow on the db, on the 

41
00:02:59,060 --> 00:03:03,800
relational database, and again it was 
sort of fast on Vertica as well. 

42
00:03:03,800 --> 00:03:08,059
So, why is it faster on Hadoop? 
Well, there's not much to the loading 

43
00:03:08,059 --> 00:03:11,030
right? 
You have to put it into this HTFS system, 

44
00:03:11,030 --> 00:03:14,750
so, it needs to be partitioned, but 
that's about it. 

45
00:03:14,750 --> 00:03:18,636
When you put things into a database, it's 
actually re-casting the data from it's 

46
00:03:18,636 --> 00:03:23,670
raw form, into internal structures in the 
database and that takes time. 

47
00:03:23,670 --> 00:03:27,040
Okay, and the process could be even 
worse. 

48
00:03:27,040 --> 00:03:30,538
Because if you're building indexes over 
the data, you actually, you know every 

49
00:03:30,538 --> 00:03:33,771
time you insert data into the index, you 
need to sort of maintain that data 

50
00:03:33,771 --> 00:03:37,260
structure. 
Okay. 

51
00:03:37,260 --> 00:03:41,448
And so load times are known to be bad. 
So, the take away here is remember that 

52
00:03:41,448 --> 00:03:44,796
load times are typically bad in 
relational databases relative to Hadoop, 

53
00:03:44,796 --> 00:03:49,758
because it has to do more work. 
Now, what is in the database you actually 

54
00:03:49,758 --> 00:03:53,677
get some benefit from that and will see 
in a second these results. 

55
00:03:53,677 --> 00:03:57,316
But actually we know, we know we can 
conform to a schema for example. 

56
00:03:57,316 --> 00:04:00,249
Hadoop is just a pile of bits. 
We don't know anything at all, actually 

57
00:04:00,249 --> 00:04:04,792
we run at out produced task on it, okay. 
And so, how much faster will [UNKNOWN] 

58
00:04:04,792 --> 00:04:11,430
their experiments for the on 25 machines, 
you know, we're up here at 25,000. 

59
00:04:11,430 --> 00:04:16,964
These are all seconds by the way. 
you know, 7500 seconds versus 25,000 and 

60
00:04:16,964 --> 00:04:23,270
a little bit less as we go to more 
servers. 

61
00:04:23,270 --> 00:04:27,928
Okay. 
Now, actually running the Grep task to 

62
00:04:27,928 --> 00:04:35,292
find things, this is what we see. 
Again maybe ignoring Vertica for now, 

63
00:04:35,292 --> 00:04:39,648
because I haven't explained to why, you 
know what the difference about Vertica 

64
00:04:39,648 --> 00:04:44,144
that allows it to, to be so fast. 
But just think about a database from what 

65
00:04:44,144 --> 00:04:50,418
we do understand. 
And Hadoop is, s, slower here and the 

66
00:04:50,418 --> 00:05:02,638
primary reason is that it doesn't have 
access to a index to search. 

67
00:05:02,638 --> 00:05:08,848
Okay. 
So again, no indexes available, Hadoop 

68
00:05:08,848 --> 00:05:17,024
has to do. 
That's wrong. 

69
00:05:17,024 --> 00:05:24,164
Okay, so Hadoop is slower than the 
database, even though both are doing a 

70
00:05:24,164 --> 00:05:29,805
full scan of the data. 
The grep task here is not something 

71
00:05:29,805 --> 00:05:33,100
amenable to any sort of indexing. 
You actually haven't touched any record, 

72
00:05:33,100 --> 00:05:36,590
so there's no fundamental reason why the 
database should be slower or faster. 

73
00:05:36,590 --> 00:05:42,330
But, partially because it gets a win out 
of the structured internal representation 

74
00:05:42,330 --> 00:05:47,578
of the data and doesn't have to re, 
re-parse the raw data from disc like 

75
00:05:47,578 --> 00:05:51,842
Hadoop does. 
And so, I said that there is no 

76
00:05:51,842 --> 00:05:53,790
fundamental reason, there is a 
fundamental reason. 

77
00:05:53,790 --> 00:05:56,826
Because it's already in,in a packed 
internal binary representation, which we 

78
00:05:56,826 --> 00:06:00,790
paid for in the loading phase, but now we 
get the benefit from. 

79
00:06:00,790 --> 00:06:05,008
Here in the query phase, even before we 
even talk about indexes. 

80
00:06:05,008 --> 00:06:10,205
Okay. 
Now, a selection task we're not having to 

81
00:06:10,205 --> 00:06:17,264
scan every record necessarily. 
You know, that is amenable to indexing as 

82
00:06:17,264 --> 00:06:22,640
we discussed in the scalability segment. 
Well, the story is even you know, more 

83
00:06:22,640 --> 00:06:24,642
extreme. 
[UNKNOWN] right? 

84
00:06:24,642 --> 00:06:30,402
The Hadoop results are just way, way, way 
higher than both the database and in 

85
00:06:30,402 --> 00:06:36,386
particular, the Vertica results. 
And so here, the reason is because you 

86
00:06:36,386 --> 00:06:39,690
can build an index on the page rank 
attribute and zoom in directly to the 

87
00:06:39,690 --> 00:06:43,485
records that your, you're interested in. 
Okay. 

88
00:06:43,485 --> 00:06:47,784
Fine. 
So, those are sort of search and 

89
00:06:47,784 --> 00:06:55,910
retrieval tasks not, arguably not exactly 
what Hadoop was designed for. 

90
00:06:55,910 --> 00:06:58,090
Hadoop was designed more for analytically 
tasks. 

91
00:06:58,090 --> 00:07:03,490
So, now consider these, so here the data 
set is 600,000 HTML documents. 

92
00:07:03,490 --> 00:07:05,730
Which works out to be 60GB of data per 
node. 

93
00:07:05,730 --> 00:07:10,876
Along with, another data set is 105, 155 
million user visit records, and 18 

94
00:07:10,876 --> 00:07:18,922
million rankings records. 
So, this is kind of a web data processing 

95
00:07:18,922 --> 00:07:21,760
task. 
Alright. 

96
00:07:21,760 --> 00:07:25,646
And so, a simple aggregate task here is 
to add up all the adRevenue corresponding 

97
00:07:25,646 --> 00:07:30,402
to a particular sub-domain. 
So, they do apply this function, SUBSTR, 

98
00:07:30,402 --> 00:07:35,300
to the source IP, to pull out the first 
seven six characters. 

99
00:07:35,300 --> 00:07:37,628
Right? 
The subnet mask of the first pretext of 

100
00:07:37,628 --> 00:07:42,412
the, of the IP address. 
Group by that, and just add up the 

101
00:07:42,412 --> 00:07:45,622
adRevenue. 
Okay, so this is, one thing is point, to 

102
00:07:45,622 --> 00:07:48,382
point out is this actually very nicely 
and simply expressed as a, as a SQL 

103
00:07:48,382 --> 00:07:51,143
query. 
You know, you don't necessarily have to 

104
00:07:51,143 --> 00:07:54,660
write a bunch of Java code in MapReduce 
to express it in this particular case. 

105
00:07:54,660 --> 00:07:56,930
Okay. 
And so here the results again you see a, 

106
00:07:56,930 --> 00:08:04,890
a row oriented database. 
Beating Hadoop, Hadoop. 

107
00:08:04,890 --> 00:08:09,020
And the reason here is maybe not quite so 
easy to explain, but essentially it's 

108
00:08:09,020 --> 00:08:14,260
the, the internal representation in the 
database, a, again wins credit. 

109
00:08:14,260 --> 00:08:18,070
There's no parsing that has to happen, 
okay. 

110
00:08:18,070 --> 00:08:31,875
Okay, on this same schema there's a join 
task. 

111
00:08:31,875 --> 00:08:37,436
which is defined, the sourceIP that 
generated the most adRevenue along with 

112
00:08:37,436 --> 00:08:41,640
its average pageRank. 
And so, this is kind of a complicated 

113
00:08:41,640 --> 00:08:44,702
thing involving multi-step, multiple 
passes over the data. 

114
00:08:44,702 --> 00:08:47,915
You know, to sort of compute the average 
pageRank and then find the source IP 

115
00:08:47,915 --> 00:08:51,166
that, 
Find the maximum adRevenue, find that 

116
00:08:51,166 --> 00:08:55,840
corresponding source IP and then compute 
its average, pageRank. 

117
00:08:55,840 --> 00:08:59,680
And so the implementations here are 
fairly complex SQL statement involving 

118
00:08:59,680 --> 00:09:03,340
the use of temporary tables, and in 
MapReduce, it has to be three separate 

119
00:09:03,340 --> 00:09:07,010
MapReduce jobs. 
Chained together, okay? 

120
00:09:07,010 --> 00:09:10,850
And so for the complicated SQL, we won't 
go through this in too much detail but 

121
00:09:10,850 --> 00:09:15,620
just notice that there's a join. 
And then there's a group buy. 

122
00:09:15,620 --> 00:09:17,348
You know we looked at some complicated 
SQL, in which were going to show you how 

123
00:09:17,348 --> 00:09:19,120
to break them down. 
And this is no different. 

124
00:09:19,120 --> 00:09:21,560
So there's a, you know, you know, you 
know you see two tables which you, you 

125
00:09:21,560 --> 00:09:24,990
should think yourself joined. 
And then you see a group buy. 

126
00:09:24,990 --> 00:09:27,970
And so, those are really the two, tasks 
going on. 

127
00:09:27,970 --> 00:09:39,730
And then the second step is to do a big 
sort, because you see the order by and 

128
00:09:39,730 --> 00:09:51,360
just find the top most record. 
Okay. 

129
00:09:51,360 --> 00:09:55,329
A join in a group, and so here are the 
results are also pretty imporessive and 

130
00:09:55,329 --> 00:09:59,710
the reason is, again because of the 
Indexing, right? 

131
00:09:59,710 --> 00:10:03,382
This join can be done very, very quickly 
because there's different kinds of way to 

132
00:10:03,382 --> 00:10:06,584
do the join. 
The one we described from [UNKNOWN] is 

133
00:10:06,584 --> 00:10:10,094
when you have no information abotu the 
scheme, all you've got are these two big 

134
00:10:10,094 --> 00:10:14,470
relations and you have to scan them both 
in parallel. 

135
00:10:14,470 --> 00:10:17,620
And shuffle them across the network on, 
with respect to the join key and then 

136
00:10:17,620 --> 00:10:20,699
perform the join. 
But if one of them is indexed on that 

137
00:10:20,699 --> 00:10:24,290
join attribute, you have other plans as 
available to you. 

138
00:10:24,290 --> 00:10:27,350
And the database is automatically 
going to figure out the right one thanks 

139
00:10:27,350 --> 00:10:31,020
to the magic of relational algebra. 
And so that's what's going on here. 

140
00:10:31,020 --> 00:10:34,976
And so both Vertica and the relational 
database can do a lot better. 

141
00:10:34,976 --> 00:10:39,198
Okay. 
Now, so that's fine, so that sort of 

142
00:10:39,198 --> 00:10:43,811
paints the picture that maybe relational 
databases are, are great. 

143
00:10:43,811 --> 00:10:47,903
And, boy, this, this, you know, 
MapReduce, framework is, is all wet, you 

144
00:10:47,903 --> 00:10:57,469
know, and why would anyone use it. 
Well, we talked about fault-tolerance but 

145
00:10:57,469 --> 00:11:08,673
a couple of other things, you know. 
There'e other ways to avoid sequential 

146
00:11:08,673 --> 00:11:17,160
scans that you can actually implement 
directly in Hadoop. 

147
00:11:17,160 --> 00:11:20,570
So, for example if you have a large 
relation and a small relation, one thing 

148
00:11:20,570 --> 00:11:23,660
you could. 
And the relation is small in the sense 

149
00:11:23,660 --> 00:11:26,685
that it, that it fits in, fits on a 
single node, it doesn't need to be 

150
00:11:26,685 --> 00:11:31,780
partitioned anymore. 
Which happens a fair amount. 

151
00:11:31,780 --> 00:11:34,508
You could actually broadcast that and 
make a copy of it, and send it to every 

152
00:11:34,508 --> 00:11:37,148
machine in the cluster, or every, at 
least every machine that has a, has a 

153
00:11:37,148 --> 00:11:41,160
copy of, of the other relation he's 
joined against, right? 

154
00:11:41,160 --> 00:11:44,360
So, you're joining R and S, and S is 
small, and R is big. 

155
00:11:44,360 --> 00:11:48,142
We'll just copy S to every partition of R 
and not you can do the join locally, 

156
00:11:48,142 --> 00:11:52,180
without having to do this sort of 
[UNKNOWN] phase. 

157
00:11:52,180 --> 00:11:55,971
Okay, and so I didn't get a chance to 
take advantage of that mechanism. 

158
00:11:55,971 --> 00:12:01,361
Moreover, there are, especially in modern 
systems, this paper was in 2009, which is 

159
00:12:01,361 --> 00:12:07,157
now a little bit old, or quite a bit old. 
There are ways to provide indexing 

160
00:12:07,157 --> 00:12:10,918
capability in [UNKNOWN] stack and so. 
Sort of dead in the water when you aren't 

161
00:12:10,918 --> 00:12:13,876
allowed to use indexing. 
So, the positive view, you know if you 

162
00:12:13,876 --> 00:12:17,205
look sort of warmly on this work, you can 
think 

163
00:12:17,205 --> 00:12:20,733
Great, you know, relational databases 
have all these benefits and G Hadoop 

164
00:12:20,733 --> 00:12:24,804
can't really compete on somebody's even 
very basic queries. 

165
00:12:24,804 --> 00:12:28,458
And another way of looking at this is, 
well these tricks that we already know 

166
00:12:28,458 --> 00:12:32,120
work really well, like indexing, do 
indeed work. 

167
00:12:32,120 --> 00:12:34,700
And so, all we gotta do is add those to 
Hadoop and we'll get the same kind of 

168
00:12:34,700 --> 00:12:39,270
benefits. 
Okay. 

169
00:12:39,270 --> 00:12:45,210
So, what's interesting here is to read 
about the response from Google when this 

170
00:12:45,210 --> 00:12:51,660
paper came out, which was a discussion 
published in CACM. 

171
00:12:51,660 --> 00:12:56,350
And one of their points was that, the 
largest known database installations were 

172
00:12:56,350 --> 00:13:01,462
both at Ebay at the time. 
Which was a Greenplum on a, on about 100 

173
00:13:01,462 --> 00:13:06,930
nodes and a Teradata system on a, on also 
about 100 nodes. 

174
00:13:06,930 --> 00:13:10,234
And the largest MapReduce in, 
installations at the time were way, way, 

175
00:13:10,234 --> 00:13:12,450
way larger. 
Right? 

176
00:13:12,450 --> 00:13:17,796
Nearly 4,000 at Yahoo and 600 plus at 
Facebook and again this is years, years 

177
00:13:17,796 --> 00:13:21,288
ago. 
So, these numbers are much higher 

178
00:13:21,288 --> 00:13:24,429
actually in both cases. 
But I think the overall point is still 

179
00:13:24,429 --> 00:13:28,408
the same. 
The, the size of even perhaps typical 

180
00:13:28,408 --> 00:13:33,302
Hadoop [UNKNOWN] is, pretty enormous, 
okay. 

181
00:13:34,530 --> 00:13:36,774
To conclude the comparison, we said this 
a couple of times, but just to wrap it up 

182
00:13:36,774 --> 00:13:39,270
one more time, what can MapReduce learn 
from databases? 

183
00:13:39,270 --> 00:13:41,615
And, in the words of the authors of this 
paper, is that declarative languages are 

184
00:13:41,615 --> 00:13:45,102
a good thing, schemas are important. 
And what can databases learn from 

185
00:13:45,102 --> 00:13:48,701
MapReduce, is this query level 
fault-tolerance support for what I'm 

186
00:13:48,701 --> 00:13:52,881
calling in situ data, 
Which is, you know, data as it lies, 

187
00:13:52,881 --> 00:13:56,941
right, supporting without, don't require 
that the database sort of transforms and 

188
00:13:56,941 --> 00:14:02,143
loaded before you can work with it. 
And then maybe embrace open-source, 

189
00:14:02,143 --> 00:14:06,049
because if they, again, if there had been 
an open source parallel database 

190
00:14:06,049 --> 00:14:11,356
available, you might not see the same 
popularity in, in MapReduce. 

191
00:14:11,356 --> 00:14:14,884
Okay, other systems that are being 
considered in the same kind of frame work 

192
00:14:14,884 --> 00:14:18,468
after the fact where Hadoop DB which 
became Hadapt, which I mentioned is now a 

193
00:14:18,468 --> 00:14:22,313
start-up. 
And this is Hadoop filesystem but as 

194
00:14:22,313 --> 00:14:26,038
Post, it's actually not Postgres anymore 
it's um, [UNKNOWN]. 

195
00:14:26,038 --> 00:14:30,866
But the point being a relational database 
on individual nodes in order to get some 

196
00:14:30,866 --> 00:14:35,836
of the indexing at some of the at lease 
local level, benefits of relational query 

197
00:14:35,836 --> 00:14:41,538
optimization. 
And then Hive also came out since then. 

198
00:14:41,538 --> 00:14:47,586
Fine, so that's the end of MapReduced 
com, [INAUDIBLE] of both, MapReduced 

199
00:14:47,586 --> 00:14:54,030
itself and MapReduced compared to 
relational databases. 

200
00:14:54,030 --> 00:14:57,432
And in the next segment we'll talk about 
no SQL systems that are solving a 

201
00:14:57,432 --> 00:14:59,810
slightly different problem. 

