1
00:00:00,730 --> 00:00:03,187
Okay, so let's look at the CoGroup 
command, so we've seen the group command, 

2
00:00:03,187 --> 00:00:05,722
which works in a single data set and the 
co-group command does the same thing on 

3
00:00:05,722 --> 00:00:09,385
multiple data sets. 
And we seen this mechanism before we 

4
00:00:09,385 --> 00:00:13,147
talked about map reduce but here in pig 
we make this operation explicit and give 

5
00:00:13,147 --> 00:00:17,880
it a specific command and we'll try to 
see why that is. 

6
00:00:17,880 --> 00:00:22,632
Okay, so the syntax looks like this 
[SOUND] you say CoGroup a dataset by some 

7
00:00:22,632 --> 00:00:28,037
field reference or multiple fields. 
And here we're referring to a field by 

8
00:00:28,037 --> 00:00:29,708
name. 
And here we're referring to a field by 

9
00:00:29,708 --> 00:00:33,359
position as we've seen before. 
And so, if this is A, a bag of tuples and 

10
00:00:33,359 --> 00:00:38,910
this is B, a bag of tuples then the 
CoGroup of them looks like this. 

11
00:00:38,910 --> 00:00:47,688
So, now we have a Group field and we have 
a dataset A a field named A and we have a 

12
00:00:47,688 --> 00:00:53,726
field named B. 
Actually we still have the same grouping 

13
00:00:53,726 --> 00:00:55,580
key but we have two groups associated 
with it. 

14
00:00:55,580 --> 00:01:00,290
One from one data set and one from the 
other, okay. 

15
00:01:00,290 --> 00:01:03,299
So, a couple things to point out. 
One is, if, if there are no tuples from 

16
00:01:03,299 --> 00:01:07,450
one of the data sets and it just get's 
the empty group. 

17
00:01:07,450 --> 00:01:11,433
here you see [INAUDIBLE] from B. 
And then the other thing to point out is 

18
00:01:11,433 --> 00:01:16,169
this is a little different then what 
happens in MapReduce directly where you 

19
00:01:16,169 --> 00:01:21,275
know, you, we,what we showed, with, with 
the, JOIN operator without produce is 

20
00:01:21,275 --> 00:01:29,265
that you could, use the JOIN key to sort 
of CoGroup two relations. 

21
00:01:29,265 --> 00:01:34,225
And then all of these tuples would all 
sort of appear in the same group in the 

22
00:01:34,225 --> 00:01:39,185
reducer but here we sort of make it 
explicit to make it two separate groups, 

23
00:01:39,185 --> 00:01:44,536
okay? 
Fine, so I want to come back to co-groups 

24
00:01:44,536 --> 00:01:49,144
in a second but first let's talk about 
another operator, which is just called 

25
00:01:49,144 --> 00:01:54,241
JOIN and so it does exactly what you'd 
expect. 

26
00:01:54,241 --> 00:01:58,643
So JOIN, you say JOIN A by some set of 
field references and B by some other set 

27
00:01:58,643 --> 00:02:04,321
of field references, okay. 
And given A and given B the result will 

28
00:02:04,321 --> 00:02:09,456
be what we've talked about in the past, 
which is you know, you look for the first 

29
00:02:09,456 --> 00:02:17,120
position of A, $0 and find all the 
corresponding first positions of B. 

30
00:02:17,120 --> 00:02:21,593
So, here's one; 1,1. 
So we have 1, 2 from A and 1, and 3, 1, 3 

31
00:02:21,593 --> 00:02:25,668
from B. 
Okay, so there's nothing stopping you 

32
00:02:25,668 --> 00:02:31,574
from having multiple datasets out here. 
As many as you want and they can all be 

33
00:02:31,574 --> 00:02:36,743
sort of processed together. 
And so you think about what's going on 

34
00:02:36,743 --> 00:02:41,701
here is this one MapReduce job underneath 
the sheets where in the map phase every 

35
00:02:41,701 --> 00:02:49,521
tuple from a variety of data sets is all 
being associated with the JOIN key. 

36
00:02:49,521 --> 00:02:55,191
represented by this field reference in 
the syntax, and then shuffled across the 

37
00:02:55,191 --> 00:03:01,388
network to arrive at the reduce a fade. 
Okay, so here's what it looks like well 

38
00:03:01,388 --> 00:03:05,046
if we have multiple relations where 
represent by color here green blue and 

39
00:03:05,046 --> 00:03:08,704
red well they can all be processed by the 
map phase and associated with keys 

40
00:03:08,704 --> 00:03:14,660
shuffle them across the network to 
produce the Join tuples here. 

41
00:03:14,660 --> 00:03:17,433
Right so there's not there's nothing 
fundamentally binary about this 

42
00:03:17,433 --> 00:03:19,210
operation. 
Operator. 

43
00:03:19,210 --> 00:03:20,640
Right? 
It's just processing tuples associating 

44
00:03:20,640 --> 00:03:22,884
them with their, with their respective 
Join keys and shuffling them all across 

45
00:03:22,884 --> 00:03:26,114
the network. 
What can go wrong with this basic 

46
00:03:26,114 --> 00:03:29,390
mechanism of associating tuples with a 
Join key, shuffling them across the 

47
00:03:29,390 --> 00:03:33,082
network and then producing the, you know, 
completing the Join on the reduce side, 

48
00:03:33,082 --> 00:03:37,210
which is what we did the assignment as 
well. 

49
00:03:37,210 --> 00:03:40,927
Well, so one example might be if one 
table is very very large and another 

50
00:03:40,927 --> 00:03:46,108
table is very very small. 
There's an opportunity to do something 

51
00:03:46,108 --> 00:03:51,628
much much faster, which is, replicate the 
small table across all partitions of the 

52
00:03:51,628 --> 00:03:58,928
large table which allows you to do all 
the work in the map phase, okay. 

53
00:03:58,928 --> 00:04:04,990
The second sort of special case here is a 
Skewed join. 

54
00:04:04,990 --> 00:04:08,329
And what I mean by that is, if there are 
many values in one table that join with 

55
00:04:08,329 --> 00:04:11,774
many, many, many values In the second 
table then you'll end up with one reducer 

56
00:04:11,774 --> 00:04:17,258
doing all the, most of the work. 
And the effects of parallelism gets sort 

57
00:04:17,258 --> 00:04:20,978
of washed out. 
And the third special case algorithm is 

58
00:04:20,978 --> 00:04:24,838
to do a Merge join. 
So this takes advantage of the fact that 

59
00:04:24,838 --> 00:04:28,198
you may have already grouped two 
relations in the same way such that the, 

60
00:04:28,198 --> 00:04:32,006
you know, for sure that on a single 
machine. 

61
00:04:32,006 --> 00:04:35,534
I've got all the tuples I need from one 
relation and all the tuples I need from 

62
00:04:35,534 --> 00:04:38,049
the other relation. 
Okay. 

63
00:04:38,049 --> 00:04:40,822
So let me see if I can make this more 
clear in, pictures in the next few 

64
00:04:40,822 --> 00:04:42,830
slides. 
Okay. 

65
00:04:42,830 --> 00:04:46,340
So for replicated joined, the situation 
is we have one large table broken into 

66
00:04:46,340 --> 00:04:50,790
pieces and one much smaller table. 
And so what we could do is just shuffle 

67
00:04:50,790 --> 00:04:54,246
all these two bulls across the network 
and do the Join on the reduced side as we 

68
00:04:54,246 --> 00:04:58,334
do normally. 
But there's an opportunity here that says 

69
00:04:58,334 --> 00:05:03,015
well look if this thing is small enough 
to fit in memory on a single machine. 

70
00:05:03,015 --> 00:05:06,540
Why don't we just copy it? 
Send it out there to every, you know, 

71
00:05:06,540 --> 00:05:12,621
every map function that wakes up will go 
pull it across the network directly. 

72
00:05:12,621 --> 00:05:15,246
Okay. 
Then we have all the information we need 

73
00:05:15,246 --> 00:05:19,683
right here in the map side. 
Every one of these tuples can be joined 

74
00:05:19,683 --> 00:05:23,010
with corresponding tuples in this 
partition. 

75
00:05:23,010 --> 00:05:25,341
Every one of these tuples could be joined 
with corresponding tuples in this 

76
00:05:25,341 --> 00:05:27,990
partition and so on. 
And so the end of the map phase you end 

77
00:05:27,990 --> 00:05:31,478
up with the right answer. 
All the joined tuples for you know this, 

78
00:05:31,478 --> 00:05:34,582
this sort of blue A and this red, red B. 
Okay. 

79
00:05:34,582 --> 00:05:37,685
And so why is this cheaper? 
Well we didn't have to shuffle everything 

80
00:05:37,685 --> 00:05:40,220
across the network. 
We did it all on the map phase, and it, 

81
00:05:40,220 --> 00:05:42,730
it makes a huge difference. 
Okay. 

82
00:05:42,730 --> 00:05:46,034
So you might see this called a Broadcast 
join, as well. 

83
00:05:46,034 --> 00:05:50,590
So the small relation must fit in memory, 
and the idea is that each mapper in, in 

84
00:05:50,590 --> 00:05:56,772
the map function pulls a copy of the 
small relation directly out of HGFS. 

85
00:05:56,772 --> 00:05:59,130
Okay. 
All right. 

86
00:05:59,130 --> 00:06:04,990
So for the Skewed join the situation is, 
is as usual. 

87
00:06:04,990 --> 00:06:07,670
You got two relations, this sort of blue 
one and this red one. 

88
00:06:07,670 --> 00:06:12,448
And the map phase associates each tuple 
with its JOIN key. 

89
00:06:12,448 --> 00:06:17,453
But the problem is that most of the data 
ends up on single reducer and the reason 

90
00:06:17,453 --> 00:06:24,626
is because maybe most of the data here is 
associated with a single Join key. 

91
00:06:24,626 --> 00:06:26,752
Okay. 
So, for example, if you're joined on 

92
00:06:26,752 --> 00:06:30,168
order ID and line item, you know, you 
just try and associate all line items 

93
00:06:30,168 --> 00:06:33,920
with their corresponding order, it could 
be that one order had millions of parts 

94
00:06:33,920 --> 00:06:39,081
and all the other orders had five parts, 
or something. 

95
00:06:39,081 --> 00:06:41,645
Okay. 
So if there's if there's significant 

96
00:06:41,645 --> 00:06:44,921
refraction in the overall data set is 
associated with a single join key then 

97
00:06:44,921 --> 00:06:48,041
this reducer will be doing all the work 
and there won't be much benefit to 

98
00:06:48,041 --> 00:06:59,662
parallelism okay. 
So this would work but there's a problem. 

99
00:06:59,662 --> 00:07:10,590
So, this is from former student here at 
UDUB who did some work on this problem. 

100
00:07:10,590 --> 00:07:15,630
And this plot shows time in seconds on 
this x axis. 

101
00:07:15,630 --> 00:07:20,970
And this is just a list of all the tasks 
and so you see that the reduced task 

102
00:07:20,970 --> 00:07:28,932
here, they are, they can't start until 
all the map tests are finished. 

103
00:07:28,932 --> 00:07:31,328
Okay. 
And the map tests mostly finished quite 

104
00:07:31,328 --> 00:07:34,791
quickly. 
Sort of, maybe, 20 seconds. 

105
00:07:34,791 --> 00:07:38,220
But one or two of these map tasks take a 
very long time. 

106
00:07:38,220 --> 00:07:40,680
You know, sort of on the order of 270 
seconds or so. 

107
00:07:40,680 --> 00:07:45,260
Okay, so all this space in here is sort 
of wasted work. 

108
00:07:45,260 --> 00:07:48,859
There's a couple of problems here, one is 
fundamentally map reduce in may 

109
00:07:48,859 --> 00:07:52,694
applications you logically could start 
doing some work early based on the map 

110
00:07:52,694 --> 00:07:59,038
output that has already been finished. 
And that's just not the way map producers 

111
00:07:59,038 --> 00:08:01,833
is designed. 
y, you can't take, you, it's not designed 

112
00:08:01,833 --> 00:08:04,143
to be able to take advantage of those 
applications because you can't guarantee 

113
00:08:04,143 --> 00:08:07,040
that it's safe to do so. 
So, for example, if you're adding up 

114
00:08:07,040 --> 00:08:09,780
numbers, you could start adding up in, in 
the reduced stage. 

115
00:08:09,780 --> 00:08:12,796
You could start adding up numbers early. 
But if you're doing something more 

116
00:08:12,796 --> 00:08:14,972
complicated you may actually need to wait 
for all the results to be present in 

117
00:08:14,972 --> 00:08:18,681
order to get the correct result. 
And so, since they can't guarantee that 

118
00:08:18,681 --> 00:08:26,050
it's safe th, they make you wait. 
Okay, so fine. 

119
00:08:26,050 --> 00:08:29,040
So, this skew problem ends up sort of 
killing parallelism in terms the job that 

120
00:08:29,040 --> 00:08:32,554
you know could take sort of 50 seconds. 
And the one that takes 300 to 350 

121
00:08:32,554 --> 00:08:35,115
seconds. 
And is not even significantly longer than 

122
00:08:35,115 --> 00:08:38,090
doing this sequentially. 
We should of added all these pieces up. 

123
00:08:38,090 --> 00:08:40,650
Well, I shouldn't say that. 
This is so worse than, in sequential. 

124
00:08:40,650 --> 00:08:44,119
But you cer, you certainly lose a lot of 
your bandwidth in parallelism. 

125
00:08:44,119 --> 00:08:45,899
Okay. 
So one task might take five times longer 

126
00:08:45,899 --> 00:08:48,350
than the average and so there's little 
benefit. 

127
00:08:48,350 --> 00:08:50,970
And so one take away here. 
So I'll tell you what one way, one way to 

128
00:08:50,970 --> 00:08:53,300
partially solve this problem on the next 
slide. 

129
00:08:53,300 --> 00:08:55,984
But the take away here is, if someone 
asks you what the, you know, one of the 

130
00:08:55,984 --> 00:08:58,844
biggest performance bottlenecks of 
MapReduce is, you should say stragglers 

131
00:08:58,844 --> 00:09:03,530
or skew. 
skew is more of a term for this in the 

132
00:09:03,530 --> 00:09:10,050
database community but stragglers is a 
little bit more common. 

133
00:09:10,050 --> 00:09:13,522
So this is a straggler task that takes a 
lot longer. 

134
00:09:13,522 --> 00:09:17,476
Alright, so what can we do about this? 
Well, here's the situation where we have 

135
00:09:17,476 --> 00:09:19,855
a lot of keys that all ended up on the 
same reducer and this guy's taking too 

136
00:09:19,855 --> 00:09:24,350
long. 
One thing we can do is split this reducer 

137
00:09:24,350 --> 00:09:30,274
up into more reducers. 
So take, allocate a few more reduce tasks 

138
00:09:30,274 --> 00:09:35,701
and move that data over here and split it 
into three more. 

139
00:09:35,701 --> 00:09:40,050
Now, broadcast that little red relation, 
right? 

140
00:09:40,050 --> 00:09:43,263
So I made it disappear from over here 
recognizing that this is perhaps small 

141
00:09:43,263 --> 00:09:46,870
and replicate it. 
To all three of these reducers. 

142
00:09:46,870 --> 00:09:52,570
So, sort of a combination of the 
replicated or Broadcast join and the 

143
00:09:52,570 --> 00:09:58,700
regular reduce side, pass join. 
Okay. 

144
00:09:58,700 --> 00:10:01,220
So, now, you've got three reducers 
working on this problem as opposed to 

145
00:10:01,220 --> 00:10:03,006
just one. 
And you get things a, a little bit more 

146
00:10:03,006 --> 00:10:05,985
balanced. 
And so this Skew Join is something the 

147
00:10:05,985 --> 00:10:11,978
pig can do if you specify it. 
Alright, it won't do it automatically 

148
00:10:11,978 --> 00:10:13,550
though. 
So, fine. 

149
00:10:13,550 --> 00:10:15,150
And so now we get our entire Join 
relation. 

150
00:10:17,730 --> 00:10:21,036
So finally Merge join is the third 
special case and this as you saw the 

151
00:10:21,036 --> 00:10:24,864
first special case was an opportunity to 
do things much faster, if one relation 

152
00:10:24,864 --> 00:10:30,580
was very large and other relation was 
very small to fit in memory. 

153
00:10:30,580 --> 00:10:34,545
Skew join is more way it is there is a 
problem that can occur and you need a 

154
00:10:34,545 --> 00:10:38,730
special trick to be able to address the 
problem. 

155
00:10:38,730 --> 00:10:43,432
Merge join is more like the formal. 
It's looking at an opportunity to use up 

156
00:10:43,432 --> 00:10:49,550
more high performance algorithm when 
certain conditions are met. 

157
00:10:49,550 --> 00:10:54,441
And so one of those conditions well, when 
you recognize that red relation and the 

158
00:10:54,441 --> 00:10:58,821
blue relation are already been 
co-partitioned on the appropriate Join 

159
00:10:58,821 --> 00:11:02,589
key. 
Right. 

160
00:11:02,589 --> 00:11:07,450
Then the map phase alone has enough 
information to just keep enjoying itself. 

161
00:11:07,450 --> 00:11:10,245
It doesn't actually need to assign it all 
to assign this tuple to a joined until 

162
00:11:10,245 --> 00:11:12,740
you shuffle it across the network. 
Right. 

163
00:11:12,740 --> 00:11:21,060
So you know that all the line items for a 
particular order. 

164
00:11:22,270 --> 00:11:24,601
You know, for every order that's here, 
all of it's line items are also here on 

165
00:11:24,601 --> 00:11:28,345
this machine. 
If you know that to be true then you can 

166
00:11:28,345 --> 00:11:33,071
do it in the map phase, okay. 
So the question maybe when is that true? 

167
00:11:33,071 --> 00:11:38,320
well, that's when we go back to the 
co-root operator. 

168
00:11:38,320 --> 00:11:44,043
It's possible that you've already 
partitioned these two tables on the 

169
00:11:44,043 --> 00:11:50,868
appropriate JOIN key because of a 
previous command in pig. 

170
00:11:50,868 --> 00:11:55,330
And so can be aware of that and use a 
merge, Merge join. 

171
00:11:55,330 --> 00:12:02,680
Sorry, I, say aware of that. 
If you still specify explicitly you want 

172
00:12:02,680 --> 00:12:08,375
to use a Merge join but, so take exactly 
that not automatically but you can take 

173
00:12:08,375 --> 00:12:15,370
advantage of the situation where when it 
arises, alright. 

174
00:12:15,370 --> 00:12:17,554
So, since each map already has local 
access to the records from both 

175
00:12:17,554 --> 00:12:19,564
relations. 
They're already grouped and assorted by 

176
00:12:19,564 --> 00:12:22,059
the Join key. 
You can just read in both relations from 

177
00:12:22,059 --> 00:12:26,931
disk in order and compute the Join. 
Alright so we had this CoGroup operation 

178
00:12:26,931 --> 00:12:31,015
and we have Join operation and why do we 
need both. 

179
00:12:31,015 --> 00:12:34,301
Well the reason we made CoGroup explicit 
is that if you think about a join as 

180
00:12:34,301 --> 00:12:37,799
really a two step process right there's 
one step to create the groups based on 

181
00:12:37,799 --> 00:12:44,350
the Join key. 
And then a second step to actually 

182
00:12:44,350 --> 00:12:47,768
produce the joined tuples. 
But that group creation step is useful 

183
00:12:47,768 --> 00:12:51,355
for a lot of applications, not just 
producing a JOIN. 

184
00:12:51,355 --> 00:12:59,539
So for example, if you want to, CoGROUP 
and add up all the contributions from 

185
00:12:59,539 --> 00:13:07,470
each relation you can do that in a single 
step. 

186
00:13:07,470 --> 00:13:09,888
You don't need to sort of join the two 
groups first then do another grouping 

187
00:13:09,888 --> 00:13:13,302
afterward, which is something you would 
have to in original databases. 

188
00:13:13,302 --> 00:13:15,734
So that's a chance to take what would be 
two bad produced jobs and combine them 

189
00:13:15,734 --> 00:13:20,495
into one. 
Okay, and so here I guess the example if 

190
00:13:20,495 --> 00:13:29,154
you just want to count the tuples. 
Okay. 

191
00:13:29,154 --> 00:13:34,000
So, but, you know, what the point out 
here is that you can't express Join. 

192
00:13:34,000 --> 00:13:37,250
JOIN is essentially just syntactic sugar, 
you can't express a JOIN in terms of 

193
00:13:37,250 --> 00:13:40,116
CoGroup. 
First is, the first step is to CoGroup on 

194
00:13:40,116 --> 00:13:42,870
the same colu, on the, on the JOIN 
columns. 

195
00:13:42,870 --> 00:13:47,595
And the second is to run this foreach 
command that would generate A flat view 

196
00:13:47,595 --> 00:13:54,370
of results of revenue. 
Sorry I guess I should say A and B here 

197
00:13:54,370 --> 00:13:58,055
that's sloppy. 
I changed the names. 

198
00:13:58,055 --> 00:14:08,870
Okay. 
Alright. 

199
00:14:08,870 --> 00:14:11,480
So, other commands that we're not 
going to talk about in detail are Store 

200
00:14:11,480 --> 00:14:15,440
which writes data out to HTFS. 
So, it's is available for. 

201
00:14:15,440 --> 00:14:18,328
You know, future commands. 
union that combines two data sets 

202
00:14:18,328 --> 00:14:22,174
together and removes duplicates. 
Cross product which finds all possible 

203
00:14:22,174 --> 00:14:25,828
pairs between two data sets which can be 
useful if you're going to compute some 

204
00:14:25,828 --> 00:14:29,700
sort of similarity function as we talked 
about. 

205
00:14:29,700 --> 00:14:32,270
and then dump which prints outputs of the 
screen. 

206
00:14:32,270 --> 00:14:37,634
and then order which sorts, sorts the 
output. 

207
00:14:37,634 --> 00:14:40,250
Okay. 
So as an example of Store here. 

208
00:14:40,250 --> 00:14:42,872
remember you could also, just, just like 
you can use your own custom function to 

209
00:14:42,872 --> 00:14:46,190
parse data, you can also use your own 
custom function to write it out. 

210
00:14:46,190 --> 00:14:48,264
Okay. 
Which means it allows the data to be more 

211
00:14:48,264 --> 00:14:52,032
compatible with, say, some other system 
that you're using. 

212
00:14:52,032 --> 00:14:54,156
Say, MapReduce itself or some other 
application that expects data in a 

213
00:14:54,156 --> 00:14:56,700
certain way. 
And so that's, again, I sort of stress 

214
00:14:56,700 --> 00:15:00,956
this with, with the load command as well. 
But this is actually pretty powerful and 

215
00:15:00,956 --> 00:15:03,945
a pretty big difference from this, you 
know, walled garden approach that 

216
00:15:03,945 --> 00:15:07,580
relational databases take. 
Where everything sort of goes in and then 

217
00:15:07,580 --> 00:15:10,058
it's, you know, under complete control, 
the database. 

218
00:15:10,058 --> 00:15:12,354
This sort of has more permeable 
boundaries where you can have kind of 

219
00:15:12,354 --> 00:15:14,590
have data lying around in whatever 
format. 

220
00:15:14,590 --> 00:15:17,876
And you're still able to process it with 
pig but you can also process it with 

221
00:15:17,876 --> 00:15:21,960
other, other systems and as well. 
And so the reasons this is crucially 

222
00:15:21,960 --> 00:15:24,360
important is not just sort of a 
performance optimization for, you know, 

223
00:15:24,360 --> 00:15:27,594
reduced load time or something although 
sometimes I can help. 

224
00:15:27,594 --> 00:15:30,426
Is because the data is too big to move 
nowadays you can't put it all in the 

225
00:15:30,426 --> 00:15:33,450
database and suck it all out and move it 
over to some other system for working 

226
00:15:33,450 --> 00:15:37,900
with, it's all two big, you can't move a 
petabyte, right? 

227
00:15:37,900 --> 00:15:40,294
You have to bring the computation to the 
data as opposed to bring the data the 

228
00:15:40,294 --> 00:15:44,584
computation, okay? 
And so these abilities to work with in 

229
00:15:44,584 --> 00:15:50,788
situ data, you know, produce sort of in 
situ data is emerging as a sort of a key 

230
00:15:50,788 --> 00:15:56,992
requirement in this big data era that was 
not such, not such a key requirement in 

231
00:15:56,992 --> 00:16:10,173
the sort of era of relational databases. 
Okay. 

