1
00:00:00,025 --> 00:00:04,575
So we showed how to implement loops in 
MapReduce and optimize them and utilize 

2
00:00:04,575 --> 00:00:09,125
that to implement the data log pattern 
language in recursive programs in that 

3
00:00:09,125 --> 00:00:15,335
language, and we showed different 
representations of graphs. 

4
00:00:15,335 --> 00:00:21,680
And now we're going to show PageRank in a 
couple of different ways at scale. 

5
00:00:21,680 --> 00:00:24,345
So remember, the context here is that 
graphs are getting bigger and bigger and 

6
00:00:24,345 --> 00:00:27,028
bigger. 
So you can imagine a social scale graph 

7
00:00:27,028 --> 00:00:31,680
having about a billion vertices, one per 
person, with maybe 100 billion edges. 

8
00:00:31,680 --> 00:00:34,416
Web scale graph is bigger, because 
there's more webpages than people, so 

9
00:00:34,416 --> 00:00:37,690
maybe 50 billion vertices and a trillion 
edges. 

10
00:00:37,690 --> 00:00:39,390
But there’s, you know, the human 
connectdome. 

11
00:00:39,390 --> 00:00:43,044
The neural network in your brain has 
maybe 100 billion edges and 100, excuse 

12
00:00:43,044 --> 00:00:46,724
me, 100 billion vertices and 100 trillion 
edges. 

13
00:00:46,724 --> 00:00:52,136
So we need large scale systems to process 
this and SNAP produces one such system. 

14
00:00:52,136 --> 00:00:55,864
There’s more coming down the pipe though. 
SNAP produces perhaps on the tail end of 

15
00:00:55,864 --> 00:00:59,576
its utility for this massive, massive 
scale, but we're going to stick with it 

16
00:00:59,576 --> 00:01:04,216
for now. 
So here's a new limitation of PageRank in 

17
00:01:04,216 --> 00:01:07,990
MapReduce. 
in the map function we have a node ID and 

18
00:01:07,990 --> 00:01:12,360
a vertex object. 
That vertex object has a couple of 

19
00:01:12,360 --> 00:01:15,496
methods that we use. 
You can get its current PageRank, 

20
00:01:15,496 --> 00:01:19,935
n.pagerank, and you can get its adjacency 
list, n.adjacencylist. 

21
00:01:19,935 --> 00:01:22,815
So N dot page rank divided by the length 
of the adjacency list gives us the 

22
00:01:22,815 --> 00:01:25,887
fraction of the rank that we're going to 
distribute to each of our out, out, each 

23
00:01:25,887 --> 00:01:30,610
of our neighbors. 
Across each of our out edges. 

24
00:01:30,610 --> 00:01:33,370
Then we're going to admit a special key 
value pair with the node I D and the 

25
00:01:33,370 --> 00:01:36,655
vertex object. 
And I'll come back to that in a bit. 

26
00:01:36,655 --> 00:01:40,313
And then for each of my outgoing 
neighbors Send that, send the appropriate 

27
00:01:40,313 --> 00:01:45,150
fraction of my age rank, which is value p 
that we just constructed. 

28
00:01:45,150 --> 00:01:52,155
[UNKNOWN] On the reduced side, we've got 
a node id m and a list of things. 

29
00:01:52,155 --> 00:01:56,867
Then we initialize a couple of values, 
this M object to null and a, a new rank 

30
00:01:56,867 --> 00:02:01,182
to zero. 
And then for every p in this list of p's 

31
00:02:01,182 --> 00:02:05,370
we're going to check to see if its a 
vertex. 

32
00:02:05,370 --> 00:02:08,583
Now why that's there is if it is a vertex 
that means it corresponds to this key 

33
00:02:08,583 --> 00:02:12,930
value pair that we've mated up there. 
And all this was was a way to pass that 

34
00:02:12,930 --> 00:02:15,810
complex vertex object through to the 
reducer. 

35
00:02:15,810 --> 00:02:20,336
Because remember the key here is just a 
note ID So we need some way of passing 

36
00:02:20,336 --> 00:02:25,227
this, this more complex object with some 
internal structure over to the reduced 

37
00:02:25,227 --> 00:02:28,955
side. 
And so now we have it and we assign it to 

38
00:02:28,955 --> 00:02:31,218
M. 
If it’s not a vertex then it’s one of 

39
00:02:31,218 --> 00:02:36,672
these other PageRank values that we 
computed and therefore we just add it up. 

40
00:02:36,672 --> 00:02:41,660
And then finally, the new PageRank For 
that vertex M, is going to be this, 

41
00:02:41,660 --> 00:02:46,650
formula that we saw. 
Well, so, sorry, I guess I'm guess I'm 

42
00:02:46,650 --> 00:02:49,904
being a little glib there. 
This is the damping factor that we 

43
00:02:49,904 --> 00:02:53,247
mentioned. 
And this is a, equivalent expression to 

44
00:02:53,247 --> 00:02:57,104
what we showed before. 
And finally now that we've constructed 

45
00:02:57,104 --> 00:03:01,073
the new page rank, we emit a key value 
pair, the note ID, and the Vertex object 

46
00:03:01,073 --> 00:03:04,979
itself, which is the same key value types 
as we need on the map side to repeat 

47
00:03:04,979 --> 00:03:08,948
this, so we're going to run this over, 
and over, and over again, with different 

48
00:03:08,948 --> 00:03:15,190
cross-multiple iterations. 
So, there's some problems with this 

49
00:03:15,190 --> 00:03:17,490
implementation. 
The main one is that the entire state of 

50
00:03:17,490 --> 00:03:22,555
the graph is shuffled on every iteration. 
Right, that little complex vertex object 

51
00:03:22,555 --> 00:03:26,449
it includes within it the adjacency list 
associated with that vertex is sent 

52
00:03:26,449 --> 00:03:31,890
across the network over to the reduced 
side and handled by the reducer. 

53
00:03:31,890 --> 00:03:34,890
When all that really needs to be singed 
is just the new rank contribuions. 

54
00:03:34,890 --> 00:03:37,015
If the vertex stake could sort of stay in 
place. 

55
00:03:37,015 --> 00:03:40,614
If you know you could address the 
neighboring vertices directly and just 

56
00:03:40,614 --> 00:03:43,940
say, hey here's your new rank 
contribution. 

57
00:03:45,120 --> 00:03:47,285
That would be useful, that would save 
some communication traffic. 

58
00:03:47,285 --> 00:03:51,311
And then also we have to control the 
iteration outside of MapReduce, including 

59
00:03:51,311 --> 00:03:55,280
determination conditions and just the 
logic itself. 

60
00:03:56,920 --> 00:04:01,752
So, for these reasons and just to explore 
programming model, Pregel was suggested. 

61
00:04:01,752 --> 00:04:07,322
also at Google, are we in 2010? 
So there's been, since then there's been 

62
00:04:07,322 --> 00:04:11,348
open sourcing limitations just like there 
as for Map Reduce and Apache Giraph, 

63
00:04:11,348 --> 00:04:16,860
Stanford there's a system called GPS, a 
system called Jpregel, Hama. 

64
00:04:16,860 --> 00:04:20,143
And this is really designed for batch 
algorithms on large graphs in particular, 

65
00:04:20,143 --> 00:04:24,310
so focused on graphs, and so the basic 
logic looks something like this. 

66
00:04:24,310 --> 00:04:28,402
It says, while any vertex is s, still 
active, or the max iterations have not 

67
00:04:28,402 --> 00:04:32,428
been reached, for each vertex, process 
all the messages you get from your 

68
00:04:32,428 --> 00:04:36,718
neighbors, update your own internal 
state, and then send messages back out to 

69
00:04:36,718 --> 00:04:42,536
your neighbor. 
And then maybe set the active flag 

70
00:04:42,536 --> 00:04:44,806
appropriately. 
You know, if you don't get any messages, 

71
00:04:44,806 --> 00:04:47,368
then you're not active, or if your 
internal state doesn't change enough, 

72
00:04:47,368 --> 00:04:49,888
you're not active, and you can control 
this logic, and so the point is now 

73
00:04:49,888 --> 00:04:53,442
you're writing. 
Instead of writing a map function or a 

74
00:04:53,442 --> 00:04:56,214
reduce function, you're essentially 
writing just one little function, which 

75
00:04:56,214 --> 00:04:58,818
is the, the, the function that I should 
run on each vertex at every, at every 

76
00:04:58,818 --> 00:05:03,362
time step. 
All right, so let's see the example, 

77
00:05:03,362 --> 00:05:09,067
which is PageRank. 
So there's just one method here, compute. 

78
00:05:09,067 --> 00:05:13,407
This code is a little awkward, this is 
straight from the Pregl paper, it's a 

79
00:05:13,407 --> 00:05:17,677
little awkward here because it's in C, so 
bear with me if you're not used to 

80
00:05:17,677 --> 00:05:23,870
looking at C. 
But this says if basically the iteration 

81
00:05:23,870 --> 00:05:27,770
number is greater than or equal to one, 
this super step then add up all the 

82
00:05:27,770 --> 00:05:32,125
messages I get from my neighbors and then 
assign my own internal state to this new 

83
00:05:32,125 --> 00:05:38,740
calculation, the new range. 
Then as long as the super step is less 

84
00:05:38,740 --> 00:05:43,452
than 30 Which is our, you know, if we go 
over 30, we're just going to assume that 

85
00:05:43,452 --> 00:05:50,822
we converge and just quit. 
Then send the contribution by rank to all 

86
00:05:50,822 --> 00:05:57,722
my neighbors, which is my current divided 
by n, and n is just the number of number 

87
00:05:57,722 --> 00:06:04,670
of neighbors. 
Otherwise, set myself to inactive, which 

88
00:06:04,670 --> 00:06:08,210
is this vote to halt, so you sort of say 
I've got no more for to do, I"m done 

89
00:06:08,210 --> 00:06:12,409
until somebody wakes me up with a new 
value. 

90
00:06:12,409 --> 00:06:16,161
And what's not in here in this 
implementation but is often added in a 

91
00:06:16,161 --> 00:06:20,583
Pregl implementation is some kind of 
condition that says, well if the internal 

92
00:06:20,583 --> 00:06:25,005
value hasn't changed by a large enough 
amount then vote to halt, then set myself 

93
00:06:25,005 --> 00:06:30,100
to inactive. 
Alright. 

94
00:06:30,100 --> 00:06:32,872
So let's see how this works. 
[UNKNOWN] the code and it may not be very 

95
00:06:32,872 --> 00:06:35,000
clear. 
So imagine we have a graph like this. 

96
00:06:35,000 --> 00:06:41,300
We initialize the rank to the total 
number of vertices divided by the number 

97
00:06:41,300 --> 00:06:48,563
of, excuse me. 
The number one divided by the number 

98
00:06:48,563 --> 00:06:52,710
vertices. 
So there's five vertices, so we just 

99
00:06:52,710 --> 00:06:58,798
divide one by five, which is 0.2, okay. 
So that's the initial rank. 

100
00:06:58,798 --> 00:07:04,510
Then we each, the program, the compute 
program runs on every vertex 

101
00:07:04,510 --> 00:07:10,936
simultaneously in parallel and it sends 
its rank contribution to all of its 

102
00:07:10,936 --> 00:07:17,566
neighbors, so this rank contribution of 
0.2 divided by two outgoing edges, so 

103
00:07:17,566 --> 00:07:26,055
each one gives 0.1. 
And here there's three outgoing edges, so 

104
00:07:26,055 --> 00:07:37,293
each one gets 0.066, and so on. 
Then on the next super step, all those 

105
00:07:37,293 --> 00:07:44,458
conpr, contributions are added up, and so 
here we add 0.1 and 0.066. 

106
00:07:44,458 --> 00:07:51,890
Multiply by 0.85, and add in 0.03, cause 
that's our formula, and we get 0.172. 

107
00:07:51,890 --> 00:07:55,460
And we do the same for all the other 
nodes. 

108
00:07:55,460 --> 00:07:59,812
Now, this one's a little funny, why did 
0.0, why did this node and this node get 

109
00:07:59,812 --> 00:08:02,588
0.03? 
Well that's because they have a zero, 

110
00:08:02,588 --> 00:08:06,692
they have no incoming edges at all. 
In which case, this term in the formula 

111
00:08:06,692 --> 00:08:12,160
provided nothing, the sum was zero and so 
they just give 0.15 divided by 5 0.03. 

112
00:08:12,160 --> 00:08:21,580
Okay, and so that's our new value, after 
that iteration. 

113
00:08:21,580 --> 00:08:30,530
Then it repeats This rank is divided in 
half, 0.015 and 0.015. 

114
00:08:30,530 --> 00:08:33,710
This rank is divided into thirds, 0.01, 
0.01, 0.01 and so on. 

115
00:08:33,710 --> 00:08:41,500
And we see what changes. 
So, we add this up. 

116
00:08:41,500 --> 00:08:46,865
This plus this times 0.85. 
Plus 0.03 gives us 0.0513 and so on and 

117
00:08:46,865 --> 00:08:51,625
then I've marked these verticies as red 
because their value didn't change from on 

118
00:08:51,625 --> 00:08:58,050
iteration to the next. 
And so they're marked as inactive. 

119
00:08:58,050 --> 00:09:01,283
And that condition actually wasn't in the 
code that I showed you, but it's 

120
00:09:01,283 --> 00:09:04,501
typically a condition you add. 
Okay. 

121
00:09:06,890 --> 00:09:10,360
And so then the process continues. 
Here's the new values for the next 

122
00:09:10,360 --> 00:09:12,950
generation. 
We start over again. 

123
00:09:14,200 --> 00:09:19,782
one optimization here is that these 
messages aren't actually resent again. 

124
00:09:19,782 --> 00:09:24,462
They can be retained on the receiving 
side you know all the incoming edges have 

125
00:09:24,462 --> 00:09:30,740
a value and it doesn't necessarily need 
to be resent across the network. 

126
00:09:30,740 --> 00:09:33,290
You can just sit there until it's 
updated. 

127
00:09:33,290 --> 00:09:35,770
So if this, if this vertex were active it 
would be resent with a new value but 

128
00:09:35,770 --> 00:09:38,410
since it's not active you just use the 
old value from the previous iteration and 

129
00:09:38,410 --> 00:09:41,929
nothing needs to be sent. 
All right. 

130
00:09:41,929 --> 00:09:46,052
So we recompute these. 
And we find out that this top vertex 

131
00:09:46,052 --> 00:09:49,750
didn't change from one iteration to the 
other. 

132
00:09:53,560 --> 00:09:56,987
So it goes red. 
And then these two values are updated 

133
00:09:56,987 --> 00:10:01,481
appropriately. 
And that's our new values for the next 

134
00:10:01,481 --> 00:10:04,280
iteration. 
And then finally, every, remember all the 

135
00:10:04,280 --> 00:10:07,660
red nodes don't actually, all the red 
vertices don't actually get executed at 

136
00:10:07,660 --> 00:10:10,990
all. 
And so they don't send any messages. 

137
00:10:10,990 --> 00:10:13,140
And they don't send any messages. 
Nobody needs them. 

138
00:10:14,460 --> 00:10:16,160
But this vertex still needs to be 
updated. 

139
00:10:16,160 --> 00:10:20,735
And so we apply it again, and then 
finally the execution halts when all the 

140
00:10:20,735 --> 00:10:25,728
vertices are inactive. 
Now it's, it might take quite a while to 

141
00:10:25,728 --> 00:10:28,596
con, to converge. 
In which case, remember there's that 

142
00:10:28,596 --> 00:10:30,870
super step less than 30 condition as 
well. 

143
00:10:30,870 --> 00:10:34,582
So you don't really need all that, in 
general, PageRank starts to, get pretty 

144
00:10:34,582 --> 00:10:38,526
close to conversion after some number of 
iterations, you don't really need to let 

145
00:10:38,526 --> 00:10:41,330
it run for to long. 

