[CS61C FA20] Lecture 36.3 - MapReduce, Spark: MapReduce — Transcript
Full transcript
- 0:00and welcome back we're finally here
- 0:02mapreduce
- 0:04very exciting to be able the first one
- 0:05to to explain this to you if you haven't
- 0:06seen it before
- 0:08so mapreduce is an abstraction it's a
- 0:10very simple
- 0:12i i say simple compared to what it used
- 0:14to have been
- 0:15we always have to do before map produced
- 0:16days data parallel programming model
- 0:19designed for scalability and fault
- 0:21tolerance the key here is fault
- 0:22tolerance if you're operating on
- 0:24not just cores cores don't fail for the
- 0:26most part
- 0:27but boy when you go to machines and
- 0:29those machines live in other countries
- 0:30they fail the network goes down things
- 0:32happen i can't just farm it out to a
- 0:33thousand thing a thousand different
- 0:35machines and hope that all thousands
- 0:36and expect i hope they'll a thousand but
- 0:38i expect that all thousands will come
- 0:39back with the
- 0:40with with with no trouble without
- 0:42failing at some level many things could
- 0:44happen
- 0:45especially when you increase the number
- 0:46of machines to a big number you're going
- 0:48to have some problems
- 0:49not usually the same in the case of core
- 0:51so there's a little bit less worry about
- 0:52fault tolerance there but
- 0:53certainly the case when you go to
- 0:54multi-machine and distribute call this
- 0:55distributed computing by the way
- 0:57again pioneered by google um and
- 1:00regularly they're processing
- 1:0125 petabytes pretty incredible right
- 1:04petabytes of data per day
- 1:06um there's a new open source project i
- 1:08would say so new it was there was an
- 1:09open source party that was new a couple
- 1:11of years ago
- 1:12hadoop used by yahoo facebook amazon and
- 1:15many people so
- 1:16and this is a java based framework we we
- 1:18love hadoop
- 1:19for what it does what's it used for well
- 1:21google google and
- 1:23google uses its mapreduce framework for
- 1:25lots and lots of things
- 1:27you got to build your indices and your
- 1:28to be able to handle google search
- 1:31clustering articles for google news
- 1:33machine translation
- 1:35street maps how do you do all that
- 1:36street maps and multi-layer and and be
- 1:38able to handle the zoom factor and all
- 1:40of those
- 1:40um yahoo search for yahoo spam detection
- 1:44is a big one
- 1:45um facebook data mining ad optimization
- 1:47spam detection
- 1:48many things that fall into the model of
- 1:50as you can imagine a mapping phase and a
- 1:52reduction phase and we'll talk about
- 1:54what those are in a moment
- 1:55here's an example of something that they
- 1:56use they used to use mapreduce for at
- 1:58facebook this is now
- 1:59this used to be available no longer
- 2:00available but the facebook lexicon
- 2:02i think google has something as well
- 2:04where you can type in a word and you can
- 2:06see
- 2:06uh when correlating events used to
- 2:10happen so when are people searching for
- 2:11different things
- 2:12so here's a funny thing of um people
- 2:15saying party town had a hangover and so
- 2:16here's party tonight in hang our party
- 2:18party tonight is a uh is in yellow and
- 2:21hanging over some blue and you notice
- 2:22the hangovers often follow
- 2:24often peak right behind the party
- 2:25tonight it's just people typing things
- 2:27but you know
- 2:27when is this happening and this is like
- 2:29the frequency that and by the way one of
- 2:31the parties happening can you see
- 2:32oh there's a pretty big party when oh i
- 2:34don't know halloween party
- 2:35new year's eve party right maybe holly
- 2:38you know holidays
- 2:39christmas kwanzaa hanukkah all happening
- 2:41here but kind of distributed so there
- 2:42but boy
- 2:43everybody seems to party on halloween
- 2:45and new year's
- 2:46and not not anything in february march
- 2:48or april so it sounds like we need to
- 2:50introduce more
- 2:51national holidays in those days just
- 2:52distribute out the parties
- 2:54um here's the design goal the design
- 2:56goal by the way
- 2:57this initially came out of work that uh
- 3:00jeffrey dean and sanjay
- 3:01gemmawatt did at google years ago they
- 3:03wrote a famous paper one of the most
- 3:05the most downloaded paper in systems um
- 3:07since then since 2004. this is like 16
- 3:10years ago as of this this video the idea
- 3:12is it's scalable to
- 3:14really massive uh data volumes thousands
- 3:17of machines each machine multi
- 3:18must might have multi-core um within
- 3:21them
- 3:22tens of thousands of disks that's a lot
- 3:24of data you can imagine processing in
- 3:25all of that
- 3:26um cost efficiency so one of the things
- 3:28they've learned we'll talk about this in
- 3:29the warehouse scale computing lecture
- 3:31is they didn't buy the highest end
- 3:33machines they went with commodity
- 3:34machines they said
- 3:35you know value per dollar right cycles
- 3:39per dollar processing cycles per dollar
- 3:40is cheaper if you go with commodity
- 3:42machines and kind of let the
- 3:44the value of volume uh really helps
- 3:46rather than well let's buy the highest
- 3:47end
- 3:48like kind of the gamer crazy high-end
- 3:49gamer machine with
- 3:51liquid cooling and the high that you
- 3:54know
- 3:54per dollar that actually doesn't isn't
- 3:56so if you don't have one game i want to
- 3:57play for one system i'll have this the
- 3:58best thing for my 8k
- 4:00screen maybe that's a good value but not
- 4:02when you try to compute
- 4:03massive bits of data that's compute
- 4:05compute processing per dollar that's not
- 4:07the right call
- 4:07a commodity machine is very interesting
- 4:10and they actually started building their
- 4:10own which is very interesting
- 4:12commodity network um here's the thing
- 4:15fault tolerance how do you do with fault
- 4:16tolerance well all you do is you resend
- 4:18it this is what happens on the internet
- 4:19as well if i send a packet i don't get
- 4:21an act back i just send it again
- 4:22that's what we learned in the internet
- 4:24um and that's what tcpip is doing for
- 4:26you
- 4:27tcp really is doing that word for you um
- 4:29here what we'll just do is we'll resend
- 4:31that request up and probably not to the
- 4:33same machine probably send it to
- 4:34somebody else
- 4:35um but we handle fault tolerance by just
- 4:37re-execution just send it out again send
- 4:39the request out i broke it into 1000
- 4:41pieces
- 4:41and 999 came back one didn't well i'll
- 4:43send it somebody else and by the way
- 4:44maybe there's a timeout earlier than
- 4:46that rather than wait till not 99.99
- 4:48came back
- 4:48and then send the thousands back out
- 4:50maybe like say you know what i'm going
- 4:51to give up on you early send it out so
- 4:52that it can all come back at the same
- 4:53time i have to wait for that
- 4:55and we find it very easy and fun to use
- 4:56we've had some assignments in in lower
- 4:58division to do this and we're excited to
- 4:59be able to share this with you
- 5:00as well um we've linked to the dean
- 5:03paper we called the
- 5:04the demon gimbal paper we've linked that
- 5:06on the 621c page
- 5:07so please do check that out here's a
- 5:09typical cluster
- 5:11you've got each machine we call them
- 5:12nodes each node might have a disk local
- 5:14to it
- 5:15hanging right off of it and that's a
- 5:17reasonable thing to do there's a rack
- 5:19switch that connects these all we'll
- 5:20talk about
- 5:20server rack and array in the next
- 5:22lecture um but for here
- 5:24they're all connected via a switch which
- 5:26is able to take a request and be able to
- 5:27know which computer to talk to and which
- 5:28it came from where it's going kind of
- 5:30like a router in some sense
- 5:31um and there's an aggregation switch
- 5:33which is doing the same thing
- 5:34but the difference is you have a
- 5:36different allow different bandwidth
- 5:37required requirements
- 5:39once you get to um more and more kind of
- 5:42central command it's almost like
- 5:43plumbing right when you have a plumbing
- 5:44from your house the pipe's really small
- 5:46once it gets to your neighborhood it's
- 5:47much bigger it's from the city it's
- 5:49massive okay
- 5:50you go to los angeles they have this
- 5:52overflow not this i'm sure the sewer's
- 5:53larger
- 5:54massive as well but you know the if you
- 5:56i'll think about uh handling flooding
- 5:57you have a little tiny flood
- 5:59this happens a lot in the bay area and
- 6:01very flat location certainly in los
- 6:02angeles
- 6:03um you know oh here's my neighborhood
- 6:05there's a little
- 6:06ditch to handle flood and that's fine
- 6:09and then once it gets to the main area
- 6:11there's this
- 6:12massive i think it's probably like 20
- 6:15lanes wide of a flow where the water is
- 6:18finally going to the ocean
- 6:19same idea here once you get away from
- 6:22each node
- 6:23higher up in the in the in the hierarchy
- 6:25if you will these switches have to be
- 6:27very fast much higher end
- 6:29and they have an eight gigabit
- 6:30connection network here across that so
- 6:32you end up having much more expensive
- 6:34and we'll talk about the cost here
- 6:35uh much more expensive switches and
- 6:37connections at the higher level because
- 6:39you're
- 6:39you're handling larger channels of data
- 6:41this is for everybody in the whole
- 6:42in the whole maybe the whole array so
- 6:45some of the numbers we have here
- 6:47are 40 nodes in a rack uh maybe up to
- 6:504 000 nodes in a cluster so you might
- 6:53hear that here the first cluster here
- 6:55um and again one gigabit per second
- 6:57bandwidth
- 6:58within the rack and then eight out of
- 6:59the rack um and the node specs this this
- 7:01is a lower-end machine this is like a
- 7:03couple years old but
- 7:04better than not that not too much faster
- 7:05computers as you saw aren't getting that
- 7:06much faster
- 7:07so maybe eight two giga two gigahertz
- 7:10cores
- 7:11eight gigs of ram probably that's 16 or
- 7:1332 by now and four disks each about four
- 7:15terabytes
- 7:16okay so this is a fun
- 7:20slide because one of the things we've
- 7:21been doing in cs10
- 7:23and 61a a little bit in 621 people not
- 7:26as much
- 7:27is to share with you the benefits of the
- 7:29programming
- 7:30paradigm known as the functional
- 7:32programming paradigm if you
- 7:33think in a functional way where your job
- 7:35is to have a block that has an inputs
- 7:37and outputs and doesn't have any side
- 7:39effect
- 7:39then you're going to be able to work
- 7:40with these this mapreduce model a little
- 7:43bit easier
- 7:44so we've been trying to teach mapreduce
- 7:45we actually teach mapreduce the idea of
- 7:47a map and a reduce
- 7:49in cs in a non-majors class very rarely
- 7:51is that happening so i think berkeley's
- 7:52really leading the way nationally
- 7:54in terms of bringing some of these
- 7:55parallel and distributed computing ideas
- 7:57into the lower division and
- 7:58even in this case into our our
- 8:01non-majors class
- 8:02which by the way is being shared with
- 8:03you know a thousand high schools
- 8:05all around the world at least we've had
- 8:07a thousand teachers go through our
- 8:08professional development in bjc
- 8:09so that means these ideas are getting
- 8:12out there into the into the winds
- 8:15here's the model map map says do a
- 8:18function
- 8:19over a set of data that's it so i want
- 8:22to square all those numbers
- 8:231 23 and 10. so here's a little picture
- 8:27on the bottom here to show you what's
- 8:28happening
- 8:28so i'm first applying the map but what
- 8:30that does is map says take a function
- 8:32it's usually a monotic function a
- 8:34function of one argument here in this
- 8:35simple case as a function of one
- 8:37argument
- 8:38takes every number and it's going to
- 8:39square it so it's applying that function
- 8:40to everybody returning a new list
- 8:42not not mutating the list it's not it's
- 8:43an immutable list returning a new list
- 8:45in which every element has had that
- 8:47function applied to it so if i square
- 8:49every element of the list 123 and 10 i
- 8:51get 1 409 and a hundred
- 8:55and now i say combine with plus and how
- 8:57combined with again
- 8:58it's below the level abstraction of how
- 9:00that works is it
- 9:02left associated right associative tree
- 9:04wise hopefully it's
- 9:05a tree structure but that's being
- 9:07combined and that's usually a dyadic
- 9:09function
- 9:09two inputs one output and that's kind of
- 9:11a smoosher together and so here i'm just
- 9:13adding them up
- 9:13so here's here's an example here's one
- 9:15and four hundred get added together to
- 9:17give
- 9:17to give you four or one not in a hundred
- 9:19it gives you 109 and then eventually
- 9:20that's five ten
- 9:21so that's kind of a we're hoping like a
- 9:23tree model for that
- 9:25who knows who knows how reduction
- 9:27happens again who knows how the mapping
- 9:28happens
- 9:29maybe you tell me oh dan i only have two
- 9:31worker bees two
- 9:32cores two machines to help you well i
- 9:34don't care this is just gonna
- 9:36work if i have two this will be split up
- 9:38you know one machine might do these two
- 9:39one machine might do these two
- 9:40i don't know maybe there's four this
- 9:42list could also be in that you know in
- 9:44the building who knows how big that list
- 9:45is
- 9:46so you know rarely gonna have a billion
- 9:48worker bees working
- 9:49cores total cores in your in your
- 9:51available set so at some point it says
- 9:53well apply it out and then come back and
- 9:54then keep working on more of that until
- 9:56you're all done but from the user's
- 9:57point of view i don't care
- 9:58this is the beauty of this abstraction
- 9:59from the user's point of view i don't
- 10:01care
- 10:02so the end of the day i get 5 10. okay
- 10:06by the way there's only two data types
- 10:07here the data type of the input and the
- 10:09data type of the output here and this is
- 10:10the
- 10:11thinking about domain and range here
- 10:14this can be done in scheme we used to
- 10:15teach this in scheme all the time this
- 10:17is a scheme model those of you taking 6
- 10:18to 1 a
- 10:18know this and so map and reduce come
- 10:22by the way we call it combined it's we
- 10:23call it map it combine
- 10:25in cs10 and 61a we teach scheme and map
- 10:28and reduce are part of that and so
- 10:30here's map square over the list 123 and
- 10:3210 and then you reduce plus over that
- 10:34how beautiful and
- 10:35tight that is python it's a little bit
- 10:37more syntax for it but not much more
- 10:40we have to import the reduce method from
- 10:42func tools
- 10:43library then i have to define the plus
- 10:46function
- 10:47a little bit and i have to then i'm
- 10:49going to define a square function
- 10:51i can obviously use a lambda here and
- 10:52you'll see it a little later we're going
- 10:53to use the lambda due to spark stuff but
- 10:55here's exactly the same thing in python
- 10:57reduce of plus over map of square of
- 10:59this and automatically i get 510
- 11:02really really clean so this is a
- 11:04simplified version of mapreduce
- 11:06it gets a little more complicated on the
- 11:07next slide so let's actually get into it
- 11:10in the full mapreduce programming model
- 11:13that dean
- 11:15and dean at all presented is we have
- 11:18a intermediate step that's the step
- 11:21between the map and the reduction phase
- 11:24and that intermediate step has
- 11:27a key value pairing you know by key
- 11:29value pairing from python and from from
- 11:31any key value model that you've seen
- 11:32before
- 11:33so here's how map works map takes
- 11:39map takes your data a key and a value
- 11:43and it's going to process the input key
- 11:45and value pair
- 11:47slicing data into shards which are
- 11:49distributed to workers
- 11:51and that produces intermediate pairs
- 11:54so the scenario pairs are set up by key
- 11:57and then
- 11:58they are sorted by key every shared key
- 12:01all the data for a shared key is sent to
- 12:03one machine you're going to see a data
- 12:04you see a
- 12:05slide the next slide this will all
- 12:06become clear that machine
- 12:09then i'm the machine has all data at
- 12:11that same key so the key is the way that
- 12:12you're
- 12:13attaching you're tagging the data in
- 12:15some way to distribute it across those
- 12:17machines
- 12:18because the system says whatever the
- 12:20keys you have i will
- 12:21coordinate it sorted by the keys so the
- 12:24intermediate stage is there so you're
- 12:25using the keys in some sense
- 12:27not to just add more data not to add
- 12:29more craft to your
- 12:31input data but to distribute it across
- 12:33nicely you get to determine what the key
- 12:35assignment is
- 12:36and that distributes it nicely so you
- 12:37get some control over how it distributes
- 12:39in the reduction phase
- 12:40and then so everybody just attaches the
- 12:43key in a parallel way
- 12:44then it's sorted by key across machines
- 12:47and then every machine gets all the data
- 12:49for a particular key
- 12:50and then processes all the data using
- 12:52your reducer and there's a reducer
- 12:53okay that's it
- 12:56and it produces a pre a set of merged
- 12:59output usually just
- 13:00usually just one value is over here so
- 13:03here's an example
- 13:06the bread and butter problem we have
- 13:08every time when we're looking at
- 13:10a new programming language is hello
- 13:12world if you have say
- 13:14my thing my my system solves uh all
- 13:17two-person abstract strategy games we'd
- 13:19say solve tic-tac-toe
- 13:21if you were to say i have this great
- 13:23system that lets you do parallel
- 13:24computing easier
- 13:25we say to you show me word count word
- 13:28count has become the problem we go to is
- 13:30the kind of go-to
- 13:31bread and butter problem the hello world
- 13:33if you will of
- 13:35program parallel and distributed systems
- 13:38word count says the following i've got a
- 13:40corpus of text
- 13:41uh maybe it's all shakespeare's work
- 13:42maybe it's the gutenberg bible maybe
- 13:44it's
- 13:44all the hillary emails who knows what it
- 13:47is and i want to be able to
- 13:48and for word count count the number of
- 13:51words
- 13:52for every word for every unique word i
- 13:54want to count the number
- 13:55the occurrences of it that's all i'm
- 13:57doing so if i give you
- 13:59i do i learn i want to get out
- 14:03i there's two eyes one do and one learn
- 14:06here
- 14:06and so what this is going to do is this
- 14:08is going to tag
- 14:10remember i told you the output of map is
- 14:11this key value pair list
- 14:13so what it's going to do is it's going
- 14:15to tag it with the key is the word
- 14:17itself
- 14:18split it by words the word itself and
- 14:20the value is going to be the number one
- 14:22so essentially how i'm taking for every
- 14:23word that comes in i just slap on a one
- 14:25to the side of it and the key by the way
- 14:27is going to be the word
- 14:28and the value is going to be the number
- 14:29one and pass it on to the next stage
- 14:31then that sorting stage is going to sort
- 14:34things there's going to be a machine for
- 14:35each of those keys
- 14:36so the keys are i do and learn it's
- 14:38going to be a machine whose job is to
- 14:40just add up those
- 14:41values here the number one for all the
- 14:43eyes there's a machine
- 14:44adding up all the values here's the
- 14:47number one again
- 14:48for the do and here's one for learn so
- 14:50be a i a do
- 14:51and a learned machine and each of them
- 14:53does the red red reducer
- 14:55what's the reducer say well look at this
- 14:58set result to one for each value for
- 15:01each v
- 15:02in the intermediate values of the
- 15:03intermediate level intermediate values
- 15:05i accumulate in some result in some
- 15:08result variable
- 15:09and i emit the result as a string that's
- 15:11it
- 15:13the data is on a distributed file system
- 15:15this means the data isn't in memory the
- 15:16data is on up files on files
- 15:19and those files who knows where they are
- 15:21they're on the distributed network
- 15:22so they're not just a local to me you
- 15:24think oh it's always my local desk it's
- 15:25not the case you're on a distributed
- 15:27file system
- 15:27which means that they they live out
- 15:29there in kind of i would say that
- 15:30the cloud in quotes so this is my
- 15:33favorite example here's the code for map
- 15:35and reduce
- 15:36and i've got um i've got a file of just
- 15:39utterances
- 15:41a-a-a-r is in the first file file two is
- 15:43blank
- 15:44r is in the third file if and or are in
- 15:47the fourth file
- 15:49or in r the fifth file or and on if
- 15:53what does the mapper do it slaps a key
- 15:56and a vowel it slaps you know it returns
- 15:58a key and a value pair for each of these
- 15:59guys
- 16:00so it says i'm going to process this by
- 16:03the way there's no reason it has to
- 16:04process everything i could
- 16:05in the mapper it could go through and
- 16:06decide to not ex not emit some things
- 16:09what it does is see watch right what i'm
- 16:11saying is for each word
- 16:13it emits that you could say for each
- 16:16word that is
- 16:17not er and therefore you take error so
- 16:20you can do some filtering in there as
- 16:21well
- 16:21it's your call what i mean you're
- 16:23writing this code my point is you're
- 16:24only emitting things you want to emit
- 16:26here it's meaning everything because
- 16:27it's for each
- 16:28w and w's emitting everything but you
- 16:30could imagine there being a filtering
- 16:31step in here if you wanted that so think
- 16:32about as you're thinking oh i might do
- 16:34this problem but i need to do some
- 16:35filtering
- 16:36here's where the filtering would happen
- 16:37where you don't you don't admit things
- 16:38that you don't want to emit
- 16:40emit so now all i've done is attached
- 16:42one to all these guys now i have a one
- 16:44at the end of every single key the key
- 16:46again
- 16:46are all the words this is this phase i
- 16:49talked about the group by key phase
- 16:51that says they're going to be a machine
- 16:53attached to every key the keys are now
- 16:55the five unique words are
- 16:57if or or uh and
- 17:01what is a reducer then all those values
- 17:03all those values are sent to the reducer
- 17:05what does reducer say basically it says
- 17:07add them all up
- 17:08add all the numbers up whatever the
- 17:10values
- 17:11okay sometimes it ignores the keys i
- 17:13don't care what the output key is
- 17:15but for each of these guys all i'm going
- 17:16to do is output a four
- 17:18emits a string and so what comes out of
- 17:21this is
- 17:22in a way it's a dictionary and think
- 17:24about this as the
- 17:26the outputted value of the whole system
- 17:28here is
- 17:29every machine says well i was i was the
- 17:32machine
- 17:32assigned to ah and my value is four so
- 17:35it's ah colon four
- 17:37why was the machine doing er and i only
- 17:39saw one of them so i'm er colon one
- 17:41so that's the idea the output is in a
- 17:43way of this whole system
- 17:44of that was the goal the goal of the
- 17:46whole problem was word count was
- 17:48return a dictionary of where the keys
- 17:50are the words
- 17:51every unique word and the value is the
- 17:53number of occurrences they occurred in
- 17:54the original thing
- 17:55well that's what's happening here this
- 17:57is the output if you turn your head
- 17:58sideways
- 17:58this is kind of like key down here and
- 18:00value up there
- 18:02that's it and that's your output of the
- 18:03of word count very simple right it
- 18:05doesn't seem that hard
- 18:06so mapper all it does is slap a one at
- 18:07the end of it and reducer says well i
- 18:09got all the ones what do i do with them
- 18:10add them all up so it's really reduce
- 18:12plus okay
- 18:14key is attach one the the mapper the
- 18:16reducer is
- 18:18just take it's just plus that's all i'm
- 18:20doing that's it
- 18:21so that is word count pretty clean
- 18:25how does this work well
- 18:29let's say the corpus is not just you
- 18:31know four files but a ton of files maybe
- 18:33it's a huge directory maybe it's a lot
- 18:35how does this work well this is very
- 18:38clever this is the part that's the
- 18:40system engineering part of it that they
- 18:41really is
- 18:42quite brilliant and by the way not all
- 18:43maps take the same time sometimes
- 18:45you know the map task might have maybe
- 18:47the file you're given is a much larger
- 18:49file than everybody else so you're gonna
- 18:50be doing more work than everybody else
- 18:51does so not all maps take the same
- 18:52amount of time i kind of you know in
- 18:54in my simplification in my in my simple
- 18:56example of
- 18:57map of square you know square takes the
- 18:59same amount of time for every number you
- 19:00can be given
- 19:02except perhaps if you have big nums and
- 19:04the numbers actually are stored as a
- 19:05string but for for now you think that's
- 19:07that's a bad model really maps can take
- 19:09variable amount of time as reducers can
- 19:11take variable amounts of time
- 19:12who knows how many you know how many of
- 19:13those ones i had to add up in that model
- 19:16so there's a there's a lead we call it a
- 19:20i guess a master controller
- 19:22which assigns the map and reduced tasks
- 19:24to these worker servers there's all
- 19:25these servers it's like a you know
- 19:26worker bees if you will
- 19:28as soon as look at this as soon as a map
- 19:31task finishes
- 19:33that server can be re assigned a new
- 19:36task so this particular worker
- 19:38was only assigned map task this
- 19:39particular worker's only time map
- 19:42map tasks as well this worker can be
- 19:44assigned read tasks
- 19:45because it's reading one meaning what is
- 19:47this doing
- 19:48this guy says well when that's done i
- 19:51can start reading those values in and be
- 19:52waiting for it
- 19:53and i also need to know where i need to
- 19:55be reading this one so this is also
- 19:57reading for that when that comes in and
- 19:59then i'm reading for this one so when
- 20:00they're all done
- 20:02now i can start my reduction task so why
- 20:04wait why wait till they're all done
- 20:06it's very clever to think about this um
- 20:08to start its reduction
- 20:10so each worker and so this is this is
- 20:13saying look these workers are the mapper
- 20:14workers these are the reducer workers
- 20:16that's not always the case i'm basically
- 20:18once map
- 20:19map one finishes i'm telling the system
- 20:21hey i'm available for more work
- 20:23maybe i'm the job maybe this becomes me
- 20:25i mean it's drawn this way to make this
- 20:26diagram clear
- 20:27but it could be that this worker one is
- 20:29assigned
- 20:30to do that read and then it's going to
- 20:31be waiting and reading then basically
- 20:33this
- 20:34could be up here is what i'm saying okay
- 20:36it doesn't have to be like oh worker one
- 20:37is only mapping and worker two was only
- 20:39mapping ever you know this basically
- 20:40they go back into the pool of workers
- 20:42and the jobs are just out there so
- 20:43there's a queue of jobs that need to be
- 20:45processed like you gotta start the
- 20:47reducer you gotta start the next mapper
- 20:49um there may be many more mappers than
- 20:50just three here so
- 20:52it could be that um
- 20:57you're you're heterogeneous in that
- 20:59sense i'm a worker bee and i was doing
- 21:00some mapping and then i'll do some
- 21:01reducing and maybe some more mapping to
- 21:03be done i finished that reducing and
- 21:04there's more mapping to be done there
- 21:05who knows what that is
- 21:06um what's done and by the way what
- 21:09typically happens
- 21:10is it's not just one map and one
- 21:11reduction you tip it this is i
- 21:13i i read the paper i found what this
- 21:15what happens is there could be a map
- 21:16reduction here
- 21:17that feeds into another map reduction to
- 21:19feed into another so there's this whole
- 21:20series of
- 21:21a map reduction not just one it's just
- 21:23one problem it's all i have to know
- 21:24there might be several phases of it and
- 21:25doing phase one phase two phase three
- 21:27and you just continue to process the
- 21:28machine and it's really fun to watch
- 21:30there's a little um there's like a
- 21:33there's like a
- 21:34dashboard that you can see how things
- 21:35are processing it's very neat to see so
- 21:37i encourage you to take a look at that
- 21:38if you ever get this working
- 21:41the reduction task begins as soon as all
- 21:42the data shuffles finish maybe you gotta
- 21:44shuffle things around and finally there
- 21:45because
- 21:46in this problem as before that last
- 21:49mapper might contribute
- 21:50a one to your reduction so here in this
- 21:52particular problem
- 21:53you can't do any reduction until you've
- 21:55finished all the mapping
- 21:56so that i i might have implied that you
- 21:58can start reducing before the mapping is
- 22:00done
- 22:00in this particular case i certainly
- 22:02can't start adding up my ones until
- 22:03all the guys are gonna be ones feed into
- 22:05my machine and that's done
- 22:06so in this particular model you have to
- 22:08wait till all of the maps get done
- 22:10before i can even start thinking about
- 22:12the
- 22:12reduction but you can start reading
- 22:13their values certainly that's kind of a
- 22:15nice thing
- 22:17here's the thing i mentioned before uh
- 22:19if a worker doesn't respond
- 22:21um network failure disk failure hardware
- 22:24failure memory failure kind of could be
- 22:25a
- 22:25lot of failures i just reassign the task
- 22:28if a worker dies that's important
- 22:30this amount of java code that does that
- 22:32same word count if we look here
- 22:34in detail what's happening is here it's
- 22:35saying uh i'm outputting
- 22:38the value and a one this is the actual
- 22:40code for this and here i'm just here's
- 22:42my sum of zero and i'm sun plus equal
- 22:44that value
- 22:44a little bit more complicated but
- 22:46essentially that's what we were doing
- 22:47before and that's the java code for it
- 22:49that's it so there's word count folks
- 22:51and by the way most of this is just set
- 22:53up we're just setting these things up
- 22:54here in this early part of it
- 22:57that's it pretty exciting word count is
- 23:00a really nice example
- 23:01and map produces a very very powerful
- 23:03paradigm
- 23:04the next lecture we're going to talk
- 23:06about maybe there's a cleaner
- 23:08lighter weight way than when i'm looking
- 23:09at maybe 35 lines of code maybe there's
- 23:11a little tighter way of doing that in
- 23:13spark we'll see that then
- 23:14take care
About this transcript
This page contains the full transcript of [CS61C FA20] Lecture 36.3 - MapReduce, Spark: MapReduce by CS 61C Departmental, generated from the public captions YouTube serves with the video. The transcript has 5,070 words across 809 segments, with the original timestamps preserved so you can click any line to jump to that moment in the embedded player.
What you can do with it
Use the transcript to take notes, quote the speaker, build a study guide, generate a summary with ChatGPT or Claude via the YouTube Summary tool, or export it as a timed subtitle file with YouTube to SRT. You can also re-open it in the transcriber to translate the transcript into 100+ languages.
Free YouTube transcript tool
YouTube2Text is a free YouTube transcript generator — no signup, no daily limit. Paste any YouTube link and get the full transcript instantly, with timestamps, click-to-jump, translation to 100+ languages, AI prompts for ChatGPT, Claude, and Gemini, and exports to TXT, SRT, VTT, or Markdown.