YouTube2Text

[CS61C FA20] Lecture 36.3 - MapReduce, Spark: MapReduce — Transcript

by CS 61C Departmental · 5,070 words · 809 segments · language en · Watch on YouTube

Full transcript

  1. 0:00and welcome back we're finally here
  2. 0:02mapreduce
  3. 0:04very exciting to be able the first one
  4. 0:05to to explain this to you if you haven't
  5. 0:06seen it before
  6. 0:08so mapreduce is an abstraction it's a
  7. 0:10very simple
  8. 0:12i i say simple compared to what it used
  9. 0:14to have been
  10. 0:15we always have to do before map produced
  11. 0:16days data parallel programming model
  12. 0:19designed for scalability and fault
  13. 0:21tolerance the key here is fault
  14. 0:22tolerance if you're operating on
  15. 0:24not just cores cores don't fail for the
  16. 0:26most part
  17. 0:27but boy when you go to machines and
  18. 0:29those machines live in other countries
  19. 0:30they fail the network goes down things
  20. 0:32happen i can't just farm it out to a
  21. 0:33thousand thing a thousand different
  22. 0:35machines and hope that all thousands
  23. 0:36and expect i hope they'll a thousand but
  24. 0:38i expect that all thousands will come
  25. 0:39back with the
  26. 0:40with with with no trouble without
  27. 0:42failing at some level many things could
  28. 0:44happen
  29. 0:45especially when you increase the number
  30. 0:46of machines to a big number you're going
  31. 0:48to have some problems
  32. 0:49not usually the same in the case of core
  33. 0:51so there's a little bit less worry about
  34. 0:52fault tolerance there but
  35. 0:53certainly the case when you go to
  36. 0:54multi-machine and distribute call this
  37. 0:55distributed computing by the way
  38. 0:57again pioneered by google um and
  39. 1:00regularly they're processing
  40. 1:0125 petabytes pretty incredible right
  41. 1:04petabytes of data per day
  42. 1:06um there's a new open source project i
  43. 1:08would say so new it was there was an
  44. 1:09open source party that was new a couple
  45. 1:11of years ago
  46. 1:12hadoop used by yahoo facebook amazon and
  47. 1:15many people so
  48. 1:16and this is a java based framework we we
  49. 1:18love hadoop
  50. 1:19for what it does what's it used for well
  51. 1:21google google and
  52. 1:23google uses its mapreduce framework for
  53. 1:25lots and lots of things
  54. 1:27you got to build your indices and your
  55. 1:28to be able to handle google search
  56. 1:31clustering articles for google news
  57. 1:33machine translation
  58. 1:35street maps how do you do all that
  59. 1:36street maps and multi-layer and and be
  60. 1:38able to handle the zoom factor and all
  61. 1:40of those
  62. 1:40um yahoo search for yahoo spam detection
  63. 1:44is a big one
  64. 1:45um facebook data mining ad optimization
  65. 1:47spam detection
  66. 1:48many things that fall into the model of
  67. 1:50as you can imagine a mapping phase and a
  68. 1:52reduction phase and we'll talk about
  69. 1:54what those are in a moment
  70. 1:55here's an example of something that they
  71. 1:56use they used to use mapreduce for at
  72. 1:58facebook this is now
  73. 1:59this used to be available no longer
  74. 2:00available but the facebook lexicon
  75. 2:02i think google has something as well
  76. 2:04where you can type in a word and you can
  77. 2:06see
  78. 2:06uh when correlating events used to
  79. 2:10happen so when are people searching for
  80. 2:11different things
  81. 2:12so here's a funny thing of um people
  82. 2:15saying party town had a hangover and so
  83. 2:16here's party tonight in hang our party
  84. 2:18party tonight is a uh is in yellow and
  85. 2:21hanging over some blue and you notice
  86. 2:22the hangovers often follow
  87. 2:24often peak right behind the party
  88. 2:25tonight it's just people typing things
  89. 2:27but you know
  90. 2:27when is this happening and this is like
  91. 2:29the frequency that and by the way one of
  92. 2:31the parties happening can you see
  93. 2:32oh there's a pretty big party when oh i
  94. 2:34don't know halloween party
  95. 2:35new year's eve party right maybe holly
  96. 2:38you know holidays
  97. 2:39christmas kwanzaa hanukkah all happening
  98. 2:41here but kind of distributed so there
  99. 2:42but boy
  100. 2:43everybody seems to party on halloween
  101. 2:45and new year's
  102. 2:46and not not anything in february march
  103. 2:48or april so it sounds like we need to
  104. 2:50introduce more
  105. 2:51national holidays in those days just
  106. 2:52distribute out the parties
  107. 2:54um here's the design goal the design
  108. 2:56goal by the way
  109. 2:57this initially came out of work that uh
  110. 3:00jeffrey dean and sanjay
  111. 3:01gemmawatt did at google years ago they
  112. 3:03wrote a famous paper one of the most
  113. 3:05the most downloaded paper in systems um
  114. 3:07since then since 2004. this is like 16
  115. 3:10years ago as of this this video the idea
  116. 3:12is it's scalable to
  117. 3:14really massive uh data volumes thousands
  118. 3:17of machines each machine multi
  119. 3:18must might have multi-core um within
  120. 3:21them
  121. 3:22tens of thousands of disks that's a lot
  122. 3:24of data you can imagine processing in
  123. 3:25all of that
  124. 3:26um cost efficiency so one of the things
  125. 3:28they've learned we'll talk about this in
  126. 3:29the warehouse scale computing lecture
  127. 3:31is they didn't buy the highest end
  128. 3:33machines they went with commodity
  129. 3:34machines they said
  130. 3:35you know value per dollar right cycles
  131. 3:39per dollar processing cycles per dollar
  132. 3:40is cheaper if you go with commodity
  133. 3:42machines and kind of let the
  134. 3:44the value of volume uh really helps
  135. 3:46rather than well let's buy the highest
  136. 3:47end
  137. 3:48like kind of the gamer crazy high-end
  138. 3:49gamer machine with
  139. 3:51liquid cooling and the high that you
  140. 3:54know
  141. 3:54per dollar that actually doesn't isn't
  142. 3:56so if you don't have one game i want to
  143. 3:57play for one system i'll have this the
  144. 3:58best thing for my 8k
  145. 4:00screen maybe that's a good value but not
  146. 4:02when you try to compute
  147. 4:03massive bits of data that's compute
  148. 4:05compute processing per dollar that's not
  149. 4:07the right call
  150. 4:07a commodity machine is very interesting
  151. 4:10and they actually started building their
  152. 4:10own which is very interesting
  153. 4:12commodity network um here's the thing
  154. 4:15fault tolerance how do you do with fault
  155. 4:16tolerance well all you do is you resend
  156. 4:18it this is what happens on the internet
  157. 4:19as well if i send a packet i don't get
  158. 4:21an act back i just send it again
  159. 4:22that's what we learned in the internet
  160. 4:24um and that's what tcpip is doing for
  161. 4:26you
  162. 4:27tcp really is doing that word for you um
  163. 4:29here what we'll just do is we'll resend
  164. 4:31that request up and probably not to the
  165. 4:33same machine probably send it to
  166. 4:34somebody else
  167. 4:35um but we handle fault tolerance by just
  168. 4:37re-execution just send it out again send
  169. 4:39the request out i broke it into 1000
  170. 4:41pieces
  171. 4:41and 999 came back one didn't well i'll
  172. 4:43send it somebody else and by the way
  173. 4:44maybe there's a timeout earlier than
  174. 4:46that rather than wait till not 99.99
  175. 4:48came back
  176. 4:48and then send the thousands back out
  177. 4:50maybe like say you know what i'm going
  178. 4:51to give up on you early send it out so
  179. 4:52that it can all come back at the same
  180. 4:53time i have to wait for that
  181. 4:55and we find it very easy and fun to use
  182. 4:56we've had some assignments in in lower
  183. 4:58division to do this and we're excited to
  184. 4:59be able to share this with you
  185. 5:00as well um we've linked to the dean
  186. 5:03paper we called the
  187. 5:04the demon gimbal paper we've linked that
  188. 5:06on the 621c page
  189. 5:07so please do check that out here's a
  190. 5:09typical cluster
  191. 5:11you've got each machine we call them
  192. 5:12nodes each node might have a disk local
  193. 5:14to it
  194. 5:15hanging right off of it and that's a
  195. 5:17reasonable thing to do there's a rack
  196. 5:19switch that connects these all we'll
  197. 5:20talk about
  198. 5:20server rack and array in the next
  199. 5:22lecture um but for here
  200. 5:24they're all connected via a switch which
  201. 5:26is able to take a request and be able to
  202. 5:27know which computer to talk to and which
  203. 5:28it came from where it's going kind of
  204. 5:30like a router in some sense
  205. 5:31um and there's an aggregation switch
  206. 5:33which is doing the same thing
  207. 5:34but the difference is you have a
  208. 5:36different allow different bandwidth
  209. 5:37required requirements
  210. 5:39once you get to um more and more kind of
  211. 5:42central command it's almost like
  212. 5:43plumbing right when you have a plumbing
  213. 5:44from your house the pipe's really small
  214. 5:46once it gets to your neighborhood it's
  215. 5:47much bigger it's from the city it's
  216. 5:49massive okay
  217. 5:50you go to los angeles they have this
  218. 5:52overflow not this i'm sure the sewer's
  219. 5:53larger
  220. 5:54massive as well but you know the if you
  221. 5:56i'll think about uh handling flooding
  222. 5:57you have a little tiny flood
  223. 5:59this happens a lot in the bay area and
  224. 6:01very flat location certainly in los
  225. 6:02angeles
  226. 6:03um you know oh here's my neighborhood
  227. 6:05there's a little
  228. 6:06ditch to handle flood and that's fine
  229. 6:09and then once it gets to the main area
  230. 6:11there's this
  231. 6:12massive i think it's probably like 20
  232. 6:15lanes wide of a flow where the water is
  233. 6:18finally going to the ocean
  234. 6:19same idea here once you get away from
  235. 6:22each node
  236. 6:23higher up in the in the in the hierarchy
  237. 6:25if you will these switches have to be
  238. 6:27very fast much higher end
  239. 6:29and they have an eight gigabit
  240. 6:30connection network here across that so
  241. 6:32you end up having much more expensive
  242. 6:34and we'll talk about the cost here
  243. 6:35uh much more expensive switches and
  244. 6:37connections at the higher level because
  245. 6:39you're
  246. 6:39you're handling larger channels of data
  247. 6:41this is for everybody in the whole
  248. 6:42in the whole maybe the whole array so
  249. 6:45some of the numbers we have here
  250. 6:47are 40 nodes in a rack uh maybe up to
  251. 6:504 000 nodes in a cluster so you might
  252. 6:53hear that here the first cluster here
  253. 6:55um and again one gigabit per second
  254. 6:57bandwidth
  255. 6:58within the rack and then eight out of
  256. 6:59the rack um and the node specs this this
  257. 7:01is a lower-end machine this is like a
  258. 7:03couple years old but
  259. 7:04better than not that not too much faster
  260. 7:05computers as you saw aren't getting that
  261. 7:06much faster
  262. 7:07so maybe eight two giga two gigahertz
  263. 7:10cores
  264. 7:11eight gigs of ram probably that's 16 or
  265. 7:1332 by now and four disks each about four
  266. 7:15terabytes
  267. 7:16okay so this is a fun
  268. 7:20slide because one of the things we've
  269. 7:21been doing in cs10
  270. 7:23and 61a a little bit in 621 people not
  271. 7:26as much
  272. 7:27is to share with you the benefits of the
  273. 7:29programming
  274. 7:30paradigm known as the functional
  275. 7:32programming paradigm if you
  276. 7:33think in a functional way where your job
  277. 7:35is to have a block that has an inputs
  278. 7:37and outputs and doesn't have any side
  279. 7:39effect
  280. 7:39then you're going to be able to work
  281. 7:40with these this mapreduce model a little
  282. 7:43bit easier
  283. 7:44so we've been trying to teach mapreduce
  284. 7:45we actually teach mapreduce the idea of
  285. 7:47a map and a reduce
  286. 7:49in cs in a non-majors class very rarely
  287. 7:51is that happening so i think berkeley's
  288. 7:52really leading the way nationally
  289. 7:54in terms of bringing some of these
  290. 7:55parallel and distributed computing ideas
  291. 7:57into the lower division and
  292. 7:58even in this case into our our
  293. 8:01non-majors class
  294. 8:02which by the way is being shared with
  295. 8:03you know a thousand high schools
  296. 8:05all around the world at least we've had
  297. 8:07a thousand teachers go through our
  298. 8:08professional development in bjc
  299. 8:09so that means these ideas are getting
  300. 8:12out there into the into the winds
  301. 8:15here's the model map map says do a
  302. 8:18function
  303. 8:19over a set of data that's it so i want
  304. 8:22to square all those numbers
  305. 8:231 23 and 10. so here's a little picture
  306. 8:27on the bottom here to show you what's
  307. 8:28happening
  308. 8:28so i'm first applying the map but what
  309. 8:30that does is map says take a function
  310. 8:32it's usually a monotic function a
  311. 8:34function of one argument here in this
  312. 8:35simple case as a function of one
  313. 8:37argument
  314. 8:38takes every number and it's going to
  315. 8:39square it so it's applying that function
  316. 8:40to everybody returning a new list
  317. 8:42not not mutating the list it's not it's
  318. 8:43an immutable list returning a new list
  319. 8:45in which every element has had that
  320. 8:47function applied to it so if i square
  321. 8:49every element of the list 123 and 10 i
  322. 8:51get 1 409 and a hundred
  323. 8:55and now i say combine with plus and how
  324. 8:57combined with again
  325. 8:58it's below the level abstraction of how
  326. 9:00that works is it
  327. 9:02left associated right associative tree
  328. 9:04wise hopefully it's
  329. 9:05a tree structure but that's being
  330. 9:07combined and that's usually a dyadic
  331. 9:09function
  332. 9:09two inputs one output and that's kind of
  333. 9:11a smoosher together and so here i'm just
  334. 9:13adding them up
  335. 9:13so here's here's an example here's one
  336. 9:15and four hundred get added together to
  337. 9:17give
  338. 9:17to give you four or one not in a hundred
  339. 9:19it gives you 109 and then eventually
  340. 9:20that's five ten
  341. 9:21so that's kind of a we're hoping like a
  342. 9:23tree model for that
  343. 9:25who knows who knows how reduction
  344. 9:27happens again who knows how the mapping
  345. 9:28happens
  346. 9:29maybe you tell me oh dan i only have two
  347. 9:31worker bees two
  348. 9:32cores two machines to help you well i
  349. 9:34don't care this is just gonna
  350. 9:36work if i have two this will be split up
  351. 9:38you know one machine might do these two
  352. 9:39one machine might do these two
  353. 9:40i don't know maybe there's four this
  354. 9:42list could also be in that you know in
  355. 9:44the building who knows how big that list
  356. 9:45is
  357. 9:46so you know rarely gonna have a billion
  358. 9:48worker bees working
  359. 9:49cores total cores in your in your
  360. 9:51available set so at some point it says
  361. 9:53well apply it out and then come back and
  362. 9:54then keep working on more of that until
  363. 9:56you're all done but from the user's
  364. 9:57point of view i don't care
  365. 9:58this is the beauty of this abstraction
  366. 9:59from the user's point of view i don't
  367. 10:01care
  368. 10:02so the end of the day i get 5 10. okay
  369. 10:06by the way there's only two data types
  370. 10:07here the data type of the input and the
  371. 10:09data type of the output here and this is
  372. 10:10the
  373. 10:11thinking about domain and range here
  374. 10:14this can be done in scheme we used to
  375. 10:15teach this in scheme all the time this
  376. 10:17is a scheme model those of you taking 6
  377. 10:18to 1 a
  378. 10:18know this and so map and reduce come
  379. 10:22by the way we call it combined it's we
  380. 10:23call it map it combine
  381. 10:25in cs10 and 61a we teach scheme and map
  382. 10:28and reduce are part of that and so
  383. 10:30here's map square over the list 123 and
  384. 10:3210 and then you reduce plus over that
  385. 10:34how beautiful and
  386. 10:35tight that is python it's a little bit
  387. 10:37more syntax for it but not much more
  388. 10:40we have to import the reduce method from
  389. 10:42func tools
  390. 10:43library then i have to define the plus
  391. 10:46function
  392. 10:47a little bit and i have to then i'm
  393. 10:49going to define a square function
  394. 10:51i can obviously use a lambda here and
  395. 10:52you'll see it a little later we're going
  396. 10:53to use the lambda due to spark stuff but
  397. 10:55here's exactly the same thing in python
  398. 10:57reduce of plus over map of square of
  399. 10:59this and automatically i get 510
  400. 11:02really really clean so this is a
  401. 11:04simplified version of mapreduce
  402. 11:06it gets a little more complicated on the
  403. 11:07next slide so let's actually get into it
  404. 11:10in the full mapreduce programming model
  405. 11:13that dean
  406. 11:15and dean at all presented is we have
  407. 11:18a intermediate step that's the step
  408. 11:21between the map and the reduction phase
  409. 11:24and that intermediate step has
  410. 11:27a key value pairing you know by key
  411. 11:29value pairing from python and from from
  412. 11:31any key value model that you've seen
  413. 11:32before
  414. 11:33so here's how map works map takes
  415. 11:39map takes your data a key and a value
  416. 11:43and it's going to process the input key
  417. 11:45and value pair
  418. 11:47slicing data into shards which are
  419. 11:49distributed to workers
  420. 11:51and that produces intermediate pairs
  421. 11:54so the scenario pairs are set up by key
  422. 11:57and then
  423. 11:58they are sorted by key every shared key
  424. 12:01all the data for a shared key is sent to
  425. 12:03one machine you're going to see a data
  426. 12:04you see a
  427. 12:05slide the next slide this will all
  428. 12:06become clear that machine
  429. 12:09then i'm the machine has all data at
  430. 12:11that same key so the key is the way that
  431. 12:12you're
  432. 12:13attaching you're tagging the data in
  433. 12:15some way to distribute it across those
  434. 12:17machines
  435. 12:18because the system says whatever the
  436. 12:20keys you have i will
  437. 12:21coordinate it sorted by the keys so the
  438. 12:24intermediate stage is there so you're
  439. 12:25using the keys in some sense
  440. 12:27not to just add more data not to add
  441. 12:29more craft to your
  442. 12:31input data but to distribute it across
  443. 12:33nicely you get to determine what the key
  444. 12:35assignment is
  445. 12:36and that distributes it nicely so you
  446. 12:37get some control over how it distributes
  447. 12:39in the reduction phase
  448. 12:40and then so everybody just attaches the
  449. 12:43key in a parallel way
  450. 12:44then it's sorted by key across machines
  451. 12:47and then every machine gets all the data
  452. 12:49for a particular key
  453. 12:50and then processes all the data using
  454. 12:52your reducer and there's a reducer
  455. 12:53okay that's it
  456. 12:56and it produces a pre a set of merged
  457. 12:59output usually just
  458. 13:00usually just one value is over here so
  459. 13:03here's an example
  460. 13:06the bread and butter problem we have
  461. 13:08every time when we're looking at
  462. 13:10a new programming language is hello
  463. 13:12world if you have say
  464. 13:14my thing my my system solves uh all
  465. 13:17two-person abstract strategy games we'd
  466. 13:19say solve tic-tac-toe
  467. 13:21if you were to say i have this great
  468. 13:23system that lets you do parallel
  469. 13:24computing easier
  470. 13:25we say to you show me word count word
  471. 13:28count has become the problem we go to is
  472. 13:30the kind of go-to
  473. 13:31bread and butter problem the hello world
  474. 13:33if you will of
  475. 13:35program parallel and distributed systems
  476. 13:38word count says the following i've got a
  477. 13:40corpus of text
  478. 13:41uh maybe it's all shakespeare's work
  479. 13:42maybe it's the gutenberg bible maybe
  480. 13:44it's
  481. 13:44all the hillary emails who knows what it
  482. 13:47is and i want to be able to
  483. 13:48and for word count count the number of
  484. 13:51words
  485. 13:52for every word for every unique word i
  486. 13:54want to count the number
  487. 13:55the occurrences of it that's all i'm
  488. 13:57doing so if i give you
  489. 13:59i do i learn i want to get out
  490. 14:03i there's two eyes one do and one learn
  491. 14:06here
  492. 14:06and so what this is going to do is this
  493. 14:08is going to tag
  494. 14:10remember i told you the output of map is
  495. 14:11this key value pair list
  496. 14:13so what it's going to do is it's going
  497. 14:15to tag it with the key is the word
  498. 14:17itself
  499. 14:18split it by words the word itself and
  500. 14:20the value is going to be the number one
  501. 14:22so essentially how i'm taking for every
  502. 14:23word that comes in i just slap on a one
  503. 14:25to the side of it and the key by the way
  504. 14:27is going to be the word
  505. 14:28and the value is going to be the number
  506. 14:29one and pass it on to the next stage
  507. 14:31then that sorting stage is going to sort
  508. 14:34things there's going to be a machine for
  509. 14:35each of those keys
  510. 14:36so the keys are i do and learn it's
  511. 14:38going to be a machine whose job is to
  512. 14:40just add up those
  513. 14:41values here the number one for all the
  514. 14:43eyes there's a machine
  515. 14:44adding up all the values here's the
  516. 14:47number one again
  517. 14:48for the do and here's one for learn so
  518. 14:50be a i a do
  519. 14:51and a learned machine and each of them
  520. 14:53does the red red reducer
  521. 14:55what's the reducer say well look at this
  522. 14:58set result to one for each value for
  523. 15:01each v
  524. 15:02in the intermediate values of the
  525. 15:03intermediate level intermediate values
  526. 15:05i accumulate in some result in some
  527. 15:08result variable
  528. 15:09and i emit the result as a string that's
  529. 15:11it
  530. 15:13the data is on a distributed file system
  531. 15:15this means the data isn't in memory the
  532. 15:16data is on up files on files
  533. 15:19and those files who knows where they are
  534. 15:21they're on the distributed network
  535. 15:22so they're not just a local to me you
  536. 15:24think oh it's always my local desk it's
  537. 15:25not the case you're on a distributed
  538. 15:27file system
  539. 15:27which means that they they live out
  540. 15:29there in kind of i would say that
  541. 15:30the cloud in quotes so this is my
  542. 15:33favorite example here's the code for map
  543. 15:35and reduce
  544. 15:36and i've got um i've got a file of just
  545. 15:39utterances
  546. 15:41a-a-a-r is in the first file file two is
  547. 15:43blank
  548. 15:44r is in the third file if and or are in
  549. 15:47the fourth file
  550. 15:49or in r the fifth file or and on if
  551. 15:53what does the mapper do it slaps a key
  552. 15:56and a vowel it slaps you know it returns
  553. 15:58a key and a value pair for each of these
  554. 15:59guys
  555. 16:00so it says i'm going to process this by
  556. 16:03the way there's no reason it has to
  557. 16:04process everything i could
  558. 16:05in the mapper it could go through and
  559. 16:06decide to not ex not emit some things
  560. 16:09what it does is see watch right what i'm
  561. 16:11saying is for each word
  562. 16:13it emits that you could say for each
  563. 16:16word that is
  564. 16:17not er and therefore you take error so
  565. 16:20you can do some filtering in there as
  566. 16:21well
  567. 16:21it's your call what i mean you're
  568. 16:23writing this code my point is you're
  569. 16:24only emitting things you want to emit
  570. 16:26here it's meaning everything because
  571. 16:27it's for each
  572. 16:28w and w's emitting everything but you
  573. 16:30could imagine there being a filtering
  574. 16:31step in here if you wanted that so think
  575. 16:32about as you're thinking oh i might do
  576. 16:34this problem but i need to do some
  577. 16:35filtering
  578. 16:36here's where the filtering would happen
  579. 16:37where you don't you don't admit things
  580. 16:38that you don't want to emit
  581. 16:40emit so now all i've done is attached
  582. 16:42one to all these guys now i have a one
  583. 16:44at the end of every single key the key
  584. 16:46again
  585. 16:46are all the words this is this phase i
  586. 16:49talked about the group by key phase
  587. 16:51that says they're going to be a machine
  588. 16:53attached to every key the keys are now
  589. 16:55the five unique words are
  590. 16:57if or or uh and
  591. 17:01what is a reducer then all those values
  592. 17:03all those values are sent to the reducer
  593. 17:05what does reducer say basically it says
  594. 17:07add them all up
  595. 17:08add all the numbers up whatever the
  596. 17:10values
  597. 17:11okay sometimes it ignores the keys i
  598. 17:13don't care what the output key is
  599. 17:15but for each of these guys all i'm going
  600. 17:16to do is output a four
  601. 17:18emits a string and so what comes out of
  602. 17:21this is
  603. 17:22in a way it's a dictionary and think
  604. 17:24about this as the
  605. 17:26the outputted value of the whole system
  606. 17:28here is
  607. 17:29every machine says well i was i was the
  608. 17:32machine
  609. 17:32assigned to ah and my value is four so
  610. 17:35it's ah colon four
  611. 17:37why was the machine doing er and i only
  612. 17:39saw one of them so i'm er colon one
  613. 17:41so that's the idea the output is in a
  614. 17:43way of this whole system
  615. 17:44of that was the goal the goal of the
  616. 17:46whole problem was word count was
  617. 17:48return a dictionary of where the keys
  618. 17:50are the words
  619. 17:51every unique word and the value is the
  620. 17:53number of occurrences they occurred in
  621. 17:54the original thing
  622. 17:55well that's what's happening here this
  623. 17:57is the output if you turn your head
  624. 17:58sideways
  625. 17:58this is kind of like key down here and
  626. 18:00value up there
  627. 18:02that's it and that's your output of the
  628. 18:03of word count very simple right it
  629. 18:05doesn't seem that hard
  630. 18:06so mapper all it does is slap a one at
  631. 18:07the end of it and reducer says well i
  632. 18:09got all the ones what do i do with them
  633. 18:10add them all up so it's really reduce
  634. 18:12plus okay
  635. 18:14key is attach one the the mapper the
  636. 18:16reducer is
  637. 18:18just take it's just plus that's all i'm
  638. 18:20doing that's it
  639. 18:21so that is word count pretty clean
  640. 18:25how does this work well
  641. 18:29let's say the corpus is not just you
  642. 18:31know four files but a ton of files maybe
  643. 18:33it's a huge directory maybe it's a lot
  644. 18:35how does this work well this is very
  645. 18:38clever this is the part that's the
  646. 18:40system engineering part of it that they
  647. 18:41really is
  648. 18:42quite brilliant and by the way not all
  649. 18:43maps take the same time sometimes
  650. 18:45you know the map task might have maybe
  651. 18:47the file you're given is a much larger
  652. 18:49file than everybody else so you're gonna
  653. 18:50be doing more work than everybody else
  654. 18:51does so not all maps take the same
  655. 18:52amount of time i kind of you know in
  656. 18:54in my simplification in my in my simple
  657. 18:56example of
  658. 18:57map of square you know square takes the
  659. 18:59same amount of time for every number you
  660. 19:00can be given
  661. 19:02except perhaps if you have big nums and
  662. 19:04the numbers actually are stored as a
  663. 19:05string but for for now you think that's
  664. 19:07that's a bad model really maps can take
  665. 19:09variable amount of time as reducers can
  666. 19:11take variable amounts of time
  667. 19:12who knows how many you know how many of
  668. 19:13those ones i had to add up in that model
  669. 19:16so there's a there's a lead we call it a
  670. 19:20i guess a master controller
  671. 19:22which assigns the map and reduced tasks
  672. 19:24to these worker servers there's all
  673. 19:25these servers it's like a you know
  674. 19:26worker bees if you will
  675. 19:28as soon as look at this as soon as a map
  676. 19:31task finishes
  677. 19:33that server can be re assigned a new
  678. 19:36task so this particular worker
  679. 19:38was only assigned map task this
  680. 19:39particular worker's only time map
  681. 19:42map tasks as well this worker can be
  682. 19:44assigned read tasks
  683. 19:45because it's reading one meaning what is
  684. 19:47this doing
  685. 19:48this guy says well when that's done i
  686. 19:51can start reading those values in and be
  687. 19:52waiting for it
  688. 19:53and i also need to know where i need to
  689. 19:55be reading this one so this is also
  690. 19:57reading for that when that comes in and
  691. 19:59then i'm reading for this one so when
  692. 20:00they're all done
  693. 20:02now i can start my reduction task so why
  694. 20:04wait why wait till they're all done
  695. 20:06it's very clever to think about this um
  696. 20:08to start its reduction
  697. 20:10so each worker and so this is this is
  698. 20:13saying look these workers are the mapper
  699. 20:14workers these are the reducer workers
  700. 20:16that's not always the case i'm basically
  701. 20:18once map
  702. 20:19map one finishes i'm telling the system
  703. 20:21hey i'm available for more work
  704. 20:23maybe i'm the job maybe this becomes me
  705. 20:25i mean it's drawn this way to make this
  706. 20:26diagram clear
  707. 20:27but it could be that this worker one is
  708. 20:29assigned
  709. 20:30to do that read and then it's going to
  710. 20:31be waiting and reading then basically
  711. 20:33this
  712. 20:34could be up here is what i'm saying okay
  713. 20:36it doesn't have to be like oh worker one
  714. 20:37is only mapping and worker two was only
  715. 20:39mapping ever you know this basically
  716. 20:40they go back into the pool of workers
  717. 20:42and the jobs are just out there so
  718. 20:43there's a queue of jobs that need to be
  719. 20:45processed like you gotta start the
  720. 20:47reducer you gotta start the next mapper
  721. 20:49um there may be many more mappers than
  722. 20:50just three here so
  723. 20:52it could be that um
  724. 20:57you're you're heterogeneous in that
  725. 20:59sense i'm a worker bee and i was doing
  726. 21:00some mapping and then i'll do some
  727. 21:01reducing and maybe some more mapping to
  728. 21:03be done i finished that reducing and
  729. 21:04there's more mapping to be done there
  730. 21:05who knows what that is
  731. 21:06um what's done and by the way what
  732. 21:09typically happens
  733. 21:10is it's not just one map and one
  734. 21:11reduction you tip it this is i
  735. 21:13i i read the paper i found what this
  736. 21:15what happens is there could be a map
  737. 21:16reduction here
  738. 21:17that feeds into another map reduction to
  739. 21:19feed into another so there's this whole
  740. 21:20series of
  741. 21:21a map reduction not just one it's just
  742. 21:23one problem it's all i have to know
  743. 21:24there might be several phases of it and
  744. 21:25doing phase one phase two phase three
  745. 21:27and you just continue to process the
  746. 21:28machine and it's really fun to watch
  747. 21:30there's a little um there's like a
  748. 21:33there's like a
  749. 21:34dashboard that you can see how things
  750. 21:35are processing it's very neat to see so
  751. 21:37i encourage you to take a look at that
  752. 21:38if you ever get this working
  753. 21:41the reduction task begins as soon as all
  754. 21:42the data shuffles finish maybe you gotta
  755. 21:44shuffle things around and finally there
  756. 21:45because
  757. 21:46in this problem as before that last
  758. 21:49mapper might contribute
  759. 21:50a one to your reduction so here in this
  760. 21:52particular problem
  761. 21:53you can't do any reduction until you've
  762. 21:55finished all the mapping
  763. 21:56so that i i might have implied that you
  764. 21:58can start reducing before the mapping is
  765. 22:00done
  766. 22:00in this particular case i certainly
  767. 22:02can't start adding up my ones until
  768. 22:03all the guys are gonna be ones feed into
  769. 22:05my machine and that's done
  770. 22:06so in this particular model you have to
  771. 22:08wait till all of the maps get done
  772. 22:10before i can even start thinking about
  773. 22:12the
  774. 22:12reduction but you can start reading
  775. 22:13their values certainly that's kind of a
  776. 22:15nice thing
  777. 22:17here's the thing i mentioned before uh
  778. 22:19if a worker doesn't respond
  779. 22:21um network failure disk failure hardware
  780. 22:24failure memory failure kind of could be
  781. 22:25a
  782. 22:25lot of failures i just reassign the task
  783. 22:28if a worker dies that's important
  784. 22:30this amount of java code that does that
  785. 22:32same word count if we look here
  786. 22:34in detail what's happening is here it's
  787. 22:35saying uh i'm outputting
  788. 22:38the value and a one this is the actual
  789. 22:40code for this and here i'm just here's
  790. 22:42my sum of zero and i'm sun plus equal
  791. 22:44that value
  792. 22:44a little bit more complicated but
  793. 22:46essentially that's what we were doing
  794. 22:47before and that's the java code for it
  795. 22:49that's it so there's word count folks
  796. 22:51and by the way most of this is just set
  797. 22:53up we're just setting these things up
  798. 22:54here in this early part of it
  799. 22:57that's it pretty exciting word count is
  800. 23:00a really nice example
  801. 23:01and map produces a very very powerful
  802. 23:03paradigm
  803. 23:04the next lecture we're going to talk
  804. 23:06about maybe there's a cleaner
  805. 23:08lighter weight way than when i'm looking
  806. 23:09at maybe 35 lines of code maybe there's
  807. 23:11a little tighter way of doing that in
  808. 23:13spark we'll see that then
  809. 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.