()
| 281 | } |
| 282 | |
| 283 | @Test |
| 284 | public void testMapCombineReduce() throws IOException, ExecException { |
| 285 | int count = 0; |
| 286 | PrintWriter pw = new PrintWriter(Util.createInputFile(cluster, file)); |
| 287 | int [] nos = new int[10]; |
| 288 | for(int i = 0; i < 10; i++) |
| 289 | nos[i] = 0; |
| 290 | |
| 291 | for(int i = 0; i < MAX; i++) { |
| 292 | int index = r.nextInt(10); |
| 293 | int value = r.nextInt(100); |
| 294 | nos[index] += value; |
| 295 | pw.println(index + "\t" + value); |
| 296 | } |
| 297 | pw.close(); |
| 298 | |
| 299 | for(int i = 0; i < 10; i++) { |
| 300 | if(nos[i] > 0) count ++; |
| 301 | } |
| 302 | |
| 303 | PigServer pigServer = new PigServer(cluster.getExecType(), cluster.getProperties()); |
| 304 | pigServer.registerQuery("a = load '" + file + "';"); |
| 305 | pigServer.registerQuery("b = group a by $0;"); |
| 306 | pigServer.registerQuery("c = foreach b generate group, SUM(a.$1);"); |
| 307 | ExecJob job = pigServer.store("c", "output"); |
| 308 | PigStats pigStats = job.getStatistics(); |
| 309 | |
| 310 | InputStream is = FileLocalizer.open(FileLocalizer.fullPath("output", |
| 311 | pigServer.getPigContext()), pigServer.getPigContext()); |
| 312 | long filesize = 0; |
| 313 | while(is.read() != -1) filesize++; |
| 314 | |
| 315 | is.close(); |
| 316 | |
| 317 | cluster.getFileSystem().delete(new Path(file), true); |
| 318 | cluster.getFileSystem().delete(new Path("output"), true); |
| 319 | |
| 320 | System.out.println("============================================"); |
| 321 | System.out.println("Test case MapCombineReduce"); |
| 322 | System.out.println("============================================"); |
| 323 | |
| 324 | JobGraph jp = pigStats.getJobGraph(); |
| 325 | Iterator<JobStats> iter = jp.iterator(); |
| 326 | while (iter.hasNext()) { |
| 327 | JobStats js = iter.next(); |
| 328 | System.out.println("Map input records : " + js.getMapInputRecords()); |
| 329 | assertEquals(MAX, js.getMapInputRecords()); |
| 330 | System.out.println("Map output records : " + js.getMapOutputRecords()); |
| 331 | assertEquals(MAX, js.getMapOutputRecords()); |
| 332 | System.out.println("Reduce input records : " + js.getReduceInputRecords()); |
| 333 | assertEquals(count, js.getReduceInputRecords()); |
| 334 | System.out.println("Reduce output records : " + js.getReduceOutputRecords()); |
| 335 | assertEquals(count, js.getReduceOutputRecords()); |
| 336 | } |
| 337 | System.out.println("Hdfs bytes written : " + pigStats.getBytesWritten()); |
| 338 | assertEquals(filesize, pigStats.getBytesWritten()); |
| 339 | } |
| 340 |
nothing calls this directly
no test coverage detected