1
00:00:00,367 --> 00:00:06,146
[MUSIC]. 

2
00:00:06,146 --> 00:00:09,721
Last time we talked about scalability, 
and we argued that scalability really 

3
00:00:09,721 --> 00:00:13,802
means working in parallel. 
And we talked about this specific task, 

4
00:00:13,802 --> 00:00:18,246
which is, you know, Read Trimming, okay? 
So this is a bunch of small genetic 

5
00:00:18,246 --> 00:00:22,856
sequences, and your task was to trim off 
the last few characters from each one. 

6
00:00:22,856 --> 00:00:28,092
Okay, and we sort of showed that this is 
pretty simple to think about in parallel. 

7
00:00:28,092 --> 00:00:32,956
You would divide the sets of, the set of 
reads into chunks and put them all onto 

8
00:00:32,956 --> 00:00:37,594
separate computers. 
And process them all in parallel. 

9
00:00:37,594 --> 00:00:41,122
Here, you know there's a function F that 
takes a single read and trims off the 

10
00:00:41,122 --> 00:00:44,428
last few characters and returns the 
prefix. 

11
00:00:44,428 --> 00:00:47,485
and you can apply this function in 
parallel. 

12
00:00:47,485 --> 00:00:51,081
And you can get out, you, the data set 
you want, which is a set of trend reads, 

13
00:00:51,081 --> 00:00:55,200
okay? 
So let's see some more examples. 

14
00:00:55,200 --> 00:00:57,853
So, a new task. 
that was new to be performed at the New 

15
00:00:57,853 --> 00:01:01,665
York Times in 2008. 
And they have a couple blog posts about 

16
00:01:01,665 --> 00:01:06,278
it that you can read was. 
In a simplified version was convert a 

17
00:01:06,278 --> 00:01:09,298
bunch of TIFF images into a different 
format. 

18
00:01:09,298 --> 00:01:14,234
And what was really going on here is that 
they had digitized. 

19
00:01:14,234 --> 00:01:18,822
Images from the newspaper, along with 
some information about the optical 

20
00:01:18,822 --> 00:01:23,056
character recognition, so some extracted 
text. 

21
00:01:23,056 --> 00:01:27,413
And they wanted to turn this into a more 
web, web friendly format. 

22
00:01:27,413 --> 00:01:30,997
And so they had to convert the images to 
a web friendly format and they also had 

23
00:01:30,997 --> 00:01:36,019
to convert the extracted text into a 
little package of JavaScript code. 

24
00:01:36,019 --> 00:01:37,589
Okay? 
So they're getting, get read to put this 

25
00:01:37,589 --> 00:01:40,391
stuff on the web. 
So this is four hundred and five thousand 

26
00:01:40,391 --> 00:01:44,081
images, which was quite a bit, especially 
at the time. 

27
00:01:44,081 --> 00:01:47,375
Okay, but the schematic looks sort of 
similar, right, you take a big set of 

28
00:01:47,375 --> 00:01:50,345
TIFF images and you split them into 
chunks and put them on a bunch of 

29
00:01:50,345 --> 00:01:54,571
different computers. 
and you have a function f that converts a 

30
00:01:54,571 --> 00:01:57,177
TIFF to a PNG, and does the other work 
too, let's say. 

31
00:01:57,177 --> 00:02:01,425
And what you get out is the dataset you 
want, a bunch of PNG images, okay? 

32
00:02:01,425 --> 00:02:05,194
And they're distributed across these, 
these machines. 

33
00:02:05,194 --> 00:02:10,168
Right, so let's look at another example. 
All right, so now we want to run 

34
00:02:10,168 --> 00:02:13,463
thousands of little simulations. 
And what we have are the parameters to 

35
00:02:13,463 --> 00:02:15,171
each one of those thousands of 
simulations. 

36
00:02:15,171 --> 00:02:18,038
And so, and example of this at the URL 
here at the bottom of the slide is from 

37
00:02:18,038 --> 00:02:20,999
simulating muscle dynamics, and this 
comes up a lot, if you have these Monte 

38
00:02:20,999 --> 00:02:24,007
Carlo simulations, they need to do they, 
they understand sort of phenomenon 

39
00:02:24,007 --> 00:02:27,156
stochastically by running lots and lots 
and lots of simulations with different 

40
00:02:27,156 --> 00:02:33,155
kind of, different inputs, and then kind 
of averaging the results. 

41
00:02:33,155 --> 00:02:35,899
Okay. 
As to modeling everything precisely. 

42
00:02:35,899 --> 00:02:37,929
So now you want to run thousands of 
simulations. 

43
00:02:37,929 --> 00:02:41,570
Well, you have a set of inputs, of 
parameters, to these simulations. 

44
00:02:41,570 --> 00:02:45,068
And you break them into chunks and put 
them all in separate machines. 

45
00:02:45,068 --> 00:02:47,446
And apply the function, and here the 
function is actually running the 

46
00:02:47,446 --> 00:02:50,601
simulation. 
And what you get out is the output of the 

47
00:02:50,601 --> 00:02:54,371
stimulation distributed across all the 
machines. 

48
00:02:54,371 --> 00:02:57,500
Okay. 
So there's, you know, a pattern should be 

49
00:02:57,500 --> 00:03:01,245
emerging here, right? 
So another example, so imagine each one 

50
00:03:01,245 --> 00:03:05,770
of these little bars is a document. 
And your task is just to find the most 

51
00:03:05,770 --> 00:03:09,745
common word in every individual document. 
Okay. 

52
00:03:09,745 --> 00:03:14,773
Well, same thing, distribute the 
documents across the k computers. 

53
00:03:14,773 --> 00:03:19,723
And then your function f now in this case 
opens up a single document, figures out 

54
00:03:19,723 --> 00:03:24,373
which word is the most common in that 
document, and then just produces that 

55
00:03:24,373 --> 00:03:28,880
word. 
And so now you have a big distributed 

56
00:03:28,880 --> 00:03:32,336
list of pairs where, you know, the first 
part of the pair is the document ID and 

57
00:03:32,336 --> 00:03:38,512
the second one is the word. 
Okay, so that could be useful but it's a 

58
00:03:38,512 --> 00:03:42,625
bit contrived. 
You know, consider a slightly more 

59
00:03:42,625 --> 00:03:46,310
general program that computes the word 
frequency of every word still in a single 

60
00:03:46,310 --> 00:03:48,424
document. 
Right? 

61
00:03:48,424 --> 00:03:51,496
So instead of just finding the most 
common one and producing that, now you're 

62
00:03:51,496 --> 00:03:54,664
going to produce [INAUDIBLE] histogram of 
the frequencies of every word in the 

63
00:03:54,664 --> 00:04:00,042
document. 
Okay, so given this Input, you produce 

64
00:04:00,042 --> 00:04:03,340
all of these items. 
Right? 

65
00:04:03,340 --> 00:04:05,420
A set of items. 
And the only reason I'm making this 

66
00:04:05,420 --> 00:04:08,470
distinction from the last one is that the 
last one, you know, took a single 

67
00:04:08,470 --> 00:04:13,740
document to produce a single word. 
And now we're taking a single [INAUDIBLE] 

68
00:04:13,740 --> 00:04:17,384
just want to make that clear that that's 
allowed. 

69
00:04:17,384 --> 00:04:20,222
Okay. 
Looks like the animation isn't here, but 

70
00:04:20,222 --> 00:04:24,658
I don't think I'll wind up fixing it. 
So you have millions of documents, you 

71
00:04:24,658 --> 00:04:28,879
distribute them again, now your function 
returns a set of word frequency pairs, 

72
00:04:28,879 --> 00:04:32,653
right? 
But that's okay, and now we have, you 

73
00:04:32,653 --> 00:04:36,729
know, lots of little lines here. 
[NOISE] a single word, let's say. 

74
00:04:36,729 --> 00:04:39,850
So they're not one to one any more. 
But that's no problem the function just 

75
00:04:39,850 --> 00:04:45,109
returns a set of things. 
So there should be a pattern here, right? 

76
00:04:45,109 --> 00:04:47,592
We have a function that maps a read to a 
trimmed read. 

77
00:04:47,592 --> 00:04:50,716
We have a function that maps a TIFF image 
to a PNG image. 

78
00:04:50,716 --> 00:04:54,533
A function that maps a set of parameters 
to the simulation result. 

79
00:04:54,533 --> 00:04:58,356
Alright, it's the simulation itself. 
We have a function that maps a document 

80
00:04:58,356 --> 00:05:02,028
to its most common word. 
And we have a function that maps a 

81
00:05:02,028 --> 00:05:06,118
document to the histogram of its word 
frequencies. 

82
00:05:06,118 --> 00:05:08,336
Okay? 
So, so good. 

83
00:05:08,336 --> 00:05:13,755
So these kinds of tasks we think we know 
how to do in parallel. 

84
00:05:13,755 --> 00:05:17,239
Given a big set of objects and a function 
that knows how to process a single object 

85
00:05:17,239 --> 00:05:21,646
you know, you should be able to think 
about how to paralyze this. 

86
00:05:21,646 --> 00:05:25,820
[INAUDIBLE] but we should abstractly 
understand how this is done. 

87
00:05:25,820 --> 00:05:28,445
Right? 
So, I'll say that one more time. 

88
00:05:28,445 --> 00:05:34,583
We have a big set of objects, we have a 
function Maps a single object [INAUDIBLE] 

89
00:05:34,583 --> 00:05:37,725
computers. 
Okay? 

90
00:05:37,725 --> 00:05:43,051
The objects among the computers and 
function that can parallel. 

91
00:05:43,051 --> 00:05:44,976
This is trivial. 
Alright. 

92
00:05:44,976 --> 00:05:49,656
Okay, so what if we want to compute the 
word frequency across all documents not 

93
00:05:49,656 --> 00:05:54,696
just uh, [INAUDIBLE] frequencies for each 
document [INAUDIBLE] frequencies for a 

94
00:05:54,696 --> 00:05:59,162
single document. 
Okay, so here you know, if we have one of 

95
00:05:59,162 --> 00:06:01,827
these three documents, now we want to get 
a single histogram that counts out the 

96
00:06:01,827 --> 00:06:04,410
number of times the word people appears 
across all three of them, the number 

97
00:06:04,410 --> 00:06:08,850
times the word, government appears across 
all three of them and so on. 

98
00:06:08,850 --> 00:06:13,457
So, let's go back to our schematic here, 
the pattern. 

99
00:06:13,457 --> 00:06:18,014
Well, now we want to compute the word, 
frequency, across five million documents. 

100
00:06:18,014 --> 00:06:21,853
And we can still distribute them among k 
computers, you know. 

101
00:06:21,853 --> 00:06:25,484
So far, so good. 
And then for each document we return a 

102
00:06:25,484 --> 00:06:30,038
set of word frequency pairs and now I've 
switched the notation here from F to map 

103
00:06:30,038 --> 00:06:36,990
since we can consider this a map, that's 
the, the, the terminology I used. 

104
00:06:36,990 --> 00:06:41,944
Okay, but now what do we do? 
So what we can get out here, what we will 

105
00:06:41,944 --> 00:06:44,623
get out here is a set of Frequency but 
that's not what we want a, one big 

106
00:06:44,623 --> 00:06:49,130
histogram. 
And to build this one big histogram we 

107
00:06:49,130 --> 00:06:54,330
have to make sure that a single computer 
has access to every occurrence of some 

108
00:06:54,330 --> 00:07:01,512
particular term, so. 
[SOUND] If the word history appears in 

109
00:07:01,512 --> 00:07:04,760
some document on this machine, and it 
appears, you know, twice in this 

110
00:07:04,760 --> 00:07:08,456
document, and three times, and two, two 
times in documents on that machine, and 

111
00:07:08,456 --> 00:07:11,984
so on, we have to sort of group those all 
up and send them to, to a single place, 

112
00:07:11,984 --> 00:07:17,890
just so we can count them. 
Okay. 

113
00:07:17,890 --> 00:07:25,001
So, let's look at this again. 
So, we distribute the documents across 

114
00:07:25,001 --> 00:07:28,100
these computers. 
We map apply our map function to each 

115
00:07:28,100 --> 00:07:32,175
document in order to, to produce a set of 
word frequency pairs. 

116
00:07:32,175 --> 00:07:36,981
And now we have a big distributed list of 
these words, sets of word frequencies. 

117
00:07:36,981 --> 00:07:40,811
And now we want to get. 
These workers involved in the process. 

118
00:07:40,811 --> 00:07:45,163
And these guys are going to be the ones 
who count the occurrences of a particular 

119
00:07:45,163 --> 00:07:48,352
word. 
Okay, and so imagine all these little 

120
00:07:48,352 --> 00:07:51,628
colored red lines are occurrences of 
words, such that all the red lines are, 

121
00:07:51,628 --> 00:07:55,060
represent a, a single word occurrences of 
a single word, and all the green lines 

122
00:07:55,060 --> 00:07:59,627
represent a different word, and so on. 
On. 

123
00:07:59,627 --> 00:08:08,055
Well so, these guys are going to sent, to 
their respective locations. 

124
00:08:08,055 --> 00:08:14,071
So, such that, this worker is in charge 
of handling all the occurrences of the 

125
00:08:14,071 --> 00:08:18,716
blue word. 
And this worker is in charge of all the 

126
00:08:18,716 --> 00:08:24,399
occurrences of the red work, and so on. 
Okay, so now instead of lines that go, 

127
00:08:24,399 --> 00:08:28,824
you know, from 1 to 1 we have lines that 
go from this 1 computer to a bunch of 

128
00:08:28,824 --> 00:08:32,921
different computers. 
Okay. 

129
00:08:32,921 --> 00:08:34,106
And so on. 
Alright? 

130
00:08:34,106 --> 00:08:38,714
They just sort of shuffle the data, they 
have to spray this data out across the 

131
00:08:38,714 --> 00:08:43,880
network, in order to regroup it. 
Fine, so now that we have the data 

132
00:08:43,880 --> 00:08:47,228
grouped the way we want and partitioned 
the right way, we can apply another 

133
00:08:47,228 --> 00:08:50,630
function which I'll call the reduced 
function which in this case it doesn't 

134
00:08:50,630 --> 00:08:53,978
really assemble it just counts them and 
that allows us to produce our final 

135
00:08:53,978 --> 00:09:00,855
result which is oh, there are 4 green 
words and 4 red words and 3 blue words. 

136
00:09:00,855 --> 00:09:06,443
Words and so on. 
Okay, so, now the schematic looks a 

137
00:09:06,443 --> 00:09:11,837
little different. 
We have sort of a two-step process. 

138
00:09:11,837 --> 00:09:16,842
So, we start with [SOUND] some large set 
of objects distributed over a bunch of 

139
00:09:16,842 --> 00:09:20,579
machines. 
And then we want to apply some function F 

140
00:09:20,579 --> 00:09:24,267
to each one of those objects. 
And that was our First step. 

141
00:09:24,267 --> 00:09:28,359
But then the output of those functions 
are all going to be redistributed across 

142
00:09:28,359 --> 00:09:31,861
the network and grouped. 
To form groups. 

143
00:09:31,861 --> 00:09:36,275
Okay, and then the second step is to 
process each one of those groups. 

144
00:09:36,275 --> 00:09:39,804
So I'll write, instead of f, I'll write 
map here. 

145
00:09:39,804 --> 00:09:44,776
And I'll write reduce here. 
Okay, and so that's exactly what Map 

146
00:09:44,776 --> 00:09:48,493
Reduce does and we'll explain this in 
more detail next time, but the key idea 

147
00:09:48,493 --> 00:09:52,446
here is that the user, the programmer is 
going to write these two functions, a map 

148
00:09:52,446 --> 00:09:56,458
function and a reduce function, which are 
serial and what I mean by serial is there 

149
00:09:56,458 --> 00:10:02,428
not parallel. 
You don't have to worry about how to 

150
00:10:02,428 --> 00:10:05,392
manage, how to program distributed 
cluster within each one of these 

151
00:10:05,392 --> 00:10:10,296
functions. 
You just write a function map that takes 

152
00:10:10,296 --> 00:10:16,496
in an object of some kind and returns 
some other object. 

153
00:10:16,496 --> 00:10:21,324
And I'm not using objects in this sort of 
object-oriented sense, this is sort of in 

154
00:10:21,324 --> 00:10:27,991
the mathematical sense, just sort of any 
input to produce any kind of output. 

155
00:10:27,991 --> 00:10:34,111
Okay, and then the reduced function takes 
a set of objects which I'll note this way 

156
00:10:34,111 --> 00:10:39,871
and returns some of the kind of object 
and actually this can be a set of objects 

157
00:10:39,871 --> 00:10:45,621
as well. 
And in fact this can be a set of objects 

158
00:10:45,621 --> 00:10:48,532
as well. 
So maybe I'll switch colors here. 

159
00:10:48,532 --> 00:10:53,554
This can actually be a set of things and 
this can actually be a set of things. 

160
00:10:53,554 --> 00:10:59,805
So for example we saw a document 
returning a set of word frequency pairs. 

161
00:10:59,805 --> 00:11:03,526
But, you can think of it. 
It's, it's essentially important to 

162
00:11:03,526 --> 00:11:07,307
understand what's going on here, okay. 
So this is, this is, this is a, an 

163
00:11:07,307 --> 00:11:12,190
interesting hypothesis, right? 
That perhaps all distributed algorithms 

164
00:11:12,190 --> 00:11:16,365
can be expressed as sequences of these 
two step operations, right? 

165
00:11:16,365 --> 00:11:21,352
A map fall by reduce, and then maybe more 
map reduce, map reduce, map reduce. 

166
00:11:21,352 --> 00:11:24,666
As, as needed. 
Okay, and we'll talk about that in more 

167
00:11:24,666 --> 00:11:26,370
detail next time. 

