1
00:00:00,845 --> 00:00:04,399
Welcome back to Mining
of Massive Datasets.

2
00:00:04,399 --> 00:00:08,180
We're going to continue our
lecture on MapReduce, and

3
00:00:08,180 --> 00:00:11,382
take a look on the MapReduce
computational model.

4
00:00:11,382 --> 00:00:15,347
So before we look at the actual
MapReduce Programming Model,

5
00:00:15,347 --> 00:00:16,956
let's do a warm up task.

6
00:00:16,956 --> 00:00:22,463
Now imagine you have a huge text document
you know maybe tera, terabytes long and

7
00:00:22,463 --> 00:00:27,260
you want to count, the number of times
each distinct word appears in the file.

8
00:00:27,260 --> 00:00:31,530
For example, we want to find out that
the word the appears 10 million times and

9
00:00:31,530 --> 00:00:36,500
the word you know apple appears 433 times.

10
00:00:36,500 --> 00:00:37,530
Right?

11
00:00:37,530 --> 00:00:42,190
And some sample applications of
this kind of toy example in real

12
00:00:42,190 --> 00:00:45,860
life are you know if you have a big
depth of a log and you want to find out,

13
00:00:45,860 --> 00:00:49,270
how often each URL is accessed that
could be a sample application.

14
00:00:49,270 --> 00:00:52,690
Or it maybe building terms,
such as text for a search engine.

15
00:00:52,690 --> 00:00:53,920
Right?
So, but for

16
00:00:53,920 --> 00:00:56,810
now let's just imagine that
we have this one big file.

17
00:00:56,810 --> 00:00:58,710
That's a huge text document and

18
00:00:58,710 --> 00:01:02,910
our task is to count, the number of times
each distinct word appears in that file.

19
00:01:04,920 --> 00:01:07,290
So, let's look at two cases,

20
00:01:07,290 --> 00:01:11,340
the first case is that the file
itself is too large for memory.

21
00:01:11,340 --> 00:01:14,100
Because remember we said it's a,
it's a, big, big file.

22
00:01:14,100 --> 00:01:18,070
But imagine that there are,
is few enough words in it so

23
00:01:18,070 --> 00:01:21,710
that all the word count pairs
actually fit in memory, right?

24
00:01:22,760 --> 00:01:24,540
How do you solve the problem in this case?

25
00:01:24,540 --> 00:01:29,480
Well it turns out, that in this
case a very simple approach works.

26
00:01:29,480 --> 00:01:31,325
You can just build a, a Hash Table.

27
00:01:35,606 --> 00:01:38,060
I'll, build the, the index by word.

28
00:01:39,330 --> 00:01:44,050
And and the Hash Table for
each word will of course,

29
00:01:44,050 --> 00:01:48,230
will restore the count,
of the number of times, that word appears.

30
00:01:48,230 --> 00:01:52,610
So you the first time you see
a word you initialize you know,

31
00:01:52,610 --> 00:01:56,330
you add an entry to the Hash Table,
with that word, and set the count to 1.

32
00:01:56,330 --> 00:01:59,900
And every subsequent time you see the
word, you, you increment the count by one.

33
00:01:59,900 --> 00:02:03,430
And you,
you make a single sweep through the file.

34
00:02:03,430 --> 00:02:06,830
And at the end of that,
you have the word count pairs for

35
00:02:06,830 --> 00:02:08,470
every unique word that
appears in the file.

36
00:02:08,470 --> 00:02:11,670
So this is a simple program,
that all of us have written, you know,

37
00:02:11,670 --> 00:02:15,030
many many times in some context or
the other.

38
00:02:15,030 --> 00:02:16,730
Now, let's make it a little
bit more complicated.

39
00:02:19,350 --> 00:02:24,438
Let's let's imagine that even the word,
count pairs don't fit in memory.

40
00:02:24,438 --> 00:02:28,010
Right, the file's too big it doesn't
fit in memory, but, there's so

41
00:02:28,010 --> 00:02:31,070
many words,
distinct words in the file that even,

42
00:02:31,070 --> 00:02:33,860
you can't even hold all
the distinct words in memory.

43
00:02:33,860 --> 00:02:34,590
Right?
Now how do

44
00:02:34,590 --> 00:02:36,740
you go about solving
the problem in this case?

45
00:02:36,740 --> 00:02:40,840
Well you can try to write some kind
of complicated code but you know,

46
00:02:40,840 --> 00:02:45,720
I'm lazy so I like to use Unix
file system commands to do this.

47
00:02:45,720 --> 00:02:49,310
And so
here's how I would go about doing this

48
00:02:53,250 --> 00:02:58,630
So this is a a Unix command line way of,
of doing this.

49
00:02:58,630 --> 00:03:04,040
You know here the the command
you know words is,

50
00:03:04,040 --> 00:03:09,633
is, is a little script that goes
through doc.txt which is the,

51
00:03:09,633 --> 00:03:12,820
which is the big text file.

52
00:03:12,820 --> 00:03:16,260
And it outputs the words
in it one per line.

53
00:03:16,260 --> 00:03:21,440
And once once those words are output
I can pipe them to to a sort.

54
00:03:21,440 --> 00:03:26,490
And the sort sorts the you know,
sorts the output of that.

55
00:03:26,490 --> 00:03:29,280
And once you sort it

56
00:03:29,280 --> 00:03:34,190
all of the all occurrences of
the same word come together.

57
00:03:34,190 --> 00:03:38,260
And once you do that you can pipe it to
another little handy utility called uniq

58
00:03:38,260 --> 00:03:46,260
and one of the one of the nifty features
of uniq is the, my, is the dash c option.

59
00:03:46,260 --> 00:03:47,980
And when you do uniq dash c

60
00:03:49,330 --> 00:03:54,920
what uniq dash c does is it takes a run
of the occurrence of the same word.

61
00:03:54,920 --> 00:03:56,800
And then just counts
the occurrences of the same word.

62
00:03:56,800 --> 00:04:01,190
So the output of this is going
to be word count pairs, right?

63
00:04:01,190 --> 00:04:08,220
So and you know I, I'm sure many of
you have done something like this.

64
00:04:08,220 --> 00:04:09,855
And if you've done something like this,

65
00:04:09,855 --> 00:04:12,510
you've actually done something
that's like MapReduce.

66
00:04:12,510 --> 00:04:15,950
Right?
So this case actually captures the essence

67
00:04:15,950 --> 00:04:16,620
of MapReduce.

68
00:04:16,620 --> 00:04:21,010
And the nice thing about this kind of
implementation is that it's, it's very,

69
00:04:21,010 --> 00:04:22,400
very naturally paralle, light,

70
00:04:22,400 --> 00:04:24,920
parallelizable as we'll see in a,
in a moment.

71
00:04:24,920 --> 00:04:32,620
So so let's look at an old view of
MapReduce using this example, right?

72
00:04:32,620 --> 00:04:35,880
So the the first step that we did.

73
00:04:35,880 --> 00:04:38,150
What we,
we took the document which was our input.

74
00:04:39,250 --> 00:04:45,320
And we wrote a script called words,
that output one word to a line, right?

75
00:04:45,320 --> 00:04:50,840
And this is what's called a Map
function in in, in, in MapReduce.

76
00:04:50,840 --> 00:04:55,440
The Map function scans the input
file record-at-a-time.

77
00:04:55,440 --> 00:04:59,760
And for each record, it, it pulls
out something that you care about.

78
00:04:59,760 --> 00:05:01,400
In this case, it was words.

79
00:05:01,400 --> 00:05:02,780
and, and the thing that the you output for

80
00:05:02,780 --> 00:05:07,060
each record you can, you, you cannot one
or multiple things for each records.

81
00:05:07,060 --> 00:05:10,050
And the things that you output,
are called keys, okay?

82
00:05:11,080 --> 00:05:14,850
The second step is is to group by key.

83
00:05:14,850 --> 00:05:18,020
And this is the what the,
the sort step was doing.

84
00:05:18,020 --> 00:05:22,060
It grouped all the keys with
the same value together.

85
00:05:22,060 --> 00:05:22,978
Right?

86
00:05:22,978 --> 00:05:30,784
And the the third step the is the,
the unique minus c step.

87
00:05:30,784 --> 00:05:33,530
That's the reduce piece of MapReduce.

88
00:05:33,530 --> 00:05:37,949
And once the reducer looks
at all the key you know, all

89
00:05:37,949 --> 00:05:43,330
the keys with the same value, and then it,
then it ru, runs some kind of function.

90
00:05:43,330 --> 00:05:48,729
In this case it counted the number of
times the each key occurred but, but

91
00:05:48,729 --> 00:05:52,040
it could be something
much more complicated.

92
00:05:52,040 --> 00:05:56,139
And once it does that kind of analysis it,
it has an answer which,

93
00:05:56,139 --> 00:05:58,190
which it then writes up.

94
00:05:58,190 --> 00:06:00,610
Okay, so this is MapReduce in a,
in a nutshell.

95
00:06:02,760 --> 00:06:07,160
Now the, the outline of this computation
actually stays the stame, same for

96
00:06:07,160 --> 00:06:08,800
any MapReduce computation.

97
00:06:08,800 --> 00:06:11,910
What changes is it that it
change the Map function,

98
00:06:11,910 --> 00:06:14,925
the Reduce function, to the fit
the problem that you're actually solving.

99
00:06:14,925 --> 00:06:15,510
Right?

100
00:06:15,510 --> 00:06:20,424
In this case for the word count the Map
and the Reduce function were quite simple.

101
00:06:20,424 --> 00:06:21,647
In some other problems the Map and

102
00:06:21,647 --> 00:06:23,489
the Reduce functions might
be more complicated.

103
00:06:24,610 --> 00:06:27,390
Here's here's,
here's another way of looking at it.

104
00:06:27,390 --> 00:06:30,396
You start with a,
with a bunch of key value pairs.

105
00:06:30,396 --> 00:06:35,860
And so here's k k v k stands for
key and v stands for value.

106
00:06:36,950 --> 00:06:40,630
And the, the, the, the, the, the Map step.

107
00:06:42,230 --> 00:06:47,570
Takes the key-value pairs and
maps them to intermediate key-value pairs.

108
00:06:47,570 --> 00:06:48,290
Okay?

109
00:06:48,290 --> 00:06:53,995
So for example, you run the Map on the
first key-value pair pair here at k v and

110
00:06:53,995 --> 00:06:58,822
it it actually outputs two
intermediate key-value pairs.

111
00:06:58,822 --> 00:07:02,666
And the, the intermediate key-value
pairs need not have the same key,

112
00:07:02,666 --> 00:07:04,210
as input key value-pairs.

113
00:07:04,210 --> 00:07:06,340
They could be different keys.

114
00:07:06,340 --> 00:07:08,550
And there could be multiple of them.

115
00:07:08,550 --> 00:07:12,900
And the values although they look
the same here, they, they both say v,

116
00:07:12,900 --> 00:07:15,940
the values could be different as well.

117
00:07:15,940 --> 00:07:21,130
And and notice in this case we started
with the one input key-value pair,

118
00:07:21,130 --> 00:07:24,830
and the Map function produced multiple
intermediate key-value pairs.

119
00:07:24,830 --> 00:07:27,490
So there can be zero, one, or

120
00:07:27,490 --> 00:07:33,440
multiple intermediate key-value pairs,
for each, input key-value pair.

121
00:07:33,440 --> 00:07:36,480
Now let's do it again, for
the second key-value pair.

122
00:07:36,480 --> 00:07:38,650
Let's apply the Map function and

123
00:07:38,650 --> 00:07:41,650
it turns out that in this case we
have the one key-value pair in the,

124
00:07:41,650 --> 00:07:44,580
the in the intermediate key-value pair,
in the output.

125
00:07:46,160 --> 00:07:46,880
And so on.
So,

126
00:07:46,880 --> 00:07:49,720
so, we,
we run through the entire input file.

127
00:07:49,720 --> 00:07:52,562
Apply the Map function
to each input record.

128
00:07:52,562 --> 00:07:55,430
And create intermediate key-value pairs.

129
00:07:55,430 --> 00:08:00,590
Now the next step, is to take these
intermediate key-value pairs,

130
00:08:00,590 --> 00:08:03,070
and group them by key.

131
00:08:03,070 --> 00:08:03,880
Right?
So,

132
00:08:03,880 --> 00:08:08,080
all the intermediate key-value pairs that
have the same key, are grouped together.

133
00:08:08,080 --> 00:08:11,570
So it turns out that
there are three values.

134
00:08:11,570 --> 00:08:14,840
With the, with the first key,
two values for the second key and so on.

135
00:08:14,840 --> 00:08:19,230
And they all get grouped together,
and this is done by sorting by key and

136
00:08:19,230 --> 00:08:23,348
then by grouping together the value of,
you know the values for the same key.

137
00:08:23,348 --> 00:08:24,964
And these are all different values,

138
00:08:24,964 --> 00:08:27,059
although I use the same
same symbol v here.

139
00:08:29,740 --> 00:08:34,600
Now, once you have once you have these
key value groups then the final step is

140
00:08:34,600 --> 00:08:36,030
the reducer.

141
00:08:36,030 --> 00:08:40,197
The reducer takes a look at a,

142
00:08:40,197 --> 00:08:47,150
a single a single key-value
group as input and.

143
00:08:47,150 --> 00:08:52,120
It produces produces an output
that has the sa, you know,

144
00:08:52,120 --> 00:08:56,930
that has the same key but
it combines the the, the, the,

145
00:08:56,930 --> 00:09:01,260
the values or the values for
a given key, into a single value.

146
00:09:01,260 --> 00:09:04,020
For example,
it could add up all the values.

147
00:09:05,470 --> 00:09:08,080
In the, or, or it could or
it could multiply them, or

148
00:09:08,080 --> 00:09:09,540
it could do, it could take the average.

149
00:09:09,540 --> 00:09:11,580
Or it can do something more complicated.

150
00:09:11,580 --> 00:09:14,250
But with all of the values for
a given key.

151
00:09:14,250 --> 00:09:17,890
And finally you, the output,
it outputs a single value for the key.

152
00:09:20,250 --> 00:09:21,620
Right?
And so, when you,

153
00:09:21,620 --> 00:09:24,830
when you apply the reducer to
the second key-value group.

154
00:09:24,830 --> 00:09:27,150
You get, you get another output and so on.

155
00:09:27,150 --> 00:09:31,290
And once you apply the reducer to all
the intermediate key-value groups.

156
00:09:31,290 --> 00:09:32,440
You get the final output.

157
00:09:37,040 --> 00:09:43,660
So more formally the input to
MapReduce is a set of key-value pairs.

158
00:09:43,660 --> 00:09:46,240
And the programmer has
to specify two methods.

159
00:09:46,240 --> 00:09:48,634
The first method is a Map method.

160
00:09:48,634 --> 00:09:52,779
And the Map method takes
an input key-value pair.

161
00:09:52,779 --> 00:09:57,920
And produces an int, an set of
intermediate key-value pair, zero or

162
00:09:57,920 --> 00:10:00,890
more intermediate key-value pairs.

163
00:10:00,890 --> 00:10:04,400
and, there is one Map call, for
every input key-value pairs.

164
00:10:05,410 --> 00:10:10,195
The Reduce function, takes an intermediate
key-value group the intermediate key-value

165
00:10:10,195 --> 00:10:13,840
group consists of a key, and
a set of values for that key.

166
00:10:13,840 --> 00:10:18,680
and, the output can consist of one,
zero, one, or

167
00:10:18,680 --> 00:10:22,250
multiple key-value pairs once again.

168
00:10:22,250 --> 00:10:27,140
The key is the same as the as
the input key but the value is, is,

169
00:10:27,140 --> 00:10:31,290
is is obtained by combining,
the input values in some manner.

170
00:10:33,480 --> 00:10:39,137
For example, you might add up the you
know, add up the input values and

171
00:10:39,137 --> 00:10:42,700
that could be the output
v double prime here.

172
00:10:47,010 --> 00:10:50,402
So let's look at the the word
count example and

173
00:10:50,402 --> 00:10:52,970
run that through
the MapReduce process again.

174
00:10:52,970 --> 00:10:54,049
Here's our big document.

175
00:10:55,081 --> 00:11:00,150
And I hope you can see the text of
this you know, the document but

176
00:11:00,150 --> 00:11:03,250
it doesn't matter,
you can see that there are words in there.

177
00:11:03,250 --> 00:11:06,850
And so
we're going to take this big document.

178
00:11:06,850 --> 00:11:10,782
And we're going to take the Map function
that's provided by the programmer.

179
00:11:10,782 --> 00:11:14,961
The Map function reads the input,
and produces a, produces a set of

180
00:11:14,961 --> 00:11:21,000
key-value pairs, and the key-value pairs
in this case are going to be the key.

181
00:11:21,000 --> 00:11:24,530
Each word is going to be a key, and
the value is going to be the number 1.

182
00:11:24,530 --> 00:11:25,762
Right?

183
00:11:25,762 --> 00:11:29,980
so, for example,
the word the and 1 crew and

184
00:11:29,980 --> 00:11:34,210
1 and so on, and
the word the appears again.

185
00:11:35,316 --> 00:11:40,210
And so there, there's another the,
1 here and so on.

186
00:11:40,210 --> 00:11:43,150
So these are the intermediate
key-value pairs,

187
00:11:43,150 --> 00:11:44,900
that are produced by the Map function.

188
00:11:44,900 --> 00:11:50,975
[SOUND] Now the next step is
the group by key step which

189
00:11:50,975 --> 00:11:56,536
collects together all
pairs with the same key.

190
00:11:56,536 --> 00:12:02,406
So we can see that the, there are two
tuples two intermediate tuples with the,

191
00:12:02,406 --> 00:12:07,290
with the key crew and
then those are collected together here.

192
00:12:07,290 --> 00:12:09,040
There's one with you know,

193
00:12:09,040 --> 00:12:12,870
with the word space, there are three
with the word the, and so on.

194
00:12:12,870 --> 00:12:15,810
And they're all sorted and
collected together.

195
00:12:15,810 --> 00:12:18,880
In this yeah, in, in this place here.

196
00:12:18,880 --> 00:12:22,640
And the,
the final step is the Reduce step.

197
00:12:22,640 --> 00:12:26,910
The Reduce, the Reduce step
collects together all the values so

198
00:12:26,910 --> 00:12:30,400
the Reduce step adds,
adds together the 2, 1 from crew.

199
00:12:31,600 --> 00:12:33,790
and, and
figures out that there are two you know,

200
00:12:33,790 --> 00:12:35,936
two occurrences of the word crew.

201
00:12:35,936 --> 00:12:37,150
Space has 1.

202
00:12:37,150 --> 00:12:40,580
There are 3 tuples with the 1 for
the there all added together.

203
00:12:40,580 --> 00:12:43,130
And the output is 3, and so on.

204
00:12:43,130 --> 00:12:48,070
Right, so this is a schematic, of
the the MapReduce word counting example.

205
00:12:48,070 --> 00:12:52,190
now, of course this, this whole
example doesn't run on a single node.

206
00:12:52,190 --> 00:12:54,680
The data is actually distributed
across multiple input nodes.

207
00:12:54,680 --> 00:12:56,250
So let's take that into account.

208
00:12:57,360 --> 00:12:59,490
And see here's here's the data.

209
00:12:59,490 --> 00:13:02,690
The data's actually divided here into,
into multiple nodes.

210
00:13:02,690 --> 00:13:06,240
Let's say the, the red the, the,

211
00:13:06,240 --> 00:13:10,600
the first portion of, of the file is
it's chunk one, and it's on one node.

212
00:13:10,600 --> 00:13:13,960
The second portion of the file here is
chunk two, which is on a different node.

213
00:13:13,960 --> 00:13:17,650
The third portion is chunk three, and
the fourth portion is chunk four, and

214
00:13:17,650 --> 00:13:19,930
each of these is on a different node.

215
00:13:19,930 --> 00:13:23,505
Now the Map tasks are going to be run
on each of these four different nodes.

216
00:13:23,505 --> 00:13:28,080
There going to be a Map task that's run on
chunk one that just looks at this portion,

217
00:13:28,080 --> 00:13:29,510
the first portion of the file.

218
00:13:29,510 --> 00:13:31,370
Map task is run on chunk two that,

219
00:13:31,370 --> 00:13:33,940
that, that just looks at the second
portion of the file and so on.

220
00:13:34,980 --> 00:13:39,600
And the the outputs of those Map
tasks will therefore be produced,

221
00:13:39,600 --> 00:13:41,440
on on four different nodes.

222
00:13:44,210 --> 00:13:49,460
li, like so so here are the,
here are the first chunk of Map output.

223
00:13:49,460 --> 00:13:52,330
The second chunk of Map output,
which is on another node.

224
00:13:52,330 --> 00:13:54,950
The third chunk of Map output,
which is on a third node.

225
00:13:54,950 --> 00:13:57,950
And the fourth chunk of Map output,
which is on yet another node.

226
00:13:57,950 --> 00:13:58,450
Right.

227
00:13:59,480 --> 00:14:03,150
Now the output of the of,
of the Map functions,

228
00:14:03,150 --> 00:14:06,812
are therefore spread
across multiple nodes.

229
00:14:06,812 --> 00:14:10,612
And what the system then does,
is that it it,

230
00:14:10,612 --> 00:14:15,059
it copies the, the Map outputs,
onto a single node.

231
00:14:17,230 --> 00:14:18,860
And then so

232
00:14:18,860 --> 00:14:23,276
you can see the data from all these four
nodes flowing into this single node here.

233
00:14:23,276 --> 00:14:26,498
And once the data has,
has flowed to the single node,

234
00:14:26,498 --> 00:14:30,590
it can then sort it by key and
then do the final radial step.

235
00:14:32,280 --> 00:14:35,613
Now it's a little bit trickier
than this unfortunately.

236
00:14:35,613 --> 00:14:39,760
Because you know,
you may not want to use you know,

237
00:14:39,760 --> 00:14:42,300
to, to move all the data
from all the Map nodes,

238
00:14:42,300 --> 00:14:45,660
going to be a lot of it, into a single
Reduce node, and sort it there.

239
00:14:45,660 --> 00:14:48,400
That might be a lot of you know,
a lot of sorting.

240
00:14:48,400 --> 00:14:53,000
So in practice you use
multiple Reduce nodes as well.

241
00:14:53,000 --> 00:14:57,840
And you, you know, when you run
a MapReduce job, you can say you know,

242
00:14:57,840 --> 00:15:00,334
you can tell the system to use
a certain number of Reduce nodes.

243
00:15:00,334 --> 00:15:04,240
Let's say you tell the system in this
case, to use three Reduce nodes.

244
00:15:04,240 --> 00:15:09,980
So if you use three reduce
nodes then the then

245
00:15:09,980 --> 00:15:15,330
the MapReduce system is smart
enough to split the the, the,

246
00:15:15,330 --> 00:15:20,660
the output of the Map into, into three,
three into three Reduce nodes.

247
00:15:20,660 --> 00:15:25,130
And it makes sure, that for
any given key in this case the,

248
00:15:25,130 --> 00:15:30,990
all instances of the, regardless of
which Map node they start out from,

249
00:15:30,990 --> 00:15:33,540
always end up at the same Reduce node,
right?

250
00:15:33,540 --> 00:15:37,250
So all instances of the,
whether it started from Map node one or

251
00:15:37,250 --> 00:15:41,450
Map node two ended up at Reduce node two,
in this case.

252
00:15:41,450 --> 00:15:45,960
And all instances of the word crew
regardless of whether they started from

253
00:15:45,960 --> 00:15:50,620
Map node one or Map node four,
ended up at Reduce node one.

254
00:15:50,620 --> 00:15:52,970
And this is done by using a hash function,
right?

255
00:15:52,970 --> 00:15:57,910
So the system uses a hash function
that hashes each Map key and

256
00:15:57,910 --> 00:16:01,460
determines a single Reduce
node to shift that tuple two.

257
00:16:02,560 --> 00:16:06,540
And this ensures that all
tuples with the same key,

258
00:16:06,540 --> 00:16:08,450
end up with the same Reduce node.

259
00:16:08,450 --> 00:16:13,149
And once once tuples end up at a Reduce
node, they get sorted as before.

260
00:16:14,540 --> 00:16:17,950
in, on each Reduce node and and, and

261
00:16:17,950 --> 00:16:21,040
the result is created now,
on multiple Reduce nodes.

262
00:16:21,040 --> 00:16:28,770
For example the result for
crew is now on is now on Reduce node one.

263
00:16:28,770 --> 00:16:33,180
The result for the is now on,
on the Reduce node two and the result for

264
00:16:33,180 --> 00:16:35,650
shuttle and
recently are on Reduce node three.

265
00:16:35,650 --> 00:16:41,560
So the final result is actually now
spread across three nodes in the system.

266
00:16:41,560 --> 00:16:45,320
Which is perfectly fine because you're
dealing with a distributed file system,

267
00:16:45,320 --> 00:16:50,040
which know, knows that your file is
spread across three nodes of the system.

268
00:16:50,040 --> 00:16:53,530
So you can still access it as
a single file in your client.

269
00:16:53,530 --> 00:16:57,889
And the system knows to access the data
from those three three independent nodes.

270
00:17:00,050 --> 00:17:04,320
One final point before we move
on from the slide is that

271
00:17:05,590 --> 00:17:11,370
all this magic in the MapReduce
magic is implemented to use

272
00:17:11,370 --> 00:17:16,680
as far as possible, only sequential scans
of disk as opposed to a random access is.

273
00:17:16,680 --> 00:17:20,130
If you think a little bit carefully,
what all the steps that I mentioned

274
00:17:20,130 --> 00:17:24,370
about how the Map function is applied
on the input file record by record.

275
00:17:24,370 --> 00:17:26,286
How the sorting is done and so on.

276
00:17:26,286 --> 00:17:30,240
A moment's thought will make it apparent
that you can actually implement,

277
00:17:30,240 --> 00:17:34,280
all of this by using only
sequential reads of disk, and

278
00:17:34,280 --> 00:17:36,710
never using random accesses of disk.

279
00:17:36,710 --> 00:17:40,344
Now this is super important
because sequential reads are much,

280
00:17:40,344 --> 00:17:44,010
much more efficient than
random accesses to disk.

281
00:17:44,010 --> 00:17:49,190
If you don't learn your basics of
of database systems, it takes much,

282
00:17:49,190 --> 00:17:51,510
much longer to do random seeks.

283
00:17:51,510 --> 00:17:53,970
Than to do a single
sequential axis of a file.

284
00:17:53,970 --> 00:17:57,450
And that's why the, the MapReduce,
the whole MapReduce system,

285
00:17:57,450 --> 00:18:02,870
is built around doing only sequential
reads of files and never random accesses.

286
00:18:02,870 --> 00:18:08,780
So here is the actual pseudocode for
for the word count using MapReduce.

287
00:18:08,780 --> 00:18:13,902
Remember, the programmer is required to
provide two functions, a Map function,

288
00:18:13,902 --> 00:18:15,592
a Reduce function.

289
00:18:15,592 --> 00:18:17,580
And this is the the Map
function right here.

290
00:18:17,580 --> 00:18:21,877
The Map function takes a key and
a value and its output has to be int,

291
00:18:21,877 --> 00:18:25,732
an intermi,
a set of intermediate key-value pairs.

292
00:18:25,732 --> 00:18:27,298
Now the key in this case is,

293
00:18:27,298 --> 00:18:30,989
is a document name and
the value is the text of the document.

294
00:18:32,010 --> 00:18:36,720
And the Map the Map function itself
is very simple in this case.

295
00:18:36,720 --> 00:18:39,680
It scans the the input document.

296
00:18:39,680 --> 00:18:44,280
And for each word in the input document
[INAUDIBLE] the input document.

297
00:18:44,280 --> 00:18:48,840
For each word in the input document
it emits that word and the number 1.

298
00:18:48,840 --> 00:18:52,480
So, so it's, it's a tuple,
whose key is the, is the word.

299
00:18:52,480 --> 00:18:54,317
And whose value is the number 1.

300
00:18:54,317 --> 00:18:58,710
And here's the reduced function.

301
00:18:58,710 --> 00:19:02,570
The reduced function, remember,
takes a key and a set of values.

302
00:19:02,570 --> 00:19:07,410
The set of values all correspond
to the same key and in this case,

303
00:19:07,410 --> 00:19:12,380
they just iterate through all
the values and and, and sums them up.

304
00:19:12,380 --> 00:19:16,631
And the output has the same key and
the value is the, is the sum.

305
00:19:21,890 --> 00:19:26,510
We looked at a very simple example
of bullet count using MapReduce.

306
00:19:26,510 --> 00:19:28,910
Now let's look at a couple more examples.

307
00:19:28,910 --> 00:19:32,480
Here's here's here's another example.

308
00:19:32,480 --> 00:19:38,065
Suppose we have a large web
corpus that we've called and, for

309
00:19:38,065 --> 00:19:43,149
each and we have a metadata file for
a, for [INAUDIBLE] and

310
00:19:43,149 --> 00:19:48,560
each record in the metadata file lo,
looks like this.

311
00:19:48,560 --> 00:19:54,080
It has a URL the size of the file,

312
00:19:54,080 --> 00:19:57,650
the date and
then various other pieces of data.

313
00:19:57,650 --> 00:20:02,090
Now the problem, is for each host we
want to find the total number of bytes.

314
00:20:03,650 --> 00:20:05,880
And not for each URL, but for each host.

315
00:20:05,880 --> 00:20:11,110
Remember, there can be multiple
many URLs with the same host name,

316
00:20:11,110 --> 00:20:15,580
in the crawl and you want to find the
number of bytes associated with each host,

317
00:20:15,580 --> 00:20:17,550
not with each URL, right?

318
00:20:17,550 --> 00:20:20,800
Clearly the the number of bytes
associated with the host,

319
00:20:20,800 --> 00:20:25,144
is just the sum of the number of bytes
associated with all the URLs for a,

320
00:20:25,144 --> 00:20:30,020
for the host and this is very easy
to implement in in, in MapReduce.

321
00:20:30,020 --> 00:20:36,360
The mapper in this case, the Map
function just looks at each record and

322
00:20:36,360 --> 00:20:41,400
it looks at the URL of, of the record and
outputs the hostname of the URL.

323
00:20:41,400 --> 00:20:44,920
and, and, and the size, right?

324
00:20:44,920 --> 00:20:51,890
And the the Reduce function just
sums the sizes for each host, right?

325
00:20:51,890 --> 00:20:57,240
And at the end of it, you will have
this the, the, the size of each host.

326
00:21:01,076 --> 00:21:02,150
Here's another example.

327
00:21:02,150 --> 00:21:05,350
Let's say you're building
a language model by year.

328
00:21:05,350 --> 00:21:07,320
You have a large collection of documents.

329
00:21:07,320 --> 00:21:12,920
And you want to build a language model and
and this language model for

330
00:21:12,920 --> 00:21:17,290
some reason requires the count
of every 5-word sequence.

331
00:21:17,290 --> 00:21:20,390
Every unique 5-word sequence that
occurs in a large corpus of document.

332
00:21:21,490 --> 00:21:25,024
Earlier we looked at accounting,
each unique word.

333
00:21:25,024 --> 00:21:28,100
This example ask for each 5-word sequence.

334
00:21:28,100 --> 00:21:31,960
It turns out that the solution
is not very different.

335
00:21:31,960 --> 00:21:33,650
the, just the Map function differs.

336
00:21:33,650 --> 00:21:36,740
The Map function extracts you know,

337
00:21:36,740 --> 00:21:40,900
goes through each document and outputs
every 5-word sequence in the document.

338
00:21:42,010 --> 00:21:45,980
And the the Reduce function
just combines those counts and

339
00:21:45,980 --> 00:21:48,010
adds them up and then you have the output.

340
00:21:49,740 --> 00:21:55,812
So I hope these simple examples
illustrate how MapReduce works.

341
00:21:55,812 --> 00:22:00,367
In the next section, we are going to
understand how the underlying system,

342
00:22:00,367 --> 00:22:03,930
actually implements some of
the magic that makes MapReduce work.

