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

2
00:00:05,678 --> 00:00:10,412
So, let's talk a little bit about what 
kind of systems MapReduce is deployed on. 

3
00:00:10,412 --> 00:00:13,016
So I won't spend to much time on the 
system internals since this is a data 

4
00:00:13,016 --> 00:00:15,830
science course not a distributed systems 
course. 

5
00:00:15,830 --> 00:00:18,520
But it's good to have an intuition for 
what's going on under the hood. 

6
00:00:18,520 --> 00:00:23,074
So, there's three types of systems to be 
aware of, architectures to be aware of: 

7
00:00:23,074 --> 00:00:27,320
shared memory, shared disc and shared 
nothing. 

8
00:00:27,320 --> 00:00:32,336
in these diagrams, these cylinders are 
the discs, these rectangles are the 

9
00:00:32,336 --> 00:00:36,577
memory and the circles are the 
processors. 

10
00:00:36,577 --> 00:00:38,826
Okay. 
So shared memory means that every 

11
00:00:38,826 --> 00:00:42,760
processor has access to all of the memory 
and all of the disc. 

12
00:00:42,760 --> 00:00:44,872
And this is what you think about when you 
have sort of a laptop that has four cores 

13
00:00:44,872 --> 00:00:47,206
in it. 
And if you have a quad core system, or 

14
00:00:47,206 --> 00:00:50,676
[INAUDIBLE] you have you know, six or 12 
cores in your laptop. 

15
00:00:50,676 --> 00:00:58,008
This is the model that's being used, the 
architecture that's being used. 

16
00:00:58,008 --> 00:01:01,912
Shared disc is somewhat less common in, 
at least in these contexts that we're 

17
00:01:01,912 --> 00:01:05,234
talking about. 
Although certainly, certainly common 

18
00:01:05,234 --> 00:01:08,498
enough overall where you'll have the 
individual machines all access a shared 

19
00:01:08,498 --> 00:01:11,738
file system. 
Now that, that setup is very, very 

20
00:01:11,738 --> 00:01:16,226
common, but for using that setup for 
parallel ana, analytics is sort of the 

21
00:01:16,226 --> 00:01:20,650
domain of, of high end commercial 
databases. 

22
00:01:20,650 --> 00:01:24,850
So, you know, your Oracles and your IBMs 
will often used a shared disc 

23
00:01:24,850 --> 00:01:27,555
architecture. 
Okay. 

24
00:01:27,555 --> 00:01:31,510
And then shared nothing is really what 
we're focusing on here. 

25
00:01:31,510 --> 00:01:33,040
And this is what MapReduce is designed 
for. 

26
00:01:33,040 --> 00:01:37,990
And what increasingly, and largely 
parallel databases are designed for as 

27
00:01:37,990 --> 00:01:41,395
well, okay. 
And shared nothing here means individual 

28
00:01:41,395 --> 00:01:43,864
machines that are only connected by a 
network. 

29
00:01:43,864 --> 00:01:46,736
Okay. 
No shared memory, no shared disc, fine. 

30
00:01:46,736 --> 00:01:51,836
So, the argument is that only the shared 
nothing architecture can scale to sort of 

31
00:01:51,836 --> 00:01:55,933
thousands of computers and beyond. 
Okay. 

32
00:01:55,933 --> 00:01:59,827
Because eventually that sh, the the 
sharing of memory or the sharing of disc 

33
00:01:59,827 --> 00:02:03,603
eventually becomes a bottleneck and 
limits how many computers it can attach 

34
00:02:03,603 --> 00:02:09,105
to the same device or, same logical, 
logical or physical device. 

35
00:02:09,105 --> 00:02:12,730
Okay. 
And so learning how to program these 

36
00:02:12,730 --> 00:02:15,634
massive shared nothing clusters is what 
MapReduce and parallel/g, this is all 

37
00:02:15,634 --> 00:02:18,713
about. 
And the shared memory machines are 

38
00:02:18,713 --> 00:02:23,001
perhaps the easiest to program, but are 
conventionally assumed to be pretty 

39
00:02:23,001 --> 00:02:26,780
expensive. 
I should point out that you know, the 

40
00:02:26,780 --> 00:02:31,526
costs are dropping fairly quickly. 
So, it's getting more and more feasible 

41
00:02:31,526 --> 00:02:35,816
to buy a pretty beefy machine with lots 
of main memory and lots of cores, and you 

42
00:02:35,816 --> 00:02:40,470
know, your problem might fit inside that 
one. 

43
00:02:40,470 --> 00:02:44,686
So, you know, when you see people 
deploying Hadoop and MapReduce on fairly 

44
00:02:44,686 --> 00:02:48,640
small clusters in the order of say ten 
nodes. 

45
00:02:49,850 --> 00:02:53,300
You should see whether the data size that 
they're processing are actually all that 

46
00:02:53,300 --> 00:02:56,084
large, right? 
It could be that the data size that 

47
00:02:56,084 --> 00:02:59,396
they're processing is something that fits 
in main memory on a similarly priced 

48
00:02:59,396 --> 00:03:01,650
amount of hardware. 
Okay. 

49
00:03:01,650 --> 00:03:08,610
So, fine, so Hadoop and MapReduce are 
designed for really, really large 

50
00:03:08,610 --> 00:03:12,288
clusters. 
Okay. 

51
00:03:12,288 --> 00:03:15,950
Alright, so this is the context we're in. 
A large amount of commodity servers 

52
00:03:15,950 --> 00:03:18,610
connected by a high-speed commodity 
network. 

53
00:03:18,610 --> 00:03:20,240
And here you know, you think about a data 
center. 

54
00:03:20,240 --> 00:03:22,520
There's a rack that has a number of 
servers and there's a data center that 

55
00:03:22,520 --> 00:03:25,816
has many racks. 
And this is how you organize your 

56
00:03:25,816 --> 00:03:29,800
thousands or tens of thousands of 
computers. 

57
00:03:29,800 --> 00:03:31,125
Okay. 
Alright. 

58
00:03:31,125 --> 00:03:34,581
So you're looking for massive scale 
parallel, parallelism that will run you 

59
00:03:34,581 --> 00:03:38,145
know, jobs that will run for many hours 
even on thousands or tens of thousands of 

60
00:03:38,145 --> 00:03:41,700
servers. 
That's really the context we're in. 

61
00:03:41,700 --> 00:03:45,251
So, when you're in this context an issue 
that comes up that does not come up all 

62
00:03:45,251 --> 00:03:49,190
that often in a much smaller context is 
failures, right? 

63
00:03:49,190 --> 00:03:52,910
So, if you're running a job for a long 
time on thousands of computers, the 

64
00:03:52,910 --> 00:03:57,250
chance of something going wrong during 
that job becomes essentially, you know, 

65
00:03:57,250 --> 00:04:02,290
100% probability, right? 
There's going to be something that fails. 

66
00:04:02,290 --> 00:04:06,754
And so, your system of processing, doing 
this sort of analytics has to just 

67
00:04:06,754 --> 00:04:11,848
tolerate this kind of failure. 
You can't just, you can't, roll back to 

68
00:04:11,848 --> 00:04:16,037
the beginning and just restart every time 
there's a failure occur, or you'd never 

69
00:04:16,037 --> 00:04:19,090
get anything done. 
Okay. 

70
00:04:19,090 --> 00:04:29,370
So, even if the mean time between failure 
for say a disk is a year. 

71
00:04:29,370 --> 00:04:32,430
If you've got 10,000 servers with 
multiple disks or 10,000 disks so say, 

72
00:04:32,430 --> 00:04:35,439
spread across 1,000 servers or any 
combination thereof, you're going to 

73
00:04:35,439 --> 00:04:39,160
start to have failures you know, once per 
hour. 

74
00:04:39,160 --> 00:04:41,722
And you can look up the mean times 
failure and actually do the math, but it, 

75
00:04:41,722 --> 00:04:44,284
there's a, there's a couple of nice 
papers out there that I'll try to put in 

76
00:04:44,284 --> 00:04:47,640
the readings if I remember. 
Okay. 

77
00:04:47,640 --> 00:04:51,418
So, failures are what we're concerned 
about here. 

78
00:04:51,418 --> 00:04:55,129
Alright. 
So, that's hardware, popping up the stack 

79
00:04:55,129 --> 00:04:59,390
one level is this distributed file 
system. 

80
00:04:59,390 --> 00:05:02,530
And you might see HDFS too, which is the 
Hadoop distributor file system. 

81
00:05:02,530 --> 00:05:06,184
So, you remember the context here, was 
that MapReduce proposed in a paper in 

82
00:05:06,184 --> 00:05:09,722
2004 by Google, and Hadoop was the 
implementation of that, of the ideas in 

83
00:05:09,722 --> 00:05:13,570
that paper. 
Okay. 

84
00:05:13,570 --> 00:05:17,794
So, we can almost use an interchangeably 
because the actual implementation in 

85
00:05:17,794 --> 00:05:24,150
Google has certainly evolved since that 
paper and is not completely known, right? 

86
00:05:24,150 --> 00:05:26,590
So, when people are talking about 
MapRoduce, they are typically talking 

87
00:05:26,590 --> 00:05:29,350
about Hadoop or other implementations of 
the program model that have nothing to do 

88
00:05:29,350 --> 00:05:34,180
with sort of the scale out. 
But shared nothing, MapRoduce can be 

89
00:05:34,180 --> 00:05:38,140
assumed to be synonymous with Hadoop. 
Okay. 

90
00:05:38,140 --> 00:05:42,160
So, this is a file system for very large 
files, and the idea here is that if 

91
00:05:42,160 --> 00:05:46,448
you're going to take a single file on 
your own computer, you can manipulate it 

92
00:05:46,448 --> 00:05:51,884
as a single unit. 
But if you're going to take a very, very 

93
00:05:51,884 --> 00:05:55,108
large file and dis, and you know, put it 
on a file system that's a, that's on a 

94
00:05:55,108 --> 00:05:59,144
cluster machine. 
Then there has to be some layer of logic 

95
00:05:59,144 --> 00:06:02,160
that knows how to split that data into 
pieces and put those pieces into 

96
00:06:02,160 --> 00:06:06,170
different machines, and keep track of 
where they are. 

97
00:06:06,170 --> 00:06:08,700
And that's what this distributive files 
software does. 

98
00:06:08,700 --> 00:06:11,607
So each file is printed you know, as 
you're uploading this data to the 

99
00:06:11,607 --> 00:06:14,667
cluster, each file is partitioned into 
chunks that are say typically 64 

100
00:06:14,667 --> 00:06:18,190
megabytes although these are 
configurable. 

101
00:06:18,190 --> 00:06:21,620
And so each chunk, and this is critical, 
is replicated several times. 

102
00:06:21,620 --> 00:06:23,370
Right? 
So there not just one copy of the chunk, 

103
00:06:23,370 --> 00:06:26,260
there might be multiple copies on 
different machines. 

104
00:06:26,260 --> 00:06:28,210
Why? 
Because if one of them goes down, you 

105
00:06:28,210 --> 00:06:31,217
want to have access to the other chunks. 
Okay. 

106
00:06:31,217 --> 00:06:33,791
And you wanted to make sure that these 
are spread across different racks in case 

107
00:06:33,791 --> 00:06:36,480
the entire rack of computers goes dark. 
You still have another copy of it. 

108
00:06:36,480 --> 00:06:40,062
Okay. 
And so the implementations here are GFS 

109
00:06:40,062 --> 00:06:43,589
and HDFS. 
DFS is the concept, and the 

110
00:06:43,589 --> 00:06:47,725
implementations are GFS and HDFS. 
Alright. 

111
00:06:47,725 --> 00:06:53,822
So, here's the phases of MapReduce that's 
a little bit more detailed than the 

112
00:06:53,822 --> 00:06:59,373
abstract phases that we were talking 
about when we were talking about the 

113
00:06:59,373 --> 00:07:03,505
program model. 
Okay. 

114
00:07:03,505 --> 00:07:07,783
So there's a file split here that's read 
from HDFS, and remember HDFS means 

115
00:07:07,783 --> 00:07:12,337
there's replicated, the chunks could be, 
you know, there's multiple copies of 

116
00:07:12,337 --> 00:07:17,144
every chunk. 
And, there's a unit of code called the 

117
00:07:17,144 --> 00:07:20,586
record reader that breaks that chunk 
into. 

118
00:07:20,586 --> 00:07:24,539
I'm using chunk and split for this 
synonymously, I'm not a big fan of the 

119
00:07:24,539 --> 00:07:30,110
term split, because it sort of sounds too 
much like a verb to me. 

120
00:07:30,110 --> 00:07:34,450
the record reader parses that splitter 
chunk into individual records. 

121
00:07:34,450 --> 00:07:38,675
Then the, the programmers, you know, the 
user's map function is called on that 

122
00:07:38,675 --> 00:07:42,338
individual record. 
And then there's a step called combine 

123
00:07:42,338 --> 00:07:46,010
that we haven't talked about, that I'll 
talk about in a moment. 

124
00:07:46,010 --> 00:07:46,964
Next, actually. 
Okay. 

125
00:07:46,964 --> 00:07:50,808
And then the output of these phases are 
written out to local storage on that 

126
00:07:50,808 --> 00:07:56,638
node, as we said. 
So then these regions in local storage, 

127
00:07:56,638 --> 00:08:04,700
one per key, are pulled across the 
network by the reduce phase. 

128
00:08:04,700 --> 00:08:08,848
Then all the regions from all the 
difference, all the different map tasks 

129
00:08:08,848 --> 00:08:14,420
that correspond to the same key are 
sorted together in parallel. 

130
00:08:14,420 --> 00:08:19,378
And finally, the users reduce function 
can be called to produce whatever output 

131
00:08:19,378 --> 00:08:22,486
it produces. 
And then the output of that step is 

132
00:08:22,486 --> 00:08:25,270
actually written back out to HTFS so that 
it's replicated. 

133
00:08:25,270 --> 00:08:31,990
So again, if something goes wrong in the 
map phase, you have to rerun the mapper. 

134
00:08:31,990 --> 00:08:35,977
And if something goes wrong in the reduce 
phase, you have to rerun. 

135
00:08:35,977 --> 00:08:39,281
You can, you can pull the output from 
the, the local storage from the map 

136
00:08:39,281 --> 00:08:43,604
phase. 
and if something goes wrong in the 

137
00:08:43,604 --> 00:08:47,492
overall job, you know you're safe because 
you don't lose data because the HDFS are, 

138
00:08:47,492 --> 00:08:51,614
are, are replicated. 
And if you, you know, if you lose a 

139
00:08:51,614 --> 00:08:55,582
reducer, and you lose a corresponding 
mappers, that's fine can sort of rerun 

140
00:08:55,582 --> 00:09:00,609
whatever you need to rerun, fine. 
So though, the points is that you're, 

141
00:09:00,609 --> 00:09:04,200
you're guaranteeing for fault tolerance 
during Java execution. 

142
00:09:04,200 --> 00:09:19,000
Alright, and so, so let's talk real 
briefly about the combiner. 

143
00:09:19,000 --> 00:09:22,300
So, to think about why we need a 
combiner, go back to this word count 

144
00:09:22,300 --> 00:09:26,764
example that we began with. 
Well, in each case we produced a word and 

145
00:09:26,764 --> 00:09:30,128
just the number one in the fir, in the 
earliest version of this, we just 

146
00:09:30,128 --> 00:09:36,080
produced the number one. 
Right. 

147
00:09:37,370 --> 00:09:40,042
Sending all of these occurrences of word 
one. 

148
00:09:40,042 --> 00:09:43,487
So word, if word one appeared twice, then 
you're going to get a key value pair with 

149
00:09:43,487 --> 00:09:47,720
w1 and a number 1, and another occurrence 
of w1 and a number 1. 

150
00:09:47,720 --> 00:09:50,572
And you're going to send both of these 
guys across the network to be sorted and 

151
00:09:50,572 --> 00:09:54,120
parallel and grouped, in order to be 
processed by the reduced. 

152
00:09:55,200 --> 00:09:59,664
Well that's sort of wasteful. 
You'd like to combine these into a single 

153
00:09:59,664 --> 00:10:04,903
record, w1 comma two. 
And then just send that, because it's 

154
00:10:04,903 --> 00:10:08,110
smaller. 
Well, you could rewrite your map function 

155
00:10:08,110 --> 00:10:13,390
to, to do this. 
But it's such a common need that you can 

156
00:10:13,390 --> 00:10:21,300
that, that, that you have this capability 
called a combiner. 

157
00:10:21,300 --> 00:10:25,448
And so a combiner identifies key value 
pairs of the same key and lumps them 

158
00:10:25,448 --> 00:10:30,000
together before sending it over to the 
reduce side. 

159
00:10:30,000 --> 00:10:33,232
So it just saves the network traffic. 
In many cases, the combiner function can 

160
00:10:33,232 --> 00:10:36,765
be literally the same function as the 
reduce function, it all works out. 

161
00:10:36,765 --> 00:10:40,305
And what needs to be true for this to 
work is that the function that you're 

162
00:10:40,305 --> 00:10:43,845
applying needs to be associative and 
commutative, but I'm not going to say 

163
00:10:43,845 --> 00:10:47,610
much, much more about that. 
Okay. 

164
00:10:47,610 --> 00:10:49,360
So here's what it looks like in pseudo 
code. 

165
00:10:49,360 --> 00:10:52,088
We saw the map function earlier, and 
we're emitting this key value pair of a 

166
00:10:52,088 --> 00:10:55,857
word and account. 
And we saw the reduce function earlier, 

167
00:10:55,857 --> 00:11:00,204
but we're adding new as a combiner 
function that has the same type signature 

168
00:11:00,204 --> 00:11:05,283
as a reducer. 
And again, often in, in many cases can 

169
00:11:05,283 --> 00:11:10,670
literally be the same same implementation 
as a reducer. 

170
00:11:10,670 --> 00:11:15,762
And the only point of this, is that it's 
being applied before sending data across 

171
00:11:15,762 --> 00:11:18,070
the network. 
Alright. 

172
00:11:18,070 --> 00:11:22,492
So here's sort of a summary of a Hadoop 
job that I like a lot, and this is from 

173
00:11:22,492 --> 00:11:27,390
Huy Vo who is, who is now in NYU Poly I 
believe. 

174
00:11:27,390 --> 00:11:32,608
He is still at NYU Poly. 
So the data begins on HTFS, and there are 

175
00:11:32,608 --> 00:11:39,359
these chunks, and in the in part 
partitions go to map tasks. 

176
00:11:39,359 --> 00:11:44,263
Now, these are again not invocations of 
the map function, these are entire tasks. 

177
00:11:44,263 --> 00:11:48,925
And the map, each individual map 
invocation may produce multiple key value 

178
00:11:48,925 --> 00:11:54,010
pairs regardless map task almost 
certainly does, right? 

179
00:11:54,010 --> 00:12:00,495
So it's going to produce these local sort 
of colored key regions. 

180
00:12:00,495 --> 00:12:06,168
And the regions are going to be sent 
across the, pulled across the network to 

181
00:12:06,168 --> 00:12:10,326
the reduce servers. 
And here in this example, there's only 

182
00:12:10,326 --> 00:12:12,999
two reduce servers. 
Now, before we gave this example, we sort 

183
00:12:12,999 --> 00:12:15,855
of showed all the blue ones going 
together and all the red ones going 

184
00:12:15,855 --> 00:12:19,500
together. 
That was a bit of a simplification. 

185
00:12:19,500 --> 00:12:21,831
What actually is going to happen is that 
if you only have two servers, well, all 

186
00:12:21,831 --> 00:12:24,199
the hundreds of, hundreds of possible 
keys need to be mapped to just those two 

187
00:12:24,199 --> 00:12:27,350
servers. 
So you're definitely going to get a mix 

188
00:12:27,350 --> 00:12:30,298
going to the same place. 
And this is where we've been a little bit 

189
00:12:30,298 --> 00:12:33,230
glib up until now. 
We, we've, we've said that you specify 

190
00:12:33,230 --> 00:12:36,479
the key and it hashes to a particular 
reducer. 

191
00:12:36,479 --> 00:12:40,901
whi, which is, which is true logically, 
but before that, you have to get it to a 

192
00:12:40,901 --> 00:12:45,257
machine, where lots of reducers, where 
lots of reduced tasks might be running, 

193
00:12:45,257 --> 00:12:50,590
or lots of reduced invocations might be 
running. 

194
00:12:50,590 --> 00:12:51,908
Okay. 
And so, you know, in this case the blue, 

195
00:12:51,908 --> 00:12:54,410
the blue guys and the red guys both end 
up on the same machine. 

196
00:12:54,410 --> 00:12:57,530
And the green guys And the orange guys 
building it up on the same machine. 

197
00:12:57,530 --> 00:13:03,051
And then this parallel sort manages that. 
Right. 

198
00:13:03,051 --> 00:13:05,040
So, puts all the red things together and 
puts all the blue things together. 

199
00:13:05,040 --> 00:13:08,805
And then for each individual color, one 
reduce indication is called. 

200
00:13:08,805 --> 00:13:13,251
Okay. 
And the reduced function is called and it 

201
00:13:13,251 --> 00:13:21,059
produces the output partition and all 
that Output is wri, written back out to 

202
00:13:21,059 --> 00:13:24,560
HDFS. 
Okay. 

203
00:13:24,560 --> 00:13:29,334
Let me stop there and I'll pick up here, 
and talk a little bit about parallel 

204
00:13:29,334 --> 00:13:33,877
databases and how they do query 
processing with the point being that it 

205
00:13:33,877 --> 00:13:38,244
looks a lot like MapReduce. 
[BLANK_AUDIO] 

