()
| 1007 | } |
| 1008 | |
| 1009 | @Test |
| 1010 | public void testMergeJoin() throws Exception { |
| 1011 | String query = "a = load '/tmp/input1';" + |
| 1012 | "b = load '/tmp/input2';" + |
| 1013 | "c = join a by $0, b by $0 using 'merge';" + |
| 1014 | "store c into '/tmp/output1';"; |
| 1015 | |
| 1016 | PhysicalPlan pp = Util.buildPp(pigServer, query); |
| 1017 | MRCompiler comp = new MRCompiler(pp, pc); |
| 1018 | comp.compile(); |
| 1019 | MROperPlan mrp = comp.getMRPlan(); |
| 1020 | assertTrue(mrp.size()==2); |
| 1021 | |
| 1022 | MapReduceOper mrOp0 = mrp.getRoots().get(0); |
| 1023 | assertTrue(mrOp0.mapPlan.size()==2); |
| 1024 | PhysicalOperator load0 = mrOp0.mapPlan.getRoots().get(0); |
| 1025 | MergeJoinIndexer func = (MergeJoinIndexer)PigContext.instantiateFuncFromSpec(((POLoad)load0).getLFile().getFuncSpec()); |
| 1026 | Field lrField = MergeJoinIndexer.class.getDeclaredField("lr"); |
| 1027 | lrField.setAccessible(true); |
| 1028 | POLocalRearrange lr = (POLocalRearrange)lrField.get(func); |
| 1029 | List<PhysicalPlan> innerPlans = lr.getPlans(); |
| 1030 | PhysicalOperator localrearrange0 = mrOp0.mapPlan.getSuccessors(load0).get(0); |
| 1031 | assertTrue(localrearrange0 instanceof POLocalRearrange); |
| 1032 | assertTrue(mrOp0.reducePlan.size()==3); |
| 1033 | PhysicalOperator pack0 = mrOp0.reducePlan.getRoots().get(0); |
| 1034 | assertTrue(pack0 instanceof POPackage); |
| 1035 | PhysicalOperator foreach0 = mrOp0.reducePlan.getSuccessors(pack0).get(0); |
| 1036 | assertTrue(foreach0 instanceof POForEach); |
| 1037 | PhysicalOperator store0 = mrOp0.reducePlan.getSuccessors(foreach0).get(0); |
| 1038 | assertTrue(store0 instanceof POStore); |
| 1039 | |
| 1040 | assertTrue(innerPlans.size()==1); |
| 1041 | PhysicalPlan innerPlan = innerPlans.get(0); |
| 1042 | assertTrue(innerPlan.size()==1); |
| 1043 | PhysicalOperator project = innerPlan.getRoots().get(0); |
| 1044 | assertTrue(project instanceof POProject); |
| 1045 | assertTrue(((POProject)project).getColumn()==0); |
| 1046 | |
| 1047 | MapReduceOper mrOp1 = mrp.getSuccessors(mrOp0).get(0); |
| 1048 | assertTrue(mrOp1.mapPlan.size()==3); |
| 1049 | PhysicalOperator load1 = mrOp1.mapPlan.getRoots().get(0); |
| 1050 | assertTrue(load1 instanceof POLoad); |
| 1051 | PhysicalOperator mergejoin1 = mrOp1.mapPlan.getSuccessors(load1).get(0); |
| 1052 | assertTrue(mergejoin1 instanceof POMergeJoin); |
| 1053 | PhysicalOperator store1 = mrOp1.mapPlan.getSuccessors(mergejoin1).get(0); |
| 1054 | assertTrue(store1 instanceof POStore); |
| 1055 | assertTrue(mrOp1.reducePlan.isEmpty()); |
| 1056 | } |
| 1057 | |
| 1058 | public static class WeirdComparator extends ComparisonFunc { |
| 1059 | @Override |
nothing calls this directly
no test coverage detected