1
00:00:00,206 --> 00:00:05,919
[MUSIC]. 

2
00:00:05,919 --> 00:00:10,404
So, last time we talked about parallel 
processing as a lead up to map reduce. 

3
00:00:10,404 --> 00:00:14,455
And we ended up with this schematic here. 
In the context of this example, where 

4
00:00:14,455 --> 00:00:16,802
we're counting words across a set of 
documents. 

5
00:00:16,802 --> 00:00:21,810
Just sort a canonical example to start 
thinking about programming a map reduce. 

6
00:00:21,810 --> 00:00:25,408
And so each one of these vertical black 
lines represented a document. 

7
00:00:25,408 --> 00:00:30,440
And we split them into smaller sets and 
send each one of those sets to a separate 

8
00:00:30,440 --> 00:00:34,212
machine. 
And then we applied our map function to 

9
00:00:34,212 --> 00:00:39,160
each one of those documents in turn. 
And so the map function if you recall, 

10
00:00:39,160 --> 00:00:45,760
took a single document and produced a set 
of pairs, and each pair was a word along 

11
00:00:45,760 --> 00:00:55,630
with a count of the number of occurrences 
of that word in that document. 

12
00:00:55,630 --> 00:00:58,785
Various variations on this that you can 
imagine. 

13
00:00:58,785 --> 00:01:03,063
Okay, now, this word may have appeared in 
multiple documents, one here, one on this 

14
00:01:03,063 --> 00:01:06,597
machine, one on this machine and so on, 
and so now we need to group them 

15
00:01:06,597 --> 00:01:12,230
altogether onto a single machine so that 
we can count them. 

16
00:01:12,230 --> 00:01:15,205
And that's exactly what this shuffle 
phase did. 

17
00:01:15,205 --> 00:01:21,390
Okay. 
So here I've written four different 

18
00:01:21,390 --> 00:01:26,626
tasks, processing sort of, looking like 
it's processing a single group at a time, 

19
00:01:26,626 --> 00:01:31,554
but you know, I want you to think about 
how many map tasks do we already have and 

20
00:01:31,554 --> 00:01:37,860
how many reviews tasks are we going to 
have? 

21
00:01:37,860 --> 00:01:42,610
Well, the map tasks are one per document. 
We have to, we have to call the Map 

22
00:01:42,610 --> 00:01:45,290
function, how many invocations of the map 
functions are going to there going to be? 

23
00:01:45,290 --> 00:01:46,848
Well, we're going to call it once per 
document. 

24
00:01:46,848 --> 00:01:49,770
How many invocations of the Reduce 
function are going to be? 

25
00:01:49,770 --> 00:01:52,830
Well, it's the number of groups that are 
produced by the output of the Map 

26
00:01:52,830 --> 00:01:56,788
function. 
In this case, it's once per unique word 

27
00:01:56,788 --> 00:02:02,192
appearing in any, in any document, okay? 
And so in some sense, the number of 

28
00:02:02,192 --> 00:02:05,426
machines we need to apply to this problem 
is maybe kind of predictable in the map 

29
00:02:05,426 --> 00:02:08,340
phase. 
It's, it's. 

30
00:02:08,340 --> 00:02:10,865
Corresponds to the size of the input data 
set. 

31
00:02:10,865 --> 00:02:14,944
Which we, you know, presume to know. 
But the number of reducers we're going to 

32
00:02:14,944 --> 00:02:17,460
need is maybe not known ahead of time, 
right? 

33
00:02:17,460 --> 00:02:20,530
It depends on the size of the Map output. 
Here we might be able to reason about it, 

34
00:02:20,530 --> 00:02:22,850
because we maybe know how many words 
there are in the English language. 

35
00:02:22,850 --> 00:02:25,187
And we can assume that with a big enough 
set, all of those words will be 

36
00:02:25,187 --> 00:02:28,200
represented. 
at least once. 

37
00:02:28,200 --> 00:02:30,335
But in general, it's, it's dependent on 
the output of the map, so you don't 

38
00:02:30,335 --> 00:02:33,397
really know. 
Okay, and so, the only point I want to 

39
00:02:33,397 --> 00:02:39,400
make is that we made a decision here to 
draw it as four different machines. 

40
00:02:39,400 --> 00:02:42,400
But it may you know, it may, it may be 
the same six machines you used the Map 

41
00:02:42,400 --> 00:02:45,720
phase, or maybe a thousand machines and 
so on. 

42
00:02:45,720 --> 00:02:50,270
You know, nothing's stopping you from 
sending all of the word occurrences to a 

43
00:02:50,270 --> 00:02:54,169
single machine. 
And having this one machine process the 

44
00:02:54,169 --> 00:02:59,320
green group then process the red group 
then process the blue group and so on. 

45
00:02:59,320 --> 00:03:03,939
Or maybe it would be four at a time 
because quarters in the machine, okay. 

46
00:03:03,939 --> 00:03:07,656
But that wouldn't be as, perhaps as 
efficient because it would be doing a lot 

47
00:03:07,656 --> 00:03:10,690
of. 
Serial work. 

48
00:03:10,690 --> 00:03:14,112
At the other extreme you might think, 
well we're going to need millions of 

49
00:03:14,112 --> 00:03:16,862
tasks. 
Let's allocate, you know, hundreds of 

50
00:03:16,862 --> 00:03:20,702
thousands of machines to process them. 
So that each machine is doing very little 

51
00:03:20,702 --> 00:03:23,360
work. 
Okay. 

52
00:03:23,360 --> 00:03:26,349
And that might make sense, but then the 
trade-off is perhaps sort of long, 

53
00:03:26,349 --> 00:03:30,510
spinning of all these machines and, and 
preparing them to do the work. 

54
00:03:30,510 --> 00:03:32,500
Okay. 
So there's a decision to make there. 

55
00:03:32,500 --> 00:03:35,414
And we'll, we'll come back to that. 
Okay. 

56
00:03:35,414 --> 00:03:38,960
But let's talk a little bit more about 
Map Reduce itself. 

57
00:03:38,960 --> 00:03:41,740
This is, I'm, I'm belaboring this for a 
reason. 

58
00:03:41,740 --> 00:03:44,280
Cause I, what I want you to do is start 
thinking in terms of Map Reduce. 

59
00:03:44,280 --> 00:03:48,351
Every problem you have, think what if the 
data set was absolutely enormous, way too 

60
00:03:48,351 --> 00:03:52,680
big for a machine? 
How we're going to split it into pieces? 

61
00:03:52,680 --> 00:03:54,912
And a, a very good way of thinking about 
how to split things into pieces is to 

62
00:03:54,912 --> 00:03:57,144
think about how you'd write a map reduce 
program to do whatever it is you're 

63
00:03:57,144 --> 00:04:00,503
trying to do. 
And so this is, yet again the same 

64
00:04:00,503 --> 00:04:05,456
example, just drawn a different way. 
So here the input is document ID followed 

65
00:04:05,456 --> 00:04:10,630
by a value, and the value here is the 
entire text of the document. 

66
00:04:10,630 --> 00:04:15,910
And the Map Function, just to make this 
clear, produces a set of things, not just 

67
00:04:15,910 --> 00:04:20,470
one thing. 
And they're shuffled to produce this. 

68
00:04:20,470 --> 00:04:25,334
So this is word one with a count of one. 
Word two with a count of one, word three 

69
00:04:25,334 --> 00:04:31,170
with a count of one and so on. 
And the other side what we get is word 

70
00:04:31,170 --> 00:04:35,174
one with a group. 
Of all the occurrences, 1, 1,1 1, 1, 1, 

71
00:04:35,174 --> 00:04:38,522
1, and then find the reduce function 
counts them all up and finds that there 

72
00:04:38,522 --> 00:04:41,280
are 25 occurrences. 
Okay. 

73
00:04:41,280 --> 00:04:45,245
So I'm probably, I guess if, if this is 
completely obvious, you can always fast 

74
00:04:45,245 --> 00:04:49,750
forward I guess one of the beauties of 
doing this online. 

75
00:04:49,750 --> 00:04:52,060
Okay. 
So fine. 

76
00:04:52,060 --> 00:04:55,042
So what is map reduce? 
That's the programming model we just 

77
00:04:55,042 --> 00:04:59,320
described and there's a paper in 2004 
that's on the reading list that describes 

78
00:04:59,320 --> 00:05:02,760
this. 
And there's a, a couple of key 

79
00:05:02,760 --> 00:05:06,969
motivations for doing this. 
In that paper that I think sometimes get 

80
00:05:06,969 --> 00:05:10,650
lost when you hear about the popularity 
of Map Reduce today. 

81
00:05:10,650 --> 00:05:15,900
and we'll talk about those two, those two 
benefits in, in, in, in a moment. 

82
00:05:17,240 --> 00:05:22,170
So. 
So one thing to realize is that 

83
00:05:22,170 --> 00:05:26,458
map-reduce refers to the abstraction and 
it's the name given to it by the authors 

84
00:05:26,458 --> 00:05:30,935
of this 2004 paper. 
Hadoop is an implementation of Map 

85
00:05:30,935 --> 00:05:33,400
Reduce. 
It came a few years later. 

86
00:05:33,400 --> 00:05:37,824
And was written by some people at Yahoo! 
Originally and then became an open source 

87
00:05:37,824 --> 00:05:42,870
product that is managed by the Apache and 
has lots of contributors, okay. 

88
00:05:42,870 --> 00:05:46,421
So, you know the key idea for Map Reduce 
was really this programming model, which 

89
00:05:46,421 --> 00:05:50,319
it says here, right. 
Now it had a system with it as well, but 

90
00:05:50,319 --> 00:05:54,288
the programming model being able to 
express lots of different tasks and you 

91
00:05:54,288 --> 00:05:58,572
know, have some sort of implementation to 
automatically turn that into parallel 

92
00:05:58,572 --> 00:06:03,740
job, turned out to be pretty powerful, 
right? 

93
00:06:03,740 --> 00:06:06,344
This was an attractive way to write 
parallel programs, again, because you 

94
00:06:06,344 --> 00:06:08,948
didn't actually have to worry about the 
parallel, all you had to do, write a 

95
00:06:08,948 --> 00:06:11,804
serial map function and a serial reduce 
function and the parallel had happens for 

96
00:06:11,804 --> 00:06:15,010
free. 
And so the evidence of this is not so 

97
00:06:15,010 --> 00:06:18,630
much about the system, as it is the 
program model, is that. 

98
00:06:18,630 --> 00:06:23,020
You see map reduce implementations appear 
in other contexts. 

99
00:06:23,020 --> 00:06:27,712
There's people who have implemented map 
reduce over uh,GPUs, There's people that 

100
00:06:27,712 --> 00:06:33,440
have implemented Map Reduce on multi-core 
machines in shared memory. 

101
00:06:33,440 --> 00:06:36,704
There are people who have implemented Map 
Reduce on high-performance computing 

102
00:06:36,704 --> 00:06:40,432
platforms on you know, mobile, groups of 
mobile phones and so on. 

103
00:06:40,432 --> 00:06:43,759
Okay. 
So, this goes back to one of the 

104
00:06:43,759 --> 00:06:47,520
motivations for this course. 
So I want to focus on abstractions where 

105
00:06:47,520 --> 00:06:50,608
possible, as opposed to tools. 
And so we're talking about Map Reduce in 

106
00:06:50,608 --> 00:06:53,695
the the programming model. 
But we'll spend a little bit less time in 

107
00:06:53,695 --> 00:06:58,150
the specific implementation, Hadoop. 
Although you will have a chance, an 

108
00:06:58,150 --> 00:07:00,475
optional assignment, to work with Hadoop 
directly. 

109
00:07:00,475 --> 00:07:06,090
Okay. 
So fine, so what is a data model of Map 

110
00:07:06,090 --> 00:07:08,948
Reduce? 
It's this bag of key value pairs, I mean 

111
00:07:08,948 --> 00:07:12,710
by bag is a set that might, might have 
duplicates in it. 

112
00:07:12,710 --> 00:07:15,457
Right, and so we've seen that before, the 
document ID with the value, that's a key 

113
00:07:15,457 --> 00:07:18,290
value pair. 
Sometimes on the input we'll be a little 

114
00:07:18,290 --> 00:07:22,049
sloppy and not worry about precisely what 
the key and the value is. 

115
00:07:22,049 --> 00:07:25,273
For example, if you're just given a 
record you can assume that its, say the, 

116
00:07:25,273 --> 00:07:29,782
the entire record is the key. 
Okay, or a document sometimes even we may 

117
00:07:29,782 --> 00:07:34,282
not have a explicit document id. 
But you can assume the URL or the file 

118
00:07:34,282 --> 00:07:38,475
name or something is the key. 
the output of the mapper, though, the, 

119
00:07:38,475 --> 00:07:41,611
the distinction between key and value 
gets really really important because 

120
00:07:41,611 --> 00:07:46,235
that's what controls the shuffle. 
As we've seen in that, in that example. 

121
00:07:46,235 --> 00:07:49,836
Okay. 
And so, both the, the data model here is 

122
00:07:49,836 --> 00:07:52,380
all about key value pairs. 
And the input is going to be a set of key 

123
00:07:52,380 --> 00:07:55,420
value pairs and the output is going to be 
a set of key value pairs. 

124
00:07:55,420 --> 00:07:59,102
And the point is that. 
The set of key value pairs can get 

125
00:07:59,102 --> 00:08:02,218
arbitrarily large, right. 
We're going to be able to process the set 

126
00:08:02,218 --> 00:08:06,571
no matter how big they get. 
There is kind of an implicit assumption 

127
00:08:06,571 --> 00:08:13,022
that the key and the value are small. 
And small here doesn't necessarily mean 

128
00:08:13,022 --> 00:08:16,310
very, very small, it just means that it 
needs to fit on one machine. 

129
00:08:16,310 --> 00:08:20,755
There's no support for a, for if value 
goes to be terrabytes. 

130
00:08:22,080 --> 00:08:25,300
It's it's not going to work. 
And so a document fits on one machine. 

131
00:08:25,300 --> 00:08:29,850
That's okay. 
You know a 

132
00:08:29,850 --> 00:08:34,640
an image that's on a machine, and so on, 
okay. 

133
00:08:34,640 --> 00:08:40,630
So, fine, so the map phase as we've said, 
you provide a map function. 

134
00:08:40,630 --> 00:08:44,780
The input is input key and input value. 
And the output is a bag of intermediate 

135
00:08:44,780 --> 00:08:47,572
keys and value. 
It doesn't have to just produce. 

136
00:08:47,572 --> 00:08:49,950
You're taking a single input, it produces 
single output. 

137
00:08:49,950 --> 00:08:52,420
It can produce a set of things. 
We saw this with the word count example. 

138
00:08:52,420 --> 00:08:54,760
A single document came in, but a set of 
things came out. 

139
00:08:54,760 --> 00:08:58,140
That's okay. 
Alright. 

140
00:08:58,140 --> 00:09:04,640
In the reduce phase what you're given is 
an intermediate, the intermediate key. 

141
00:09:04,640 --> 00:09:08,410
This would be the same intermediate key 
that was produced by the map phase. 

142
00:09:08,410 --> 00:09:12,585
Alright. 
One instance of the same intermediate key 

143
00:09:12,585 --> 00:09:15,469
produced by the math pairs. 
And then a bag of values that were 

144
00:09:15,469 --> 00:09:18,330
associated with that intermediate key. 
And they may have. 

145
00:09:18,330 --> 00:09:20,100
And the key. 
The important thing here is that they may 

146
00:09:20,100 --> 00:09:23,940
have come from any mapper whatsoever. 
The grouping into this bag of all, of 

147
00:09:23,940 --> 00:09:27,125
everything that shares this same 
intermediate key is handled automatically 

148
00:09:27,125 --> 00:09:29,424
by the system. 
Okay. 

149
00:09:30,880 --> 00:09:33,285
So the system will group all pairs with 
the same intermediate key and then pass 

150
00:09:33,285 --> 00:09:35,900
that bag of values to the reduced 
function. 

151
00:09:35,900 --> 00:09:39,116
And the implementation details of whether 
this is actually handed to you as a bag 

152
00:09:39,116 --> 00:09:42,236
or whether it's an iterator, if you're 
new, if you're familiar with that term, 

153
00:09:42,236 --> 00:09:45,644
that you could step over is implement-, 
implementation independent but it's not 

154
00:09:45,644 --> 00:09:50,958
important to think about. 
It's a collection of values. 

155
00:09:50,958 --> 00:09:55,780
Okay. 
Fine, so here it is all on one slide. 

156
00:09:56,990 --> 00:10:01,321
The map function takes a n key and an n 
value and produces a list, a bag of out 

157
00:10:01,321 --> 00:10:07,534
key and intermediate value pairs. 
And the reduce takes an out key and a 

158
00:10:07,534 --> 00:10:12,298
list of intermediate values and produces 
a list of out values. 

159
00:10:12,298 --> 00:10:16,420
I don't think I like this line, I think I 
prefer the other one. 

160
00:10:16,420 --> 00:10:21,388
The one thing I will mention is that the 
terms map and reduce, I tried to motivate 

161
00:10:21,388 --> 00:10:26,140
that in the last segment where you can 
think about converting a tif image to a 

162
00:10:26,140 --> 00:10:30,532
pmg image you can make a mapping, a 
function that maps Every TIF image into 

163
00:10:30,532 --> 00:10:36,720
some PNG image. 
And that's where the term came form. 

164
00:10:36,720 --> 00:10:39,236
And if look back this kind of came from 
the functional programming community that 

165
00:10:39,236 --> 00:10:41,728
use this term. 
It doesn't precisely mean the same thing 

166
00:10:41,728 --> 00:10:44,764
but it's inspired by that. 
Okay. 

167
00:10:44,764 --> 00:10:48,592
All right, so here's maybe the 
implementation for the example we have 

168
00:10:48,592 --> 00:10:52,420
and a lot of times what I like to do is 
sort of ask you to pause and stare at 

169
00:10:52,420 --> 00:10:58,620
this code for a little bit and think 
about what it does. 

170
00:10:58,620 --> 00:11:02,219
Here we've kind of gone through the 
example a lot so I'll reveal the secret 

171
00:11:02,219 --> 00:11:05,936
but it's still instructive to work 
through this for a moment yourself and in 

172
00:11:05,936 --> 00:11:09,771
fact, maybe I'll end this segment there, 
and you can stare at this and make sure 

173
00:11:09,771 --> 00:11:14,390
that you understand what it does. 

