MCPcopy Create free account
hub / github.com/apache/pig / testMapCombineReduce

Method testMapCombineReduce

test/org/apache/pig/test/TestCounters.java:283–339  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

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

Callers

nothing calls this directly

Calls 15

createInputFileMethod · 0.95
registerQueryMethod · 0.95
storeMethod · 0.95
getStatisticsMethod · 0.95
openMethod · 0.95
fullPathMethod · 0.95
getPigContextMethod · 0.95
getJobGraphMethod · 0.95
iteratorMethod · 0.95
getMapInputRecordsMethod · 0.95
getMapOutputRecordsMethod · 0.95
getReduceInputRecordsMethod · 0.95

Tested by

no test coverage detected