1
23:59:59,916 --> 00:00:05,531
[MUSIC]. 

2
00:00:05,531 --> 00:00:07,761
So let's take a step back and see where 
we are. 

3
00:00:07,761 --> 00:00:10,946
We talked about graph tasks in general, 
give an example of constructing these 

4
00:00:10,946 --> 00:00:14,910
histograms, talked about structural, 
traversal, and pattern tasks. 

5
00:00:14,910 --> 00:00:18,155
Gave page rank as an example of 
structural and of, both a structural and 

6
00:00:18,155 --> 00:00:22,726
a traversal task. 
And then talk about pattern languages in 

7
00:00:22,726 --> 00:00:26,845
SPARQL and SQL and Datalog. 
Gave a few more detailed examples of 

8
00:00:26,845 --> 00:00:31,904
Datalog motivated by what's been going on 
in the news in June 2013. 

9
00:00:31,904 --> 00:00:36,192
Related to the Prism system, but what we 
haven't talked about is how to implement 

10
00:00:36,192 --> 00:00:41,160
much of this, especially at scale. 
So that's what I want to do next. 

11
00:00:41,160 --> 00:00:43,950
We're going to stick with Datalog for a 
moment and think about how to actually 

12
00:00:43,950 --> 00:00:46,660
implement those kinds of queries in 
MapReduce. 

13
00:00:46,660 --> 00:00:51,666
And then talk about how to implement 
PageRank in MapReduce. 

14
00:00:51,666 --> 00:00:54,546
So we're back towards the beginning of 
the course, we talked about how to 

15
00:00:54,546 --> 00:00:58,380
implement a relational algebra operations 
in MapReduce at scale. 

16
00:00:58,380 --> 00:01:01,395
And we showed this sort a connection, you 
know, between languages like pig in the 

17
00:01:01,395 --> 00:01:04,438
relational algebra. 
What databases do or is obviously 

18
00:01:04,438 --> 00:01:08,630
relation algebra and what you can 
implement directly MapReduce. 

19
00:01:08,630 --> 00:01:12,846
So there seems to be this kind of 
relational algebra DNA in a lot of these 

20
00:01:12,846 --> 00:01:19,149
large large scale programming platforms. 
However, all of those cases, we never had 

21
00:01:19,149 --> 00:01:24,180
anything that had to loop, that had to do 
recursion, that had to do iteration. 

22
00:01:24,180 --> 00:01:28,665
And now that we're talking about graphs 
as we sort of pointed out that traversal 

23
00:01:28,665 --> 00:01:34,240
of these graphs is really the salient 
future of working with them. 

24
00:01:34,240 --> 00:01:38,456
And so it can't have been traversal tasks 
and it came up in pattern languages in 

25
00:01:38,456 --> 00:01:41,569
Datalog. 
We saw examples of recursive data log 

26
00:01:41,569 --> 00:01:45,000
programs without talking much about how 
they'll be implemented. 

27
00:01:45,000 --> 00:01:48,821
So I want to cover that a little bit now. 
So, given if you have a recursive Datalog 

28
00:01:48,821 --> 00:01:53,233
program that looks like this. 
So you can say this is an edge relation 

29
00:01:53,233 --> 00:01:58,317
that's just linking an arbitrary vertex x 
with an arbitrary vertex y. 

30
00:01:58,317 --> 00:02:02,029
Every pair of edges is stored in a 
relation called edge and I write this 

31
00:02:02,029 --> 00:02:05,970
program. 
This program says find me all the edges 

32
00:02:05,970 --> 00:02:10,740
connected to the vertex labeled w. 
And here we have a vertex labeled w. 

33
00:02:12,300 --> 00:02:17,340
Input those in a new relation called A. 
Then also, find me additional values of A 

34
00:02:17,340 --> 00:02:21,207
in this manner. 
For every vertex that's already in A, go 

35
00:02:21,207 --> 00:02:25,303
find me an edge that leads out from that 
vertex, and give me a new generation of 

36
00:02:25,303 --> 00:02:29,170
vertices. 
So this is now a recursive Datalog 

37
00:02:29,170 --> 00:02:33,190
program but which you can think about it 
intuitively what is doing is traversing 

38
00:02:33,190 --> 00:02:37,520
this graph in kind of a bad breath first 
way right? 

39
00:02:37,520 --> 00:02:41,922
So first we fire this rule which just 
says out of all the vertices select the 

40
00:02:41,922 --> 00:02:45,640
one with w. 
And we know how to do this, we could 

41
00:02:45,640 --> 00:02:49,555
write this in SQL, I would hope, right? 
We can write this in MapReduce, we can 

42
00:02:49,555 --> 00:02:53,464
write this in Python. 
Finding all edges where the first 

43
00:02:53,464 --> 00:02:57,660
position is, has a labeling of w, and 
call that A0. 

44
00:02:57,660 --> 00:03:02,960
Or find me everything that's connected to 
w call all those collectively A0. 

45
00:03:02,960 --> 00:03:06,929
Next we go to this rule and now we have 
to join All the ones we've already found 

46
00:03:06,929 --> 00:03:12,200
with the edge relation to produce a new 
generation of vertices. 

47
00:03:12,200 --> 00:03:14,645
For each node in A0 find all of its 
neighbors. 

48
00:03:14,645 --> 00:03:21,320
And then in step two we fire this rule 
again, for each node in A1. 

49
00:03:21,320 --> 00:03:24,738
Find the neighbors, and so on. 
So we keep finding more and more, 

50
00:03:24,738 --> 00:03:28,022
vertices in this graph as we go one 
generation away. 

51
00:03:28,022 --> 00:03:33,720
And that's, naively how we would, how we 
might evaluate this recursive program. 

52
00:03:33,720 --> 00:03:35,794
That's what this, that's what the meaning 
of this program is. 

53
00:03:35,794 --> 00:03:40,050
Okay, but there's a problem here. 
What if there's a cycle in the graph? 

54
00:03:40,050 --> 00:03:43,230
Well, at generation one, we find Sue and 
generation two, we find Tom, and 

55
00:03:43,230 --> 00:03:47,646
generation three, we find Lin. 
And in generation five we find Sue again, 

56
00:03:47,646 --> 00:03:51,440
and so we just keep going. 
So we need one more piece, which is to 

57
00:03:51,440 --> 00:03:54,715
remove all the ones we've already found 
previously. 

58
00:03:54,715 --> 00:03:59,610
So let's look at this again, so this is 
called semi-naive evaluation. 

59
00:03:59,610 --> 00:04:03,150
So here's this reachability query that we 
were just looking at. 

60
00:04:03,150 --> 00:04:07,830
Find the everybody that's connected to 
everybody that's in the, in the social 

61
00:04:07,830 --> 00:04:14,560
network of some vertex labeled A, right. 
Everybody is reachable from A. 

62
00:04:14,560 --> 00:04:17,726
So here's how it would work. 
We're going to have these delta 

63
00:04:17,726 --> 00:04:20,837
relations, which are just the new gener, 
just everything found in the new 

64
00:04:20,837 --> 00:04:24,441
generation, right? 
Just like those colored levels were in 

65
00:04:24,441 --> 00:04:29,298
the previous pictorial representation. 
So while we're still finding new 

66
00:04:29,298 --> 00:04:33,092
vertices, were going to keep going. 
If at any, if at any point we don't find 

67
00:04:33,092 --> 00:04:35,570
anything new then we know we're done and 
we can stop and we can return the 

68
00:04:35,570 --> 00:04:40,474
results. 
Okay so while we're finding new things, 

69
00:04:40,474 --> 00:04:48,720
let's call A to Ai the union of all the 
deltas that we found previously. 

70
00:04:48,720 --> 00:04:53,972
The iteration zero, union with iteration 
1, union with iteration 2, and so on. 

71
00:04:53,972 --> 00:04:56,110
And if you're not familiar with this 
notation, this just means. 

72
00:04:56,110 --> 00:04:58,761
Think of it as concatenation in a 
programming language, right? 

73
00:04:58,761 --> 00:05:03,440
Except you're also going to remove 
duplicates, right? 

74
00:05:03,440 --> 00:05:10,645
Then delta Ai is, the previous new 
results joined with the relation R. 

75
00:05:10,645 --> 00:05:13,756
Alright, so this gives me one more 
generation, and I'm expressing this in 

76
00:05:13,756 --> 00:05:18,221
the relational algebra notation of join. 
So that's what we showed pictorially a 

77
00:05:18,221 --> 00:05:20,808
moment ago. 
However now what we also want to do is 

78
00:05:20,808 --> 00:05:24,343
remove any duplicates. 
We want to subtract out anything we found 

79
00:05:24,343 --> 00:05:29,250
in any previous generation at all, and 
this is in order to break cycles. 

80
00:05:29,250 --> 00:05:32,102
If we've already seen that, if we've 
already found Sue, then we don't want to 

81
00:05:32,102 --> 00:05:35,431
find her again, alright? 
And then we increment this calendar so we 

82
00:05:35,431 --> 00:05:38,944
can keep track of everything. 
And then our final result, when we're 

83
00:05:38,944 --> 00:05:42,100
done, is this Ai for whatever i is at the 
end. 

84
00:05:43,110 --> 00:05:47,014
And so I'd claim that, if we're going to 
do this in MapReduce, which is what we're 

85
00:05:47,014 --> 00:05:51,700
trying to build up to, then this line is 
the important line. 

86
00:05:52,760 --> 00:05:55,640
This line is actually pretty easy in map 
reduce because of an implementation 

87
00:05:55,640 --> 00:05:59,058
detail. 
A union of a bunch of different data 

88
00:05:59,058 --> 00:06:02,698
sets, if they're known to be disjoint, is 
actually literally just concatenation of 

89
00:06:02,698 --> 00:06:06,033
the files. 
Which sort of happens implicitly in, in 

90
00:06:06,033 --> 00:06:08,040
MapReduce. 
And I'm not going to go into much more 

91
00:06:08,040 --> 00:06:10,312
detail than that. 
But, I just to say that it to say that 

92
00:06:10,312 --> 00:06:13,390
this line is, can be arranged so that its 
not very expensive. 

93
00:06:13,390 --> 00:06:17,000
Its essentially free, but this line is 
where there's a lot of work going on. 

94
00:06:17,000 --> 00:06:20,900
So lets look at, lets break this down and 
pretend like and show how this is going 

95
00:06:20,900 --> 00:06:26,581
to be implemented in MapReduce. 
Alright, so its actually two different 

96
00:06:26,581 --> 00:06:32,447
MapReduce joins, one to do the join and 
one to do the difference. 

97
00:06:32,447 --> 00:06:37,970
And we saw this in the lecture on pig 
where join was one operation. 

98
00:06:37,970 --> 00:06:43,054
And we didn't show a pig operation 
corresponding to difference but you can 

99
00:06:43,054 --> 00:06:48,309
implement it using grouping. 
Okay, but we did see that generally one 

100
00:06:48,309 --> 00:06:51,957
operation and another operation would 
turn into two MapReduce jobs chained 

101
00:06:51,957 --> 00:06:55,491
together. 
And that's whats going on here, so what 

102
00:06:55,491 --> 00:07:00,240
happens with the join is, we have the 
Delta A from the previous iteration. 

103
00:07:00,240 --> 00:07:03,430
Along with the, the hu, the hu, the very 
large edge relation, right? 

104
00:07:03,430 --> 00:07:07,460
This is billions of edges, perhaps, 
spread across all sorts of computers. 

105
00:07:08,870 --> 00:07:12,740
And we process them all in parallel, 
spray them across the network. 

106
00:07:12,740 --> 00:07:15,870
And per, and implement the reduce phase 
in order to compute the join. 

107
00:07:17,330 --> 00:07:20,036
Then do the difference we find everything 
we just got, everything we just found 

108
00:07:20,036 --> 00:07:24,446
from the join. 
Spread across the network again, along 

109
00:07:24,446 --> 00:07:29,715
with everything that we've ever seen 
before. 

110
00:07:29,715 --> 00:07:33,789
Everything in this Ai, everything in the 
union. 

111
00:07:33,789 --> 00:07:36,310
And make sure, make sure to remove 
duplicates. 

112
00:07:36,310 --> 00:07:38,686
And you implemented both of these things, 
or specified how to implement both of 

113
00:07:38,686 --> 00:07:41,975
these things in the MapReduce assignment. 
You showed how to do a join and you 

114
00:07:41,975 --> 00:07:44,855
showed how to do a, duplicate 
elimination, essentially, which is, which 

115
00:07:44,855 --> 00:07:48,891
is what this is doing. 
Okay, and then the only new piece we 

116
00:07:48,891 --> 00:07:53,423
need, which is where MapReduce gets a 
little bit, 

117
00:07:53,423 --> 00:07:57,632
Becomes a little bit of a poor fit for 
this task is that you need some sort of 

118
00:07:57,632 --> 00:08:02,358
external driver program to control the 
iteration. 

119
00:08:02,358 --> 00:08:05,821
To check and see if there's anything new. 
Okay, so MapReduce can't do loops in and 

120
00:08:05,821 --> 00:08:10,120
of itself, so you have to do something 
outside of the system. 

121
00:08:10,120 --> 00:08:13,698
Alright, so again two big steps. 
Compute the next generation of vertices 

122
00:08:13,698 --> 00:08:16,375
or nodes and remove the ones we have 
already seen. 

123
00:08:16,375 --> 00:08:21,475
Alright, so this is fine and this works, 
but there's some performance issues that 

124
00:08:21,475 --> 00:08:27,050
are pretty significant, and we'll talk 
about that next time. 

