[MUSIC]. So, let's talk a little bit about what kind of systems MapReduce is deployed on. So I won't spend to much time on the system internals since this is a data science course not a distributed systems course. But it's good to have an intuition for what's going on under the hood. So, there's three types of systems to be aware of, architectures to be aware of: shared memory, shared disc and shared nothing. in these diagrams, these cylinders are the discs, these rectangles are the memory and the circles are the processors. Okay. So shared memory means that every processor has access to all of the memory and all of the disc. And this is what you think about when you have sort of a laptop that has four cores in it. And if you have a quad core system, or [INAUDIBLE] you have you know, six or 12 cores in your laptop. This is the model that's being used, the architecture that's being used. Shared disc is somewhat less common in, at least in these contexts that we're talking about. Although certainly, certainly common enough overall where you'll have the individual machines all access a shared file system. Now that, that setup is very, very common, but for using that setup for parallel ana, analytics is sort of the domain of, of high end commercial databases. So, you know, your Oracles and your IBMs will often used a shared disc architecture. Okay. And then shared nothing is really what we're focusing on here. And this is what MapReduce is designed for. And what increasingly, and largely parallel databases are designed for as well, okay. And shared nothing here means individual machines that are only connected by a network. Okay. No shared memory, no shared disc, fine. So, the argument is that only the shared nothing architecture can scale to sort of thousands of computers and beyond. Okay. Because eventually that sh, the the sharing of memory or the sharing of disc eventually becomes a bottleneck and limits how many computers it can attach to the same device or, same logical, logical or physical device. Okay. And so learning how to program these massive shared nothing clusters is what MapReduce and parallel/g, this is all about. And the shared memory machines are perhaps the easiest to program, but are conventionally assumed to be pretty expensive. I should point out that you know, the costs are dropping fairly quickly. So, it's getting more and more feasible to buy a pretty beefy machine with lots of main memory and lots of cores, and you know, your problem might fit inside that one. So, you know, when you see people deploying Hadoop and MapReduce on fairly small clusters in the order of say ten nodes. You should see whether the data size that they're processing are actually all that large, right? It could be that the data size that they're processing is something that fits in main memory on a similarly priced amount of hardware. Okay. So, fine, so Hadoop and MapReduce are designed for really, really large clusters. Okay. Alright, so this is the context we're in. A large amount of commodity servers connected by a high-speed commodity network. And here you know, you think about a data center. There's a rack that has a number of servers and there's a data center that has many racks. And this is how you organize your thousands or tens of thousands of computers. Okay. Alright. So you're looking for massive scale parallel, parallelism that will run you know, jobs that will run for many hours even on thousands or tens of thousands of servers. That's really the context we're in. So, when you're in this context an issue that comes up that does not come up all that often in a much smaller context is failures, right? So, if you're running a job for a long time on thousands of computers, the chance of something going wrong during that job becomes essentially, you know, 100% probability, right? There's going to be something that fails. And so, your system of processing, doing this sort of analytics has to just tolerate this kind of failure. You can't just, you can't, roll back to the beginning and just restart every time there's a failure occur, or you'd never get anything done. Okay. So, even if the mean time between failure for say a disk is a year. If you've got 10,000 servers with multiple disks or 10,000 disks so say, spread across 1,000 servers or any combination thereof, you're going to start to have failures you know, once per hour. And you can look up the mean times failure and actually do the math, but it, there's a, there's a couple of nice papers out there that I'll try to put in the readings if I remember. Okay. So, failures are what we're concerned about here. Alright. So, that's hardware, popping up the stack one level is this distributed file system. And you might see HDFS too, which is the Hadoop distributor file system. So, you remember the context here, was that MapReduce proposed in a paper in 2004 by Google, and Hadoop was the implementation of that, of the ideas in that paper. Okay. So, we can almost use an interchangeably because the actual implementation in Google has certainly evolved since that paper and is not completely known, right? So, when people are talking about MapRoduce, they are typically talking about Hadoop or other implementations of the program model that have nothing to do with sort of the scale out. But shared nothing, MapRoduce can be assumed to be synonymous with Hadoop. Okay. So, this is a file system for very large files, and the idea here is that if you're going to take a single file on your own computer, you can manipulate it as a single unit. But if you're going to take a very, very large file and dis, and you know, put it on a file system that's a, that's on a cluster machine. Then there has to be some layer of logic that knows how to split that data into pieces and put those pieces into different machines, and keep track of where they are. And that's what this distributive files software does. So each file is printed you know, as you're uploading this data to the cluster, each file is partitioned into chunks that are say typically 64 megabytes although these are configurable. And so each chunk, and this is critical, is replicated several times. Right? So there not just one copy of the chunk, there might be multiple copies on different machines. Why? Because if one of them goes down, you want to have access to the other chunks. Okay. And you wanted to make sure that these are spread across different racks in case the entire rack of computers goes dark. You still have another copy of it. Okay. And so the implementations here are GFS and HDFS. DFS is the concept, and the implementations are GFS and HDFS. Alright. So, here's the phases of MapReduce that's a little bit more detailed than the abstract phases that we were talking about when we were talking about the program model. Okay. So there's a file split here that's read from HDFS, and remember HDFS means there's replicated, the chunks could be, you know, there's multiple copies of every chunk. And, there's a unit of code called the record reader that breaks that chunk into. I'm using chunk and split for this synonymously, I'm not a big fan of the term split, because it sort of sounds too much like a verb to me. the record reader parses that splitter chunk into individual records. Then the, the programmers, you know, the user's map function is called on that individual record. And then there's a step called combine that we haven't talked about, that I'll talk about in a moment. Next, actually. Okay. And then the output of these phases are written out to local storage on that node, as we said. So then these regions in local storage, one per key, are pulled across the network by the reduce phase. Then all the regions from all the difference, all the different map tasks that correspond to the same key are sorted together in parallel. And finally, the users reduce function can be called to produce whatever output it produces. And then the output of that step is actually written back out to HTFS so that it's replicated. So again, if something goes wrong in the map phase, you have to rerun the mapper. And if something goes wrong in the reduce phase, you have to rerun. You can, you can pull the output from the, the local storage from the map phase. and if something goes wrong in the overall job, you know you're safe because you don't lose data because the HDFS are, are, are replicated. And if you, you know, if you lose a reducer, and you lose a corresponding mappers, that's fine can sort of rerun whatever you need to rerun, fine. So though, the points is that you're, you're guaranteeing for fault tolerance during Java execution. Alright, and so, so let's talk real briefly about the combiner. So, to think about why we need a combiner, go back to this word count example that we began with. Well, in each case we produced a word and just the number one in the fir, in the earliest version of this, we just produced the number one. Right. Sending all of these occurrences of word one. So word, if word one appeared twice, then you're going to get a key value pair with w1 and a number 1, and another occurrence of w1 and a number 1. And you're going to send both of these guys across the network to be sorted and parallel and grouped, in order to be processed by the reduced. Well that's sort of wasteful. You'd like to combine these into a single record, w1 comma two. And then just send that, because it's smaller. Well, you could rewrite your map function to, to do this. But it's such a common need that you can that, that, that you have this capability called a combiner. And so a combiner identifies key value pairs of the same key and lumps them together before sending it over to the reduce side. So it just saves the network traffic. In many cases, the combiner function can be literally the same function as the reduce function, it all works out. And what needs to be true for this to work is that the function that you're applying needs to be associative and commutative, but I'm not going to say much, much more about that. Okay. So here's what it looks like in pseudo code. We saw the map function earlier, and we're emitting this key value pair of a word and account. And we saw the reduce function earlier, but we're adding new as a combiner function that has the same type signature as a reducer. And again, often in, in many cases can literally be the same same implementation as a reducer. And the only point of this, is that it's being applied before sending data across the network. Alright. So here's sort of a summary of a Hadoop job that I like a lot, and this is from Huy Vo who is, who is now in NYU Poly I believe. He is still at NYU Poly. So the data begins on HTFS, and there are these chunks, and in the in part partitions go to map tasks. Now, these are again not invocations of the map function, these are entire tasks. And the map, each individual map invocation may produce multiple key value pairs regardless map task almost certainly does, right? So it's going to produce these local sort of colored key regions. And the regions are going to be sent across the, pulled across the network to the reduce servers. And here in this example, there's only two reduce servers. Now, before we gave this example, we sort of showed all the blue ones going together and all the red ones going together. That was a bit of a simplification. What actually is going to happen is that if you only have two servers, well, all the hundreds of, hundreds of possible keys need to be mapped to just those two servers. So you're definitely going to get a mix going to the same place. And this is where we've been a little bit glib up until now. We, we've, we've said that you specify the key and it hashes to a particular reducer. whi, which is, which is true logically, but before that, you have to get it to a machine, where lots of reducers, where lots of reduced tasks might be running, or lots of reduced invocations might be running. Okay. And so, you know, in this case the blue, the blue guys and the red guys both end up on the same machine. And the green guys And the orange guys building it up on the same machine. And then this parallel sort manages that. Right. So, puts all the red things together and puts all the blue things together. And then for each individual color, one reduce indication is called. Okay. And the reduced function is called and it produces the output partition and all that Output is wri, written back out to HDFS. Okay. Let me stop there and I'll pick up here, and talk a little bit about parallel databases and how they do query processing with the point being that it looks a lot like MapReduce. [BLANK_AUDIO]