Okay, so let's look at the CoGroup command, so we've seen the group command, which works in a single data set and the co-group command does the same thing on multiple data sets. And we seen this mechanism before we talked about map reduce but here in pig we make this operation explicit and give it a specific command and we'll try to see why that is. Okay, so the syntax looks like this [SOUND] you say CoGroup a dataset by some field reference or multiple fields. And here we're referring to a field by name. And here we're referring to a field by position as we've seen before. And so, if this is A, a bag of tuples and this is B, a bag of tuples then the CoGroup of them looks like this. So, now we have a Group field and we have a dataset A a field named A and we have a field named B. Actually we still have the same grouping key but we have two groups associated with it. One from one data set and one from the other, okay. So, a couple things to point out. One is, if, if there are no tuples from one of the data sets and it just get's the empty group. here you see [INAUDIBLE] from B. And then the other thing to point out is this is a little different then what happens in MapReduce directly where you know, you, we,what we showed, with, with the, JOIN operator without produce is that you could, use the JOIN key to sort of CoGroup two relations. And then all of these tuples would all sort of appear in the same group in the reducer but here we sort of make it explicit to make it two separate groups, okay? Fine, so I want to come back to co-groups in a second but first let's talk about another operator, which is just called JOIN and so it does exactly what you'd expect. So JOIN, you say JOIN A by some set of field references and B by some other set of field references, okay. And given A and given B the result will be what we've talked about in the past, which is you know, you look for the first position of A, $0 and find all the corresponding first positions of B. So, here's one; 1,1. So we have 1, 2 from A and 1, and 3, 1, 3 from B. Okay, so there's nothing stopping you from having multiple datasets out here. As many as you want and they can all be sort of processed together. And so you think about what's going on here is this one MapReduce job underneath the sheets where in the map phase every tuple from a variety of data sets is all being associated with the JOIN key. represented by this field reference in the syntax, and then shuffled across the network to arrive at the reduce a fade. Okay, so here's what it looks like well if we have multiple relations where represent by color here green blue and red well they can all be processed by the map phase and associated with keys shuffle them across the network to produce the Join tuples here. Right so there's not there's nothing fundamentally binary about this operation. Operator. Right? It's just processing tuples associating them with their, with their respective Join keys and shuffling them all across the network. What can go wrong with this basic mechanism of associating tuples with a Join key, shuffling them across the network and then producing the, you know, completing the Join on the reduce side, which is what we did the assignment as well. Well, so one example might be if one table is very very large and another table is very very small. There's an opportunity to do something much much faster, which is, replicate the small table across all partitions of the large table which allows you to do all the work in the map phase, okay. The second sort of special case here is a Skewed join. And what I mean by that is, if there are many values in one table that join with many, many, many values In the second table then you'll end up with one reducer doing all the, most of the work. And the effects of parallelism gets sort of washed out. And the third special case algorithm is to do a Merge join. So this takes advantage of the fact that you may have already grouped two relations in the same way such that the, you know, for sure that on a single machine. I've got all the tuples I need from one relation and all the tuples I need from the other relation. Okay. So let me see if I can make this more clear in, pictures in the next few slides. Okay. So for replicated joined, the situation is we have one large table broken into pieces and one much smaller table. And so what we could do is just shuffle all these two bulls across the network and do the Join on the reduced side as we do normally. But there's an opportunity here that says well look if this thing is small enough to fit in memory on a single machine. Why don't we just copy it? Send it out there to every, you know, every map function that wakes up will go pull it across the network directly. Okay. Then we have all the information we need right here in the map side. Every one of these tuples can be joined with corresponding tuples in this partition. Every one of these tuples could be joined with corresponding tuples in this partition and so on. And so the end of the map phase you end up with the right answer. All the joined tuples for you know this, this sort of blue A and this red, red B. Okay. And so why is this cheaper? Well we didn't have to shuffle everything across the network. We did it all on the map phase, and it, it makes a huge difference. Okay. So you might see this called a Broadcast join, as well. So the small relation must fit in memory, and the idea is that each mapper in, in the map function pulls a copy of the small relation directly out of HGFS. Okay. All right. So for the Skewed join the situation is, is as usual. You got two relations, this sort of blue one and this red one. And the map phase associates each tuple with its JOIN key. But the problem is that most of the data ends up on single reducer and the reason is because maybe most of the data here is associated with a single Join key. Okay. So, for example, if you're joined on order ID and line item, you know, you just try and associate all line items with their corresponding order, it could be that one order had millions of parts and all the other orders had five parts, or something. Okay. So if there's if there's significant refraction in the overall data set is associated with a single join key then this reducer will be doing all the work and there won't be much benefit to parallelism okay. So this would work but there's a problem. So, this is from former student here at UDUB who did some work on this problem. And this plot shows time in seconds on this x axis. And this is just a list of all the tasks and so you see that the reduced task here, they are, they can't start until all the map tests are finished. Okay. And the map tests mostly finished quite quickly. Sort of, maybe, 20 seconds. But one or two of these map tasks take a very long time. You know, sort of on the order of 270 seconds or so. Okay, so all this space in here is sort of wasted work. There's a couple of problems here, one is fundamentally map reduce in may applications you logically could start doing some work early based on the map output that has already been finished. And that's just not the way map producers is designed. y, you can't take, you, it's not designed to be able to take advantage of those applications because you can't guarantee that it's safe to do so. So, for example, if you're adding up numbers, you could start adding up in, in the reduced stage. You could start adding up numbers early. But if you're doing something more complicated you may actually need to wait for all the results to be present in order to get the correct result. And so, since they can't guarantee that it's safe th, they make you wait. Okay, so fine. So, this skew problem ends up sort of killing parallelism in terms the job that you know could take sort of 50 seconds. And the one that takes 300 to 350 seconds. And is not even significantly longer than doing this sequentially. We should of added all these pieces up. Well, I shouldn't say that. This is so worse than, in sequential. But you cer, you certainly lose a lot of your bandwidth in parallelism. Okay. So one task might take five times longer than the average and so there's little benefit. And so one take away here. So I'll tell you what one way, one way to partially solve this problem on the next slide. But the take away here is, if someone asks you what the, you know, one of the biggest performance bottlenecks of MapReduce is, you should say stragglers or skew. skew is more of a term for this in the database community but stragglers is a little bit more common. So this is a straggler task that takes a lot longer. Alright, so what can we do about this? Well, here's the situation where we have a lot of keys that all ended up on the same reducer and this guy's taking too long. One thing we can do is split this reducer up into more reducers. So take, allocate a few more reduce tasks and move that data over here and split it into three more. Now, broadcast that little red relation, right? So I made it disappear from over here recognizing that this is perhaps small and replicate it. To all three of these reducers. So, sort of a combination of the replicated or Broadcast join and the regular reduce side, pass join. Okay. So, now, you've got three reducers working on this problem as opposed to just one. And you get things a, a little bit more balanced. And so this Skew Join is something the pig can do if you specify it. Alright, it won't do it automatically though. So, fine. And so now we get our entire Join relation. So finally Merge join is the third special case and this as you saw the first special case was an opportunity to do things much faster, if one relation was very large and other relation was very small to fit in memory. Skew join is more way it is there is a problem that can occur and you need a special trick to be able to address the problem. Merge join is more like the formal. It's looking at an opportunity to use up more high performance algorithm when certain conditions are met. And so one of those conditions well, when you recognize that red relation and the blue relation are already been co-partitioned on the appropriate Join key. Right. Then the map phase alone has enough information to just keep enjoying itself. It doesn't actually need to assign it all to assign this tuple to a joined until you shuffle it across the network. Right. So you know that all the line items for a particular order. You know, for every order that's here, all of it's line items are also here on this machine. If you know that to be true then you can do it in the map phase, okay. So the question maybe when is that true? well, that's when we go back to the co-root operator. It's possible that you've already partitioned these two tables on the appropriate JOIN key because of a previous command in pig. And so can be aware of that and use a merge, Merge join. Sorry, I, say aware of that. If you still specify explicitly you want to use a Merge join but, so take exactly that not automatically but you can take advantage of the situation where when it arises, alright. So, since each map already has local access to the records from both relations. They're already grouped and assorted by the Join key. You can just read in both relations from disk in order and compute the Join. Alright so we had this CoGroup operation and we have Join operation and why do we need both. Well the reason we made CoGroup explicit is that if you think about a join as really a two step process right there's one step to create the groups based on the Join key. And then a second step to actually produce the joined tuples. But that group creation step is useful for a lot of applications, not just producing a JOIN. So for example, if you want to, CoGROUP and add up all the contributions from each relation you can do that in a single step. You don't need to sort of join the two groups first then do another grouping afterward, which is something you would have to in original databases. So that's a chance to take what would be two bad produced jobs and combine them into one. Okay, and so here I guess the example if you just want to count the tuples. Okay. So, but, you know, what the point out here is that you can't express Join. JOIN is essentially just syntactic sugar, you can't express a JOIN in terms of CoGroup. First is, the first step is to CoGroup on the same colu, on the, on the JOIN columns. And the second is to run this foreach command that would generate A flat view of results of revenue. Sorry I guess I should say A and B here that's sloppy. I changed the names. Okay. Alright. So, other commands that we're not going to talk about in detail are Store which writes data out to HTFS. So, it's is available for. You know, future commands. union that combines two data sets together and removes duplicates. Cross product which finds all possible pairs between two data sets which can be useful if you're going to compute some sort of similarity function as we talked about. and then dump which prints outputs of the screen. and then order which sorts, sorts the output. Okay. So as an example of Store here. remember you could also, just, just like you can use your own custom function to parse data, you can also use your own custom function to write it out. Okay. Which means it allows the data to be more compatible with, say, some other system that you're using. Say, MapReduce itself or some other application that expects data in a certain way. And so that's, again, I sort of stress this with, with the load command as well. But this is actually pretty powerful and a pretty big difference from this, you know, walled garden approach that relational databases take. Where everything sort of goes in and then it's, you know, under complete control, the database. This sort of has more permeable boundaries where you can have kind of have data lying around in whatever format. And you're still able to process it with pig but you can also process it with other, other systems and as well. And so the reasons this is crucially important is not just sort of a performance optimization for, you know, reduced load time or something although sometimes I can help. Is because the data is too big to move nowadays you can't put it all in the database and suck it all out and move it over to some other system for working with, it's all two big, you can't move a petabyte, right? You have to bring the computation to the data as opposed to bring the data the computation, okay? And so these abilities to work with in situ data, you know, produce sort of in situ data is emerging as a sort of a key requirement in this big data era that was not such, not such a key requirement in the sort of era of relational databases. Okay.