1
00:00:00,598 --> 00:00:07,547
[MUSIC]. 

2
00:00:07,547 --> 00:00:11,036
Okay. 
So, let's talk about another example that 

3
00:00:11,036 --> 00:00:14,185
is sort of similar to this very simple 
word count example but, but has a slight 

4
00:00:14,185 --> 00:00:18,047
change. 
Okay, so I want to know/g, maybe the 

5
00:00:18,047 --> 00:00:23,118
makeup of a Corpus of documents instead 
of documents. 

6
00:00:23,118 --> 00:00:26,960
and we're trying to understand the 
characteristics of word length. 

7
00:00:26,960 --> 00:00:29,968
So now instead of a histogram on word 
usage, we're going to group things based 

8
00:00:29,968 --> 00:00:46,932
on the length of the word. 
You know, here we might group words into 

9
00:00:46,932 --> 00:00:50,700
big words, medium, small words, and tiny 
words. 

10
00:00:50,700 --> 00:00:53,750
Where the big words are everything that's 
ten plus letters. 

11
00:00:53,750 --> 00:00:57,526
And the medium words are in red here, and 
they're everything from five to nine 

12
00:00:57,526 --> 00:01:02,150
letters, and so on. 
And you could define your own sort of 

13
00:01:02,150 --> 00:01:08,256
bucketing scheme, or just not even try to 
bucket them and use the exact number, 

14
00:01:08,256 --> 00:01:12,614
okay. 
So, you know, what, what, what we're 

15
00:01:12,614 --> 00:01:16,254
basically showing here, so I'm hoping 
that some of you are, are already kind of 

16
00:01:16,254 --> 00:01:19,726
seeing how this is a, a essentially 
trivial variant of what we already just 

17
00:01:19,726 --> 00:01:23,310
did and this is some, one of the points I 
wanted to make is that you'll see these 

18
00:01:23,310 --> 00:01:26,950
patterns in designing map reduce albums 
come up again and again so I want you to 

19
00:01:26,950 --> 00:01:35,402
get a feel for how to do it. 
It won't be that every new problem looks, 

20
00:01:35,402 --> 00:01:38,522
looks different but the other, the other 
thing I want to talk about, in next few 

21
00:01:38,522 --> 00:01:42,989
slides, I think, is that. 
I-, it shows a little bit more detail in 

22
00:01:42,989 --> 00:01:48,532
how these things are broken up. 
So, for example, if this is a document 

23
00:01:48,532 --> 00:01:53,900
well, here, I guess the point I want to 
make is that before we sort of assume 

24
00:01:53,900 --> 00:01:59,620
that every document was a small but it's 
not impossible that you might have one 

25
00:01:59,620 --> 00:02:07,896
document that's, very very large. 
This typically wouldn't happen with a 

26
00:02:07,896 --> 00:02:12,975
document, and, for, various reasons. 
but imagine these weren't just documents 

27
00:02:12,975 --> 00:02:19,750
but these were, big data sets of, words. 
And so every document itself. 

28
00:02:19,750 --> 00:02:23,837
each an individual document may be so 
large that it can't be processed by a 

29
00:02:23,837 --> 00:02:28,280
single map function at a time. 
And so the question is are we s, are we 

30
00:02:28,280 --> 00:02:30,722
stuck here? 
Is this is the, is the map reduce 

31
00:02:30,722 --> 00:02:34,255
framework broken and is going to crash? 
And the answer is no. 

32
00:02:34,255 --> 00:02:39,506
What will happen is when you. 
Load this large document into the 

33
00:02:39,506 --> 00:02:43,095
system-supporting map reduce. 
And we'll talk a little bit more about 

34
00:02:43,095 --> 00:02:46,020
what that sys, what this lower-level 
system is, wha, what I mean by loading 

35
00:02:46,020 --> 00:02:50,660
the file system underneath map reduce, 
underneath implementations of map reduce. 

36
00:02:51,720 --> 00:02:55,941
When you do that loading, the file, the 
data set, the file, will automatically be 

37
00:02:55,941 --> 00:02:59,973
split into chunks and so we can pretend 
that this document was so large that it 

38
00:02:59,973 --> 00:03:06,024
needed to be split into chunks, okay? 
And so chunk one is this top part and 

39
00:03:06,024 --> 00:03:10,602
chunk two is this small part. 
Now, if the small document is underneath 

40
00:03:10,602 --> 00:03:14,472
the chunk size then it won't get split 
but a large document will. 

41
00:03:14,472 --> 00:03:17,238
Okay. 
And so, this could have happened with the 

42
00:03:17,238 --> 00:03:19,684
word count example, too. 
This is not something specific to word 

43
00:03:19,684 --> 00:03:25,371
link, obviously. 
But it's, it's another twist that we're, 

44
00:03:25,371 --> 00:03:32,768
that we're exploring, okay.. 
So fine. 

45
00:03:34,320 --> 00:03:40,000
So now, this top chunk gets assigned to 
map task 1. 

46
00:03:40,000 --> 00:03:44,900
And it produces, this little histogram, 
of of the counts of yellow words, red 

47
00:03:44,900 --> 00:03:49,924
words, blue words, and. 
Pink words, and we can imagine, you know, 

48
00:03:49,924 --> 00:03:53,188
you should think about how, how the code 
might look if you were to, if you had to 

49
00:03:53,188 --> 00:03:56,630
write this. 
Right, you would need to take the length, 

50
00:03:56,630 --> 00:03:59,600
and you know, iterate over all the words 
in the document just like we did before, 

51
00:03:59,600 --> 00:04:02,390
but now instead of emitting a key being 
the word itself, you would emit, you 

52
00:04:02,390 --> 00:04:06,640
would count its length And put in the 
case statement. 

53
00:04:06,640 --> 00:04:09,105
And figure out what color it is in this, 
in this notation. 

54
00:04:09,105 --> 00:04:14,200
And add that to account. 
Okay, [INAUDIBLE]. 

55
00:04:14,200 --> 00:04:16,240
So the output is these four key value 
pairs. 

56
00:04:19,200 --> 00:04:21,783
And in chunk two, we do the same thing, 
and produce a different set of four key 

57
00:04:21,783 --> 00:04:28,799
value pairs. 
Then in the shuffle step, all the yellow 

58
00:04:28,799 --> 00:04:36,010
key value pairs will be grouped together. 
And there's two of them. 

59
00:04:36,010 --> 00:04:37,900
And all the red ones are grouped together 
and so on. 

60
00:04:37,900 --> 00:04:41,595
And then in the reduce phase, you'll add 
these two numbers together. 

61
00:04:41,595 --> 00:04:45,355
To produce 37. 
Okay. 

62
00:04:45,355 --> 00:04:49,200
So the structure here is really identical 
to word count. 

63
00:04:49,200 --> 00:04:52,450
It's just basically a change to the map 
function. 

64
00:04:52,450 --> 00:04:56,046
And in fact, in this case you can 
actually literally use the exact same 

65
00:04:56,046 --> 00:04:59,762
reduce function. 
Alright, for every key, add up all the 

66
00:04:59,762 --> 00:05:02,218
contributions of it. 
You wouldn't need to change or reduce at 

67
00:05:02,218 --> 00:05:03,890
all. 
And this is something else you'll see, is 

68
00:05:03,890 --> 00:05:06,290
that, you know, certain reduced 
functions, certain functions are more 

69
00:05:06,290 --> 00:05:09,470
general than others, and you'll reuse 
them time and again. 

70
00:05:09,470 --> 00:05:12,053
For example, counting things and adding 
things up is pretty common in these 

71
00:05:12,053 --> 00:05:16,260
map-reduced, in, in these map-reduced 
algorithms, and so reduced. 

72
00:05:16,260 --> 00:05:20,440
General reduce functions that add things 
and count things come up time and again. 

73
00:05:20,440 --> 00:05:25,816
Okay, fine. 
Word count is the economical place to 

74
00:05:25,816 --> 00:05:31,492
start when thinking about map_reduce, 
word length is a very minor variation on 

75
00:05:31,492 --> 00:05:34,530
that. 
So. 

76
00:05:34,530 --> 00:05:37,150
let's think of, of another sort of minor 
variation. 

77
00:05:37,150 --> 00:05:40,530
And this one is arguably even simpler 
than, than word count. 

78
00:05:40,530 --> 00:05:43,080
So, here we want to build an inverted 
index. 

79
00:05:43,080 --> 00:05:46,986
And what an inverted index is, is when 
you have a corpus of documents, you can 

80
00:05:46,986 --> 00:05:51,830
presumably efficiently access any given 
document by it's name. 

81
00:05:51,830 --> 00:05:56,385
Alright, so if you want to look up a url 
on the web. 

82
00:05:56,385 --> 00:05:59,590
You, you can do so. 
But if you, but for a search engine, you 

83
00:05:59,590 --> 00:06:03,110
know very primitive search engine, you 
might want to look up, documents that 

84
00:06:03,110 --> 00:06:07,814
contain a particular word. 
And so building this, index to support 

85
00:06:07,814 --> 00:06:11,844
word look-up to provide documents, is 
called an inverted index, and its one of 

86
00:06:11,844 --> 00:06:18,160
the primitive steps in doing any kind of. 
Text retrieval kind of system. 

87
00:06:18,160 --> 00:06:20,360
Okay. 
So now, you know, imagine we just had 

88
00:06:20,360 --> 00:06:26,516
tweets instead of documents. 
and the input here is that the keys are 

89
00:06:26,516 --> 00:06:31,556
these tweet IDs that I've invented, and 
the value is the text of the tweet 

90
00:06:31,556 --> 00:06:39,240
itself. 
And so the desired output here is the 

91
00:06:39,240 --> 00:06:47,850
word along with a collection of Tweet 
ID's, okay? 

92
00:06:47,850 --> 00:06:51,170
So how do we do this? 
Well, the, the you know, again the code 

93
00:06:51,170 --> 00:06:57,015
is actually simpler than it was Before. 
Because now in the reduced phase, instead 

94
00:06:57,015 --> 00:07:01,870
of, well so we'll, we'll think about it 
for a moment. 

95
00:07:01,870 --> 00:07:07,950
The map phase instead of, instead of 
producing each word, pancake, in one for 

96
00:07:07,950 --> 00:07:12,870
an occurrence. 
We won't do that. 

97
00:07:12,870 --> 00:07:19,740
Instead we'll put out Pancake and the 
document ID itself. 

98
00:07:19,740 --> 00:07:21,070
Right? 
And all of these guys will be produced. 

99
00:07:21,070 --> 00:07:29,782
And so this tweet1, the map, the map task 
that processes this tweet1 will put out 

100
00:07:29,782 --> 00:07:42,396
pancake1, or tweet1. 
Love, tweet1 and so on. 

101
00:07:42,396 --> 00:07:45,748
Okay. 
Further, you know an optimization is, if 

102
00:07:45,748 --> 00:07:48,615
you see the word, I guess I should have 
had an example of this in here, but if 

103
00:07:48,615 --> 00:07:51,529
you see the word pancake twice in the 
same tweet, do you need to put it out, do 

104
00:07:51,529 --> 00:07:55,603
you need to emit the key-value pair 
twice? 

105
00:07:57,170 --> 00:08:00,574
Probably not build this index. 
Because all you're trying to record is of 

106
00:08:00,574 --> 00:08:06,220
the tweet contains the word pancake, not 
that it, not how many times it appears. 

107
00:08:06,220 --> 00:08:08,154
So fine. 
So these key value [INAUDIBLE] get 

108
00:08:08,154 --> 00:08:12,370
shuffled across the network, and now the 
reduce task, what does it do. 

109
00:08:12,370 --> 00:08:20,055
Well, it's going to get a key. 
Pancakes, which should have an s on it, I 

110
00:08:20,055 --> 00:08:24,226
guess. 
And it's going to have an iterator over 

111
00:08:24,226 --> 00:08:29,960
all the document IDs that contain 
pancakes, which in this case is, tweet1 

112
00:08:29,960 --> 00:08:36,940
and tweet. 
2. 

113
00:08:36,940 --> 00:08:41,586
So, do we need to do any processing on 
this key in group of values? 

114
00:08:41,586 --> 00:08:45,247
The answer in this case is no. 
The reduced function is complete is, is 

115
00:08:45,247 --> 00:08:48,270
nonexistent, there's nothing to do, 
alright? 

116
00:08:48,270 --> 00:08:51,210
From what the group is exactly what you 
want in this case. 

117
00:08:51,210 --> 00:08:53,260
And this pattern actually shown, shows 
up. 

118
00:08:53,260 --> 00:08:55,468
Somewhat often as well. 
Where you do want the map and you do want 

119
00:08:55,468 --> 00:08:57,925
the shuffle phase to do the grouping but 
all you wanted to do is to perform the 

120
00:08:57,925 --> 00:09:01,490
grouping, you didn't actually care about 
the reduce function. 

121
00:09:01,490 --> 00:09:03,625
So you're not counting these things, 
you're not adding them up in any way, 

122
00:09:03,625 --> 00:09:06,710
you're not doing and processing on the 
tweets, you just want to omit that. 

123
00:09:06,710 --> 00:09:11,198
And that's perfectly fine because this 
group of values is perfectly serviceable 

124
00:09:11,198 --> 00:09:15,150
as a value itself. 
OK. 

125
00:09:15,150 --> 00:09:19,162
So if you are used to, say, relational 
databases, you know, nested structures 

126
00:09:19,162 --> 00:09:22,761
collections of values within a single 
cell in a table, in a single row, are 

127
00:09:22,761 --> 00:09:27,728
generally disallowed. 
Right, and this is actually first normal 

128
00:09:27,728 --> 00:09:32,740
form, if you're familiar with that. 
but here we don't care. 

129
00:09:32,740 --> 00:09:35,824
If you need, it's, it's. 
For any key in any value, a key can have 

130
00:09:35,824 --> 00:09:39,420
sub structure and a value can have sub 
structure. 

131
00:09:39,420 --> 00:09:41,965
Okay fine so that's how to build an 
inverted index in map produce. 

132
00:09:41,965 --> 00:09:46,861
let me stop there next time we'll talk 
about this relational joint example, so 

133
00:09:46,861 --> 00:09:51,469
this will be how to implement a Join from 
a relational database as a map reduce 

134
00:09:51,469 --> 00:09:54,570
program. 

