1
00:00:00,135 --> 00:00:06,257
[MUSIC]. 

2
00:00:06,257 --> 00:00:11,508
So let's look at an example that doesn't 
involve processing a corpus of documents. 

3
00:00:11,508 --> 00:00:15,278
So, let's think about how to implement a 
relational the join operation that we 

4
00:00:15,278 --> 00:00:18,791
learn from the relation algebra in Map 
Reduce. 

5
00:00:18,791 --> 00:00:23,731
Okay, so here recall that you're given 
two relations, and a relation is a set of 

6
00:00:23,731 --> 00:00:28,930
tuples, all right? 
And you're trying to find every record in 

7
00:00:28,930 --> 00:00:34,900
one relation that corresponds to a record 
in the other relation, right? 

8
00:00:34,900 --> 00:00:48,620
And so here, we're going to join on SSN 
equal to EmpSSN. 

9
00:00:48,620 --> 00:00:52,710
And actually, it wasn't quite correct for 
me to leave this unspecified. 

10
00:00:52,710 --> 00:00:56,660
There's the notion of a natural join, 
that would match up with the field names. 

11
00:00:56,660 --> 00:00:58,670
they had to remade was mash but here they 
don't. 

12
00:00:58,670 --> 00:01:00,700
So, I need be explicit on what I'm 
joining on. 

13
00:01:00,700 --> 00:01:05,022
Okay, and so the output here is, these 
three records. 

14
00:01:05,022 --> 00:01:09,548
This record with joins with these two, 
and this record joins with this one, 

15
00:01:09,548 --> 00:01:14,445
fine. 
Now, image that both of these relations 

16
00:01:14,445 --> 00:01:19,467
are huge. 
Okay, so the map phase here is going to 

17
00:01:19,467 --> 00:01:27,906
process every two [INAUDIBLE] in general. 
And we have a problem right off the bat 

18
00:01:27,906 --> 00:01:34,055
that join is a binary operation. 
It has a two-input relations a left 

19
00:01:34,055 --> 00:01:38,477
relation and the right relation and 
you're trying to corresponding two equals 

20
00:01:38,477 --> 00:01:43,410
one and match the other. 
Is a unary operation that processes a 

21
00:01:43,410 --> 00:01:47,610
single data set. 
So how can you, you know, at first, to 

22
00:01:47,610 --> 00:01:54,310
first approximation you can't express 
join in map reduce period. 

23
00:01:54,310 --> 00:01:56,920
But that's okay, there's a, there's a, 
there's a bit of a trick here. 

24
00:01:56,920 --> 00:02:00,376
The trick is look you know, if it's okay 
that you know imagine the data set here 

25
00:02:00,376 --> 00:02:04,410
is just a big jumble of tuples. 
It's doesn't matter what table they came 

26
00:02:04,410 --> 00:02:10,219
from we just lump them all together. 
Into a single data set called tuples. 

27
00:02:10,219 --> 00:02:13,215
Okay, and that will be what we process 
with map reduce. 

28
00:02:13,215 --> 00:02:16,610
Okay, and so here I'm asking the question 
you know, what is this for? 

29
00:02:16,610 --> 00:02:22,506
Well, this is a label that we've attached 
to every tuple, so that we can know where 

30
00:02:22,506 --> 00:02:29,813
it came from and that's used later. 
Okay, now I 'm not being specific about 

31
00:02:29,813 --> 00:02:34,730
how, where you get this label or how you 
do it in MapReduce. 

32
00:02:34,730 --> 00:02:37,610
But I want you to think logically that 
it's necessary and in practice it's not 

33
00:02:37,610 --> 00:02:41,036
that difficult. 
Probably because you know, for example, 

34
00:02:41,036 --> 00:02:44,694
if you're processing, you have 
distributed files from this directory 

35
00:02:44,694 --> 00:02:48,588
representing all the employee chunks And 
you've got distributed files in this 

36
00:02:48,588 --> 00:02:53,730
directory representing the assigned 
departments chunks. 

37
00:02:53,730 --> 00:02:57,230
You can look at the file name to tell 
what table it is. 

38
00:02:57,230 --> 00:03:00,158
And so in the map function that you 
write, you can determine that aha, this 

39
00:03:00,158 --> 00:03:03,182
has file name, you know, this has the 
employee relation name in it's, in it's 

40
00:03:03,182 --> 00:03:08,380
file path and I'm going to go ahead and 
attach that as part of the record, okay? 

41
00:03:08,380 --> 00:03:10,510
And so we could write pseudo code for 
this. 

42
00:03:10,510 --> 00:03:17,178
And in fact, you might, as part of the an 
upcoming assignment, alright? 

43
00:03:17,178 --> 00:03:23,750
[INAUDIBLE], so what's the map phrase 
look like of this relational join? 

44
00:03:23,750 --> 00:03:31,110
Well, for every record on the input 
you're going to produce a key value pair. 

45
00:03:31,110 --> 00:03:33,126
We know we have to produce key value 
pairs, because that's how MapReduce 

46
00:03:33,126 --> 00:03:35,082
works. 
So what's the key going to be? 

47
00:03:35,082 --> 00:03:38,778
Well, I'm going to give you, the key is 
going to be the join attribute, the 

48
00:03:38,778 --> 00:03:43,068
attribute that you're joining on, and the 
value is going to be sort of everything 

49
00:03:43,068 --> 00:03:47,439
else, all right? 
The rest of the tuple, and we can 

50
00:03:47,439 --> 00:03:50,895
actually get away with removing this guy 
to save some space but typically 

51
00:03:50,895 --> 00:03:58,730
typically you wouldn't bother. 
Okay, so fine. 

52
00:03:58,730 --> 00:04:03,670
So, given a tuple that looks like this 
with, with three fields, we produce a key 

53
00:04:03,670 --> 00:04:08,350
of seven seven seven seven seven seven so 
on. 

54
00:04:08,350 --> 00:04:11,220
And the value is the entire tuple, and so 
on. 

55
00:04:11,220 --> 00:04:16,236
And again notice that both the tuple from 
both relations are all lumped into the 

56
00:04:16,236 --> 00:04:19,600
same input. 
Alright, I may be belaboring this but I 

57
00:04:19,600 --> 00:04:22,852
want to make sure that it's clear. 
Okay, so so far we've done two tricks. 

58
00:04:22,852 --> 00:04:26,362
One is lump everything together and two 
is produce a key value pair with a key as 

59
00:04:26,362 --> 00:04:30,502
the joint attribute. 
Now, what happens in the magic shuffle 

60
00:04:30,502 --> 00:04:33,098
phase? 
Well Everything with the same key gets 

61
00:04:33,098 --> 00:04:37,870
lumped together, as we've talked about. 
And so now you get a reducer invocation 

62
00:04:37,870 --> 00:04:42,930
that has you know, one reducer invocation 
for every unique key. 

63
00:04:42,930 --> 00:04:46,895
So in this case, it would be one for 999 
so on, one for 777 so on and the list of 

64
00:04:46,895 --> 00:04:51,607
values. 
Associated with that key will be all the 

65
00:04:51,607 --> 00:04:56,702
tuple that share the same join key. 
Now it doesn't matter what relation they 

66
00:04:56,702 --> 00:05:01,106
came from they'll all be in this list. 
So in this case we get one tuple from the 

67
00:05:01,106 --> 00:05:06,090
employee relation, and we get one tuple 
from the department relation. 

68
00:05:06,090 --> 00:05:09,338
And here we get one from the employee 
relation, and two tuples from the 

69
00:05:09,338 --> 00:05:13,044
department relation. 
Alright, so now in the context of a 

70
00:05:13,044 --> 00:05:17,850
single computer, we have everything we 
need to compute a single join. 

71
00:05:17,850 --> 00:05:23,691
And further, each of these reduce 
invocations can be done on a separate 

72
00:05:23,691 --> 00:05:30,780
machine, and that's how we scale. 
So now this reduce function that the 

73
00:05:30,780 --> 00:05:36,092
program would have to write if you are 
implementing a join phase would have to 

74
00:05:36,092 --> 00:05:44,520
take this key, and takes all these tuples 
and produce the joined tuple, right? 

75
00:05:44,520 --> 00:05:56,232
Where these two attributes came from 
employee and these two attributes came 

76
00:05:56,232 --> 00:06:10,340
from department. 
And same thing here employee, department. 

77
00:06:10,340 --> 00:06:14,309
Now, if you want to think about what kind 
of operation you need to implement in 

78
00:06:14,309 --> 00:06:18,713
this reduce function. 
Well, you've got a set of tuples from one 

79
00:06:18,713 --> 00:06:22,808
relation and you need to associate it 
with every possible tuple from the other 

80
00:06:22,808 --> 00:06:28,310
relation, right because we need to know 
they all join. 

81
00:06:28,310 --> 00:06:31,855
They're all by definition, they all join, 
they have the same key, they all match. 

82
00:06:31,855 --> 00:06:35,890
So if you have every member of a set 
paired with. 

83
00:06:35,890 --> 00:06:40,562
Every possible member of another set, if 
you remember that's a cross product 

84
00:06:40,562 --> 00:06:44,094
operation. 
And so here again we see relation 

85
00:06:44,094 --> 00:06:47,866
[INAUDIBLE] popping up sort in a, in a 
different context. 

86
00:06:47,866 --> 00:06:50,701
Okay, so I don't want to confuse you to 
much with the overall operation we're 

87
00:06:50,701 --> 00:06:53,354
trying to do is implement a parallel 
join. 

88
00:06:53,354 --> 00:06:57,460
It happens to be that locally that right 
here inside one reduce function. 

89
00:06:57,460 --> 00:07:00,754
We recognize, aha, that's another 
relational algebra operator, a cross 

90
00:07:00,754 --> 00:07:03,360
product. 
But this is more for abstraction purposes 

91
00:07:03,360 --> 00:07:06,208
than for algorithm purposes, I just 
want to point that out. 

92
00:07:06,208 --> 00:07:08,704
All right, so again, you put on your 
relation algebra colored glasses and the 

93
00:07:08,704 --> 00:07:12,150
whole world, you know, these operators 
start to pop up, pop up everywhere. 

94
00:07:12,150 --> 00:07:17,480
All right, so fine, so let's do this one 
more time just to make sure it's clear. 

95
00:07:17,480 --> 00:07:22,572
So, I'm giving, giving you two relations, 
order and line item and they have these 

96
00:07:22,572 --> 00:07:27,132
fields; an order ID, an account and a 
date and here there's an order ID, an 

97
00:07:27,132 --> 00:07:34,159
item ID and a quantity and we're going to 
join on order ID. 

98
00:07:34,159 --> 00:07:39,640
So what's the map phase look like? 
Well, once again, the key will be the 

99
00:07:39,640 --> 00:07:45,844
join key in this case order ID and the 
value will be this pair. 

100
00:07:45,844 --> 00:07:48,284
The relation name, along with the 
original tuple. 

101
00:07:48,284 --> 00:07:51,322
And before, we just sort of lumped these 
together in one tuple and it doesn't 

102
00:07:51,322 --> 00:07:54,438
matter very much. 
Here I've kind of structured it slightly 

103
00:07:54,438 --> 00:07:57,570
differently, but the point is, is all the 
information is here. 

104
00:07:57,570 --> 00:08:01,840
The key must be the join key and value is 
the tuple as well as some sort of 

105
00:08:01,840 --> 00:08:06,290
indicator for of what relation it came 
from. 

106
00:08:06,290 --> 00:08:08,718
So, maybe I'll ask real quick. 
Why do I need that indicator of the 

107
00:08:08,718 --> 00:08:13,930
relation? 
Well, in the reduced phase, I wouldn't be 

108
00:08:13,930 --> 00:08:16,450
able to produce the joint tuples properly 
if we didn't know which, which two 

109
00:08:16,450 --> 00:08:20,075
peoples went with which relation. 
If I just had a big bundle of tuples and 

110
00:08:20,075 --> 00:08:23,652
I couldn't really figure it out very 
easily, that would be a problem. 

111
00:08:23,652 --> 00:08:28,074
Theoretically, you might be able to avoid 
tagging explicitly with employee and 

112
00:08:28,074 --> 00:08:32,562
department because you could deduce that 
well, the employee table is the one that 

113
00:08:32,562 --> 00:08:36,852
has a string first followed by the join, 
followed by the join key and the 

114
00:08:36,852 --> 00:08:44,150
department relation has the join key 
followed by some other string. 

115
00:08:44,150 --> 00:08:46,392
But it's a little dangerous to rely on 
that kind of information, because you 

116
00:08:46,392 --> 00:08:49,840
could, you could be joining two tables 
that have exactly the same schema. 

117
00:08:49,840 --> 00:08:52,450
And there wouldn't be any real obvious 
way to figure it out. 

118
00:08:52,450 --> 00:08:55,070
So an explicit tag is a little bit safer 
here, alright, fine. 

119
00:08:55,070 --> 00:08:59,541
So that's the map phase. 
Process all the, we, we lump all the 

120
00:08:59,541 --> 00:09:02,966
records together. 
And produce key value pairs, where the 

121
00:09:02,966 --> 00:09:08,250
key is the join key from the, from the 
corresponding relation. 

122
00:09:08,250 --> 00:09:11,194
Join attribute from the corresponding 
relation, and the value is the entire 

123
00:09:11,194 --> 00:09:14,133
tuple tag with this, with the relation 
name, okay? 

124
00:09:14,133 --> 00:09:17,342
And there it is [LAUGH], right there on 
the slide. 

125
00:09:17,342 --> 00:09:23,239
Okay, so what's the reducer look like? 
Well Now we've got all the tuples that 

126
00:09:23,239 --> 00:09:31,080
have this share the same key together and 
it will produce these joined tuples. 

127
00:09:31,080 --> 00:09:35,916
It will join this order tuple with both 
of these line item tuples to produce 

128
00:09:35,916 --> 00:09:43,684
these two joined tuples. 
Okay, fine. 

129
00:09:43,684 --> 00:09:47,518
So, now lets go and do a different 
example, and so here may be we're 

130
00:09:47,518 --> 00:09:52,370
going to get started with analyzing the 
social networks. 

131
00:09:52,370 --> 00:09:57,149
So here we have a graph where every edge 
represents A, let's say a friend 

132
00:09:57,149 --> 00:10:00,792
relationship. 
Or if you're thinking about Twitter, you 

133
00:10:00,792 --> 00:10:05,132
can always be a followers relationship. 
So, the input here is a set of edges with 

134
00:10:05,132 --> 00:10:10,478
the semantics of Jim is friends with Sue, 
or follows Sue, and Sue is friends with 

135
00:10:10,478 --> 00:10:17,064
Jim, or follows Jim, and so on. 
And so, one point I'm making here is 

136
00:10:17,064 --> 00:10:21,420
you'll notice that if, if Lin follow, if 
Jim, if Lin points to Joe then Joe points 

137
00:10:21,420 --> 00:10:27,510
to Lin, and so this is there's a 
symmetric relationship here, okay? 

138
00:10:27,510 --> 00:10:30,262
And I've done that for simplicity sake to 
sort of avoid the confusion that can 

139
00:10:30,262 --> 00:10:33,920
result from thinking about undirected 
graphs versus directed graphs. 

140
00:10:36,570 --> 00:10:38,130
So, any time you see an edge going one 
way you'll see the other edge coming 

141
00:10:38,130 --> 00:10:40,825
back. 
Now, what we want to do here is, well 

142
00:10:40,825 --> 00:10:46,360
before, before I say what we're going to 
do, is make sure task is clear. 

143
00:10:46,360 --> 00:10:48,260
We're going to count the friends. 
We want something very simple. 

144
00:10:48,260 --> 00:10:50,520
We just want to say, how many friends 
does Jim have? 

145
00:10:50,520 --> 00:10:52,900
And how many friends does Sue have, and 
so on. 

146
00:10:52,900 --> 00:10:56,440
And so the desired output here is Jim 
with three, because Jim has Sue as a 

147
00:10:56,440 --> 00:11:01,238
friend and Jim has Kai as a friend and 
Jim has Lin as a friend. 

148
00:11:01,238 --> 00:11:06,728
And we want to have Lin equals two, or 
Lin has two friends, becaue Joe and Jim 

149
00:11:06,728 --> 00:11:15,235
and Kai is just one, right here. 
Jim and Joe is just one, which is right 

150
00:11:15,235 --> 00:11:19,655
here, Lin, okay? 
So, this should already look like 

151
00:11:19,655 --> 00:11:23,000
something we've already done. 
Like before so it happens to be social 

152
00:11:23,000 --> 00:11:27,355
network analysis and the records happen 
to be these you know, pairs of people but 

153
00:11:27,355 --> 00:11:31,515
the algorithm you're going to use should 
start to stand out to you okay and maybe 

154
00:11:31,515 --> 00:11:37,100
you'll see why in a second if it's not 
clear. 

155
00:11:37,100 --> 00:11:41,325
So, in the math phase how are we going to 
do this well for every you know, friend 

156
00:11:41,325 --> 00:11:45,290
on the left hand side produce a key value 
pair with the key is the name of the 

157
00:11:45,290 --> 00:11:49,515
friend or the name of the person and the 
value is just the number one indicating 

158
00:11:49,515 --> 00:11:56,630
that there is exactly one friend 
associated with that. 

159
00:11:56,630 --> 00:11:59,192
So exactly there is, there is that we've, 
that we've encountered 1 friend 

160
00:11:59,192 --> 00:12:02,730
associated with. 
Say Jim in this case alright? 

161
00:12:02,730 --> 00:12:06,492
So this record gets turned into Jim and 1 
and this record gets turned into Sue and 

162
00:12:06,492 --> 00:12:10,617
1 and this record gets turned into Lin 
and 1 and so on. 

163
00:12:10,617 --> 00:12:15,102
Okay and then you know, through the magic 
of map produce the shuffle phase produces 

164
00:12:15,102 --> 00:12:18,937
this new media result where we have Jim 
the key associated with a list of 

165
00:12:18,937 --> 00:12:26,165
occurrences 1, 1 and 1. 
And Lin associated with two occurrences, 

166
00:12:26,165 --> 00:12:28,340
and so on. 
And in the reduced phase, just sim-, 

167
00:12:28,340 --> 00:12:32,200
simply adds up all of these occurrences. 
And so, what does this remind you of? 

168
00:12:32,200 --> 00:12:34,550
Well, it's a lot like the word count 
example, right? 

169
00:12:34,550 --> 00:12:38,930
Instead of documents, producing words and 
counts. 

170
00:12:38,930 --> 00:12:43,389
It's even simpler. 
It's for every record it just produces a 

171
00:12:43,389 --> 00:12:48,426
single count for that for that person 
that the record represents, the person on 

172
00:12:48,426 --> 00:12:54,926
the left hand side. 
And after that the reducer is literally 

173
00:12:54,926 --> 00:12:59,265
exactly the same, right? 
It'll, it'll add up all the occurrences 

174
00:12:59,265 --> 00:13:04,153
and produce a single number. 
Okay, so in the next segment, we'll talk 

175
00:13:04,153 --> 00:13:10,480
about something a little more, a little 
trickier. 

176
00:13:10,480 --> 00:13:17,814
which is implementing matrix 
multiplication in MapReduce, alright? 

