[MUSIC]. So, last time we talked about parallel processing as a lead up to map reduce. And we ended up with this schematic here. In the context of this example, where we're counting words across a set of documents. Just sort a canonical example to start thinking about programming a map reduce. And so each one of these vertical black lines represented a document. And we split them into smaller sets and send each one of those sets to a separate machine. And then we applied our map function to each one of those documents in turn. And so the map function if you recall, took a single document and produced a set of pairs, and each pair was a word along with a count of the number of occurrences of that word in that document. Various variations on this that you can imagine. Okay, now, this word may have appeared in multiple documents, one here, one on this machine, one on this machine and so on, and so now we need to group them altogether onto a single machine so that we can count them. And that's exactly what this shuffle phase did. Okay. So here I've written four different tasks, processing sort of, looking like it's processing a single group at a time, but you know, I want you to think about how many map tasks do we already have and how many reviews tasks are we going to have? Well, the map tasks are one per document. We have to, we have to call the Map function, how many invocations of the map functions are going to there going to be? Well, we're going to call it once per document. How many invocations of the Reduce function are going to be? Well, it's the number of groups that are produced by the output of the Map function. In this case, it's once per unique word appearing in any, in any document, okay? And so in some sense, the number of machines we need to apply to this problem is maybe kind of predictable in the map phase. It's, it's. Corresponds to the size of the input data set. Which we, you know, presume to know. But the number of reducers we're going to need is maybe not known ahead of time, right? It depends on the size of the Map output. Here we might be able to reason about it, because we maybe know how many words there are in the English language. And we can assume that with a big enough set, all of those words will be represented. at least once. But in general, it's, it's dependent on the output of the map, so you don't really know. Okay, and so, the only point I want to make is that we made a decision here to draw it as four different machines. But it may you know, it may, it may be the same six machines you used the Map phase, or maybe a thousand machines and so on. You know, nothing's stopping you from sending all of the word occurrences to a single machine. And having this one machine process the green group then process the red group then process the blue group and so on. Or maybe it would be four at a time because quarters in the machine, okay. But that wouldn't be as, perhaps as efficient because it would be doing a lot of. Serial work. At the other extreme you might think, well we're going to need millions of tasks. Let's allocate, you know, hundreds of thousands of machines to process them. So that each machine is doing very little work. Okay. And that might make sense, but then the trade-off is perhaps sort of long, spinning of all these machines and, and preparing them to do the work. Okay. So there's a decision to make there. And we'll, we'll come back to that. Okay. But let's talk a little bit more about Map Reduce itself. This is, I'm, I'm belaboring this for a reason. Cause I, what I want you to do is start thinking in terms of Map Reduce. Every problem you have, think what if the data set was absolutely enormous, way too big for a machine? How we're going to split it into pieces? And a, a very good way of thinking about how to split things into pieces is to think about how you'd write a map reduce program to do whatever it is you're trying to do. And so this is, yet again the same example, just drawn a different way. So here the input is document ID followed by a value, and the value here is the entire text of the document. And the Map Function, just to make this clear, produces a set of things, not just one thing. And they're shuffled to produce this. So this is word one with a count of one. Word two with a count of one, word three with a count of one and so on. And the other side what we get is word one with a group. Of all the occurrences, 1, 1,1 1, 1, 1, 1, and then find the reduce function counts them all up and finds that there are 25 occurrences. Okay. So I'm probably, I guess if, if this is completely obvious, you can always fast forward I guess one of the beauties of doing this online. Okay. So fine. So what is map reduce? That's the programming model we just described and there's a paper in 2004 that's on the reading list that describes this. And there's a, a couple of key motivations for doing this. In that paper that I think sometimes get lost when you hear about the popularity of Map Reduce today. and we'll talk about those two, those two benefits in, in, in, in a moment. So. So one thing to realize is that map-reduce refers to the abstraction and it's the name given to it by the authors of this 2004 paper. Hadoop is an implementation of Map Reduce. It came a few years later. And was written by some people at Yahoo! Originally and then became an open source product that is managed by the Apache and has lots of contributors, okay. So, you know the key idea for Map Reduce was really this programming model, which it says here, right. Now it had a system with it as well, but the programming model being able to express lots of different tasks and you know, have some sort of implementation to automatically turn that into parallel job, turned out to be pretty powerful, right? This was an attractive way to write parallel programs, again, because you didn't actually have to worry about the parallel, all you had to do, write a serial map function and a serial reduce function and the parallel had happens for free. And so the evidence of this is not so much about the system, as it is the program model, is that. You see map reduce implementations appear in other contexts. There's people who have implemented map reduce over uh,GPUs, There's people that have implemented Map Reduce on multi-core machines in shared memory. There are people who have implemented Map Reduce on high-performance computing platforms on you know, mobile, groups of mobile phones and so on. Okay. So, this goes back to one of the motivations for this course. So I want to focus on abstractions where possible, as opposed to tools. And so we're talking about Map Reduce in the the programming model. But we'll spend a little bit less time in the specific implementation, Hadoop. Although you will have a chance, an optional assignment, to work with Hadoop directly. Okay. So fine, so what is a data model of Map Reduce? It's this bag of key value pairs, I mean by bag is a set that might, might have duplicates in it. Right, and so we've seen that before, the document ID with the value, that's a key value pair. Sometimes on the input we'll be a little sloppy and not worry about precisely what the key and the value is. For example, if you're just given a record you can assume that its, say the, the entire record is the key. Okay, or a document sometimes even we may not have a explicit document id. But you can assume the URL or the file name or something is the key. the output of the mapper, though, the, the distinction between key and value gets really really important because that's what controls the shuffle. As we've seen in that, in that example. Okay. And so, both the, the data model here is all about key value pairs. And the input is going to be a set of key value pairs and the output is going to be a set of key value pairs. And the point is that. The set of key value pairs can get arbitrarily large, right. We're going to be able to process the set no matter how big they get. There is kind of an implicit assumption that the key and the value are small. And small here doesn't necessarily mean very, very small, it just means that it needs to fit on one machine. There's no support for a, for if value goes to be terrabytes. It's it's not going to work. And so a document fits on one machine. That's okay. You know a an image that's on a machine, and so on, okay. So, fine, so the map phase as we've said, you provide a map function. The input is input key and input value. And the output is a bag of intermediate keys and value. It doesn't have to just produce. You're taking a single input, it produces single output. It can produce a set of things. We saw this with the word count example. A single document came in, but a set of things came out. That's okay. Alright. In the reduce phase what you're given is an intermediate, the intermediate key. This would be the same intermediate key that was produced by the map phase. Alright. One instance of the same intermediate key produced by the math pairs. And then a bag of values that were associated with that intermediate key. And they may have. And the key. The important thing here is that they may have come from any mapper whatsoever. The grouping into this bag of all, of everything that shares this same intermediate key is handled automatically by the system. Okay. So the system will group all pairs with the same intermediate key and then pass that bag of values to the reduced function. And the implementation details of whether this is actually handed to you as a bag or whether it's an iterator, if you're new, if you're familiar with that term, that you could step over is implement-, implementation independent but it's not important to think about. It's a collection of values. Okay. Fine, so here it is all on one slide. The map function takes a n key and an n value and produces a list, a bag of out key and intermediate value pairs. And the reduce takes an out key and a list of intermediate values and produces a list of out values. I don't think I like this line, I think I prefer the other one. The one thing I will mention is that the terms map and reduce, I tried to motivate that in the last segment where you can think about converting a tif image to a pmg image you can make a mapping, a function that maps Every TIF image into some PNG image. And that's where the term came form. And if look back this kind of came from the functional programming community that use this term. It doesn't precisely mean the same thing but it's inspired by that. Okay. All right, so here's maybe the implementation for the example we have and a lot of times what I like to do is sort of ask you to pause and stare at this code for a little bit and think about what it does. Here we've kind of gone through the example a lot so I'll reveal the secret but it's still instructive to work through this for a moment yourself and in fact, maybe I'll end this segment there, and you can stare at this and make sure that you understand what it does.