1
23:59:59,500 --> 00:00:05,603
[MUSIC]. 

2
00:00:05,603 --> 00:00:08,515
Okay so, I want to talk a little bit 
about the system Pig, and there's a 

3
00:00:08,515 --> 00:00:12,268
couple reasons for this. 
One is it represents this language layer 

4
00:00:12,268 --> 00:00:14,923
on top of map produce, which I think is 
pretty important as I've, as I've 

5
00:00:14,923 --> 00:00:19,029
mentioned before. 
And it also represents a more direct 

6
00:00:19,029 --> 00:00:23,785
interface to relational algebra. 
Which I've argued is really the important 

7
00:00:23,785 --> 00:00:27,043
abstraction from databases. 
SQL happens to be sort-of intergalactic 

8
00:00:27,043 --> 00:00:30,349
data-speak, and it is you know the 
universal language of programming 

9
00:00:30,349 --> 00:00:35,170
databases and manipulating databases. 
But really, the secret sauce, the 

10
00:00:35,170 --> 00:00:38,998
technology that makes it effective, is 
this underlying formulas and relation 

11
00:00:38,998 --> 00:00:42,260
algebra. 
And Pig and a few other systems sort of 

12
00:00:42,260 --> 00:00:45,305
exposes as a, as, as a first class 
citizen. 

13
00:00:45,305 --> 00:00:48,953
And it, it, it's not necessarily the only 
way to do it or the better way to do it, 

14
00:00:48,953 --> 00:00:52,715
but since you can see relational algebra 
I think it's a nice thing to be familiar 

15
00:00:52,715 --> 00:00:57,076
with, okay. 
and also I think makes the point that you 

16
00:00:57,076 --> 00:01:00,680
know, you, you can cherry pick a little 
bit what features you bring. 

17
00:01:00,680 --> 00:01:04,061
So, while it's very obviously inspired by 
relational algebra, and of course if you 

18
00:01:04,061 --> 00:01:07,860
look in the sort of people that developed 
it, you can see why. 

19
00:01:07,860 --> 00:01:11,094
It definitely changes some things about 
the relational model and I, and I think 

20
00:01:11,094 --> 00:01:14,083
that's an important point too, is that we 
can use this relation algebra of 

21
00:01:14,083 --> 00:01:18,900
abstraction independently of some of the 
systems that it came from. 

22
00:01:18,900 --> 00:01:21,070
We can make whatever changes we need and 
keep, you know keep what we like and 

23
00:01:21,070 --> 00:01:23,734
throw out what we don't like. 
And so I want, you know I want you to 

24
00:01:23,734 --> 00:01:26,318
sort of put your relation algebra goggles 
on and see the world in terms of relation 

25
00:01:26,318 --> 00:01:29,570
algebra. 
I think it's pretty effective for working 

26
00:01:29,570 --> 00:01:33,154
with not just data in general, but 
especially big data and reason they are 

27
00:01:33,154 --> 00:01:37,440
very scalable algorithms. 
Okay, so fine, what is Pig? 

28
00:01:37,440 --> 00:01:42,750
So, Pig is a engine for executing 
programs on top of Hadoop. 

29
00:01:42,750 --> 00:01:46,089
Okay, so when you are going to write a 
pig program in Pig, and it's going to 

30
00:01:46,089 --> 00:01:51,205
generate a sequence of map produced jobs, 
that implement that program, okay. 

31
00:01:51,205 --> 00:01:54,438
And the language here is called Pig 
Latin, and it's an Apache open sourced 

32
00:01:54,438 --> 00:01:59,730
project and you can read more about it. 
Okay, and so, why bother doing this? 

33
00:01:59,730 --> 00:02:02,246
As well, if you're coming from, if you're 
a map produce program, or if you're happy 

34
00:02:02,246 --> 00:02:04,710
with map produce, why would you bother 
doing this? 

35
00:02:04,710 --> 00:02:09,396
Well, suppose you have data in one file 
and data from websites, sorry, user data 

36
00:02:09,396 --> 00:02:13,940
in one file and website data in another, 
and you need find the top most visited 

37
00:02:13,940 --> 00:02:20,910
sites by user in a particular age range. 
And so if you're thinking in terms of 

38
00:02:20,910 --> 00:02:26,000
SQL, SQL you can probably think about how 
you express this task as a query. 

39
00:02:26,000 --> 00:02:34,070
and here's a pictorial, rep, 
representation of it. 

40
00:02:34,070 --> 00:02:37,920
That sort of shows a, a kind of data flow 
here. 

41
00:02:37,920 --> 00:02:41,365
We're going to load users, we'll load the 
pages, we'll filter by age, do some kind 

42
00:02:41,365 --> 00:02:44,598
of join on name, group, count up the 
clicks, order by clicks, and then take 

43
00:02:44,598 --> 00:02:48,070
the top five. 
This may or may not be the way you draw 

44
00:02:48,070 --> 00:02:50,302
it on the whiteboard if you were to 
express it yourself, but it's a, but it's 

45
00:02:50,302 --> 00:02:52,880
a reasonable way of describing the 
problem. 

46
00:02:52,880 --> 00:02:56,325
Well in MapReduce there's 170 lines of 
code, and at least in one example coming 

47
00:02:56,325 --> 00:03:00,460
from the papers here, it took someone 
four hours to write it. 

48
00:03:00,460 --> 00:03:03,764
These kind of metrics of how long 
development time are not never too 

49
00:03:03,764 --> 00:03:09,480
compelling, because it depends on their 
background and skill set, and so on. 

50
00:03:09,480 --> 00:03:13,000
Still, a decent amount of time, while the 
same program is, in Pig Latin is just 

51
00:03:13,000 --> 00:03:16,355
nine lines of code and takes, you know, 
at least in this case, 15 minutes to 

52
00:03:16,355 --> 00:03:19,730
write. 
And further, you know, it does sort of 

53
00:03:19,730 --> 00:03:22,750
capture some of this more abstract 
description of what's going on. 

54
00:03:22,750 --> 00:03:25,654
Now, this is a little bit contrived, 
because these boxes obviously correspond 

55
00:03:25,654 --> 00:03:29,746
to sort of particular pig commands. 
But I'd argue that this, that this even 

56
00:03:29,746 --> 00:03:33,841
non, people that don't know Pig, if they 
were to ask sort of what the tasks are in 

57
00:03:33,841 --> 00:03:37,369
eval, in evaluating this English 
question, they might come up with 

58
00:03:37,369 --> 00:03:42,896
something along these lines. 
They may or may not say join, but some of 

59
00:03:42,896 --> 00:03:46,746
these steps would be, would be present. 
So, I think it, you can wave your hands a 

60
00:03:46,746 --> 00:03:49,820
little bit and say that it, that it 
captures a little bit of the natural 

61
00:03:49,820 --> 00:03:54,135
language description of the task, okay. 
You may or may not buy that. 

62
00:03:54,135 --> 00:03:59,811
Fine, so here is what it looks like, load 
the user data according to a particular 

63
00:03:59,811 --> 00:04:05,573
schema, filter the user database on the 
age range, load the pages data, perform a 

64
00:04:05,573 --> 00:04:11,766
join group. 
You say, well, for each group, count the 

65
00:04:11,766 --> 00:04:16,693
number of clicks and we'll talk. 
This one's probably look the least like 

66
00:04:16,693 --> 00:04:19,670
relation algebra. 
And we'll talk about it. 

67
00:04:19,670 --> 00:04:23,702
And then sort the thing, and then just 
take the top five, and store that result 

68
00:04:23,702 --> 00:04:27,514
into a new file, okay. 
And so, the point here is that you could 

69
00:04:27,514 --> 00:04:30,731
also think, well boy I could just write 
all this in python, right. 

70
00:04:30,731 --> 00:04:32,865
That would be much easier, why do I need 
to use this new language. 

71
00:04:32,865 --> 00:04:37,610
Well the point that each one of these 
steps well not actually not each one. 

72
00:04:37,610 --> 00:04:40,823
Groups of these steps correspond to 
individual map reduce jobs, which mean 

73
00:04:40,823 --> 00:04:45,076
they scale really, really well, right. 
Does, doesn't matter how big your user's 

74
00:04:45,076 --> 00:04:48,212
table is, and it doesn't matter how big 
your page's table is this program will 

75
00:04:48,212 --> 00:04:52,219
work. 
Okay, so how does the system work? 

76
00:04:52,219 --> 00:04:56,855
Well the programmer's going to enter a 
Pig Latin program by writing these 

77
00:04:56,855 --> 00:05:03,090
commands [SOUND]. 
and assigning the results to variables. 

78
00:05:03,090 --> 00:05:05,430
And then you can refer to the variables 
in future commands. 

79
00:05:05,430 --> 00:05:09,374
So this is, you know, we've, we've loaded 
data into this variable A, loaded it into 

80
00:05:09,374 --> 00:05:13,260
variable B, and then you can filter by 
reference to A then sort of the result 

81
00:05:13,260 --> 00:05:18,051
and C, and so on, okay. 
And then, we'll make this point again 

82
00:05:18,051 --> 00:05:21,552
later, but nothing actually happens. 
No work is actually done until you try to 

83
00:05:21,552 --> 00:05:24,501
actually write out the results. 
So, this is sort of what we call lazy, 

84
00:05:24,501 --> 00:05:27,480
lazy evaluation. 
Alright, so fine. 

85
00:05:27,480 --> 00:05:28,970
So, what happens when you write this 
program? 

86
00:05:28,970 --> 00:05:33,572
Well the Pig parser will turn the syntax 
of the program into an abstract 

87
00:05:33,572 --> 00:05:39,360
representation of the, program called an 
execution plan. 

88
00:05:39,360 --> 00:05:43,844
And so this is using parlance from 
databases, and the lead of the Pig 

89
00:05:43,844 --> 00:05:49,760
project is a, is a great database guy 
named Chris Olson, okay. 

90
00:05:49,760 --> 00:05:53,855
So, this abstract presentation of the 
plan, you know, these operators, these 

91
00:05:53,855 --> 00:05:57,510
operations that I'm going to call 
operators your going to use in database 

92
00:05:57,510 --> 00:06:01,980
[UNKNOWN] are sort of one-to-one with the 
commands. 

93
00:06:01,980 --> 00:06:05,360
Although they need not always be. 
Okay, but we're not done there. 

94
00:06:05,360 --> 00:06:08,240
That's just the execution plan, it's 
still sort of abstract we can't actually 

95
00:06:08,240 --> 00:06:12,220
execute that, what we need to do is 
compile that down into map reduced jobs. 

96
00:06:12,220 --> 00:06:15,052
And the game here is going to sort of be 
to minimize the number of map reduced 

97
00:06:15,052 --> 00:06:19,390
jobs you need, because there's a lot of 
overhead to executing one of these. 

98
00:06:19,390 --> 00:06:22,835
So, you don't want to run a whole map 
reduced job, just to filter a data set. 

99
00:06:22,835 --> 00:06:26,130
You might as well lump that in with other 
work that you're already doing. 

100
00:06:26,130 --> 00:06:28,242
As long as you're scanning the data and 
reading off the disk, you might as well 

101
00:06:28,242 --> 00:06:31,990
do as much as you can with it. 
And so for example, in this case, this 

102
00:06:31,990 --> 00:06:37,130
program can all be implemented as just 
one single map reduce job. 

103
00:06:37,130 --> 00:06:41,156
One, in the map phase you load the data, 
you filter it, and you also load the 

104
00:06:41,156 --> 00:06:46,261
other data set, and the reduce phase you 
do the join, okay. 

105
00:06:46,261 --> 00:06:50,060
And we'll see another example of this, a 
little later on. 

106
00:06:50,060 --> 00:06:54,344
And then finally, this, these map reduce 
jobs are scheduled on a Hadoop cluster as 

107
00:06:54,344 --> 00:07:01,640
usual, end run as usual, okay. 
So, these performance [LAUGH] results 

108
00:07:01,640 --> 00:07:08,336
are, a, a little funny but I sort of 
like, I sort of, I, I still like the 

109
00:07:08,336 --> 00:07:16,579
argument to be made. 
This is, Pig Performance versus 

110
00:07:16,579 --> 00:07:19,688
Map-Reduce. 
Okay, so this is a little funny, because 

111
00:07:19,688 --> 00:07:22,884
Pig is built on top of Map-Reduce. 
But the point is, is that the first 

112
00:07:22,884 --> 00:07:26,480
version of Pig in September 11th, 2008, 
was a lot slower that map, than just 

113
00:07:26,480 --> 00:07:32,690
writing the thing raw against Map-Reduce. 
But with some various improvements it got 

114
00:07:32,690 --> 00:07:36,702
better and better over time, until 
eventually it was just as fast as the 

115
00:07:36,702 --> 00:07:42,000
handwritten Map-Reduce. 
And I think you see this sort of pattern 

116
00:07:42,000 --> 00:07:44,489
quite a bit. 
that, that there is a cost to 

117
00:07:44,489 --> 00:07:48,739
abstraction. 
But that you can, you typically can 

118
00:07:48,739 --> 00:07:52,070
recover a lot of the performance of the 
hand coded stuff. 

119
00:07:52,070 --> 00:07:55,742
Meanwhile having gained some measure of 
program productivity by offering a high 

120
00:07:55,742 --> 00:07:59,224
level interface, okay. 
And so I just like the fact that they are 

121
00:07:59,224 --> 00:08:02,641
actually doesn't bothered telling the 
story, admitting they were quite slow in 

122
00:08:02,641 --> 00:08:06,430
the beginning and they got faster over 
time, okay. 

123
00:08:06,430 --> 00:08:13,140
So, what's the data model here, so there 
is four types involved here. 

124
00:08:13,140 --> 00:08:16,668
One is the atom, which is just a, a 
primitive, right, an integer, a string 

125
00:08:16,668 --> 00:08:19,326
and so on. 
And then there's three different 

126
00:08:19,326 --> 00:08:21,434
collection types. 
One is a tuple, and these are, should be 

127
00:08:21,434 --> 00:08:24,010
familiar from our first assignment, we 
worked with Python. 

128
00:08:24,010 --> 00:08:26,830
Alright, so there's, these three all 
exist in, well I shouldn't say that 

129
00:08:26,830 --> 00:08:30,210
actually, bag is a little funny thing in 
terms of Python. 

130
00:08:30,210 --> 00:08:38,819
Lemme just say, the first is a tuple 
which is a sequence of fields. 

131
00:08:38,819 --> 00:08:42,900
And every field can be of any type. 
They need not be atoms, okay. 

132
00:08:42,900 --> 00:08:46,788
And then there's a bag which is a 
collection of tuples, always tuples. 

133
00:08:46,788 --> 00:08:49,190
But those tuples need not be the same 
type. 

134
00:08:49,190 --> 00:08:51,770
This is different than a table, different 
than a relation, okay. 

135
00:08:51,770 --> 00:08:55,155
And it's a bag not a set. 
And, you know, I guess I'll, I can ask 

136
00:08:55,155 --> 00:08:58,780
you what's the difference between a bag 
and a set? 

137
00:08:58,780 --> 00:09:04,694
Well, a bag allows duplicates, okay. 
And then the third type is a map, which 

138
00:09:04,694 --> 00:09:06,890
is a little bit confusing because we're 
talking about Map Reduce, but this is a 

139
00:09:06,890 --> 00:09:12,063
dictionary in Python, right? 
So, string literal keys map to any other 

140
00:09:12,063 --> 00:09:15,030
type, okay. 
So, let's look at an example of this. 

141
00:09:15,030 --> 00:09:20,169
So, what is this thing? 
Well depends on my notation I guess, but 

142
00:09:20,169 --> 00:09:26,996
assuming that you don't mind my angle 
brackets referring to tuple. 

143
00:09:26,996 --> 00:09:33,171
Then you can see that this outer thing is 
a tuple, and it has three fields, it has 

144
00:09:33,171 --> 00:09:39,061
a integer which is just an atom, and then 
it has this thing which is a bag, the 

145
00:09:39,061 --> 00:09:46,755
curly braces I'm going to use to indicate 
bag. 

146
00:09:46,755 --> 00:09:51,635
And then it's got this thing, whoop 
excuse me, which is a map, a dictionary 

147
00:09:51,635 --> 00:09:57,190
that has just one element in it. 
Mapping the key apache to the value, 

148
00:09:57,190 --> 00:10:01,508
search, okay. 
So, let's name these guys f1, f2 and f3, 

149
00:10:01,508 --> 00:10:08,684
and we'll point out in the, in a couple, 
in a couple slides where those names come 

150
00:10:08,684 --> 00:10:13,169
from. 
There's a few different places they can 

151
00:10:13,169 --> 00:10:15,300
come from, and I'll point out one of 
those sources soon. 

152
00:10:15,300 --> 00:10:18,510
But right now, let's just assume that we 
have them, okay. 

153
00:10:18,510 --> 00:10:21,010
Okay, so let's consider these expressions 
over this type. 

154
00:10:21,010 --> 00:10:30,430
So, we can write $0, and what that means 
is give me the first field in the tuple. 

155
00:10:30,430 --> 00:10:31,960
All right, so, in this case, it's just 
the number 1. 

156
00:10:31,960 --> 00:10:35,437
We can also access it by name, so if we 
access f2 that will give us the second 

157
00:10:35,437 --> 00:10:39,440
field, just because we gave it that 
explicit name. 

158
00:10:39,440 --> 00:10:41,880
There's nothing magic about f, nothing 
magic about 2. 

159
00:10:41,880 --> 00:10:45,460
And that'll give us a bag with these 
values in it. 

160
00:10:45,460 --> 00:10:47,550
By the way we should say what are the 
elements of this bag? 

161
00:10:47,550 --> 00:10:49,905
Well, they're tuples with two integers 
each. 

162
00:10:49,905 --> 00:10:56,290
Okay, now you can also access f2 to give 
me, so what does this expression do? 

163
00:10:56,290 --> 00:11:00,774
Well, f2 gives me the bag, and then you 
can write .$0 and I don't love this 

164
00:11:00,774 --> 00:11:05,470
notation, it's a bit of an abusive 
notation. 

165
00:11:05,470 --> 00:11:10,780
But what it means is for every element of 
the bag, project out the first element. 

166
00:11:10,780 --> 00:11:15,060
And so, this will give me another bag, 
but it only has 2 and 4 and 5. 

167
00:11:15,060 --> 00:11:22,070
It has the zeroth element of every tuple, 
okay. 

168
00:11:22,070 --> 00:11:24,998
Then you've got this magic hash symbol 
here which means finding the value 

169
00:11:24,998 --> 00:11:27,890
associated with the key I'm about to 
provide. 

170
00:11:27,890 --> 00:11:31,379
So, look at f3 which is the, the third 
field. 

171
00:11:32,930 --> 00:11:34,870
And hash into it, and give me the key 
apache. 

172
00:11:34,870 --> 00:11:38,272
If it's not a map, you're going to get an 
error, but in this case it is a map, and 

173
00:11:38,272 --> 00:11:42,000
so what gets returned is the atom search, 
okay. 

174
00:11:42,000 --> 00:11:45,366
And you can also write functions, you can 
say sum up all of these values, and this 

175
00:11:45,366 --> 00:11:50,267
expression is the same one we saw here, 
which will be a sequence of integers. 

176
00:11:50,267 --> 00:11:53,561
And so it gives you 2 plus 4 plus 5, and 
there's a few different other ways to 

177
00:11:53,561 --> 00:11:57,606
manipulate these things. 
But the first thing to notice is this is 

178
00:11:57,606 --> 00:12:00,605
non-relational, right? 
You've got these nested data structures, 

179
00:12:00,605 --> 00:12:02,300
and you got a few different data types. 

