fix unit test
diff --git a/parquet-hadoop/src/test/java/parquet/hadoop/data/TestInputOutputFormat.java b/parquet-hadoop/src/test/java/parquet/hadoop/data/TestInputOutputFormat.java index 98715c3..9748d33 100644 --- a/parquet-hadoop/src/test/java/parquet/hadoop/data/TestInputOutputFormat.java +++ b/parquet-hadoop/src/test/java/parquet/hadoop/data/TestInputOutputFormat.java
@@ -18,8 +18,10 @@ import static java.lang.Thread.sleep; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertTrue; -import static parquet.hadoop.metadata.CompressionCodecName.GZIP; +import static parquet.hadoop.api.ReadSupport.PARQUET_READ_SCHEMA; +import static parquet.hadoop.data.GroupOutputFormat.setSchema; +import static parquet.io.api.Binary.fromString; +import static parquet.schema.MessageTypeParser.parseMessageType; import java.io.BufferedReader; import java.io.File; @@ -41,24 +43,15 @@ import parquet.data.Group; import parquet.data.GroupBuilder; import parquet.data.materializer.GroupBuilderImpl; -import parquet.hadoop.api.ReadSupport; -import parquet.hadoop.metadata.CompressionCodecName; import parquet.hadoop.util.ContextUtil; -import parquet.io.api.Binary; -import parquet.schema.MessageTypeParser; public class TestInputOutputFormat { private static final Log LOG = Log.getLog(TestInputOutputFormat.class); - private final Path parquetPath = new Path("target/test/parquet/hadoop/data/TestInputOutputFormat/parquet"); private final Path inputPath = new Path("src/test/java/parquet/hadoop/data/TestInputOutputFormat.java"); - private final Path outputPath = new Path("target/test/parquet/hadoop/data/TestInputOutputFormat/out"); - private Job writeJob; - private Job readJob; - private final String writeSchema = "message example {\n" + - "required int32 line;\n" + - "required binary content;\n" + - "}"; - private final String readSchema = "message example {\n" + + private final String outputBase = "target/test/parquet/hadoop/data/TestInputOutputFormat/"; + private final Path parquetPath = new Path(outputBase + "parquet"); + private final Path outputPath = new Path(outputBase + "out"); + private final String schema = "message example {\n" + "required int32 line;\n" + "required binary content;\n" + "}"; @@ -76,18 +69,20 @@ protected void map(LongWritable key, Text value, Mapper<LongWritable, Text, Void, Group>.Context context) throws java.io.IOException, InterruptedException { Group group = builder.startMessage() .addIntValue("line", (int)key.get()) - .addBinaryValue("content", Binary.fromByteArray(value.getBytes(), 0, value.getLength())) + .addBinaryValue("content", fromString(value.toString())) .endMessage(); // System.out.println("line: " + key.get()); // System.out.println("content: " + value.toString()); // System.out.println(group); +// System.out.println(group.getBinary("content").toStringUsingUTF8()); context.write(null, group); } } public static class WriteMapper extends Mapper<Void, Group, LongWritable, Text> { protected void map(Void key, Group value, Mapper<Void, Group, LongWritable, Text>.Context context) throws IOException, InterruptedException { - context.write(new LongWritable(value.getInt("line")), new Text(value.getBinary("content").getBytes())); +// System.out.println(value.getBinary("content").toStringUsingUTF8()); + context.write(new LongWritable(value.getInt("line")), new Text(value.getBinary("content").toStringUsingUTF8())); } } public static class PartialWriteMapper extends Mapper<Void, Group, LongWritable, Text> { @@ -96,39 +91,27 @@ } } - private void runMapReduceJob(CompressionCodecName codec, Class<? extends Mapper> writeMapperClass, String readSchema) throws IOException, ClassNotFoundException, InterruptedException { - runMapReduceJob(codec, new Configuration(), writeMapperClass, readSchema); - } - - private void runMapReduceJob(CompressionCodecName codec, Configuration conf, Class<? extends Mapper> writeMapperClass, String readSchema) throws IOException, ClassNotFoundException, InterruptedException { - + private void runMapReduceJob(Class<? extends Mapper<Void, Group, LongWritable, Text>> writeMapperClass, String readSchema) throws IOException, ClassNotFoundException, InterruptedException { + Configuration conf = new Configuration(); final FileSystem fileSystem = parquetPath.getFileSystem(conf); fileSystem.delete(parquetPath, true); fileSystem.delete(outputPath, true); { - writeJob = new Job(conf, "write"); + Job writeJob = new Job(conf, "write"); TextInputFormat.addInputPath(writeJob, inputPath); writeJob.setInputFormatClass(TextInputFormat.class); writeJob.setNumReduceTasks(0); - GroupOutputFormat.setCompression(writeJob, codec); GroupOutputFormat.setOutputPath(writeJob, parquetPath); writeJob.setOutputFormatClass(GroupOutputFormat.class); writeJob.setMapperClass(ReadMapper.class); - - GroupOutputFormat.setSchema( - writeJob, - MessageTypeParser.parseMessageType( - writeSchema)); + setSchema(writeJob, parseMessageType(schema)); writeJob.submit(); waitForJob(writeJob); } { - - conf.set(ReadSupport.PARQUET_READ_SCHEMA, readSchema); - readJob = new Job(conf, "read"); - + conf.set(PARQUET_READ_SCHEMA, readSchema); + Job readJob = new Job(conf, "read"); readJob.setInputFormatClass(GroupInputFormat.class); - GroupInputFormat.setInputPaths(readJob, parquetPath); readJob.setOutputFormatClass(TextOutputFormat.class); TextOutputFormat.setOutputPath(readJob, outputPath); @@ -139,8 +122,9 @@ } } - private void testReadWrite(CompressionCodecName codec) throws IOException, ClassNotFoundException, InterruptedException { - runMapReduceJob(codec, WriteMapper.class, readSchema); + @Test + public void testReadWrite() throws IOException, ClassNotFoundException, InterruptedException { + runMapReduceJob(WriteMapper.class, schema); final BufferedReader in = new BufferedReader(new FileReader(new File(inputPath.toString()))); final BufferedReader out = new BufferedReader(new FileReader(new File(outputPath.toString(), "part-m-00000"))); String lineIn; @@ -158,35 +142,8 @@ } @Test - public void testReadWrite() throws IOException, ClassNotFoundException, InterruptedException { - testReadWrite(CompressionCodecName.UNCOMPRESSED); - } - - @Test public void testProjection() throws Exception{ - runMapReduceJob(CompressionCodecName.GZIP, PartialWriteMapper.class, partialSchema); - } - - @Test - public void testReadWriteWithCounter() throws Exception { - runMapReduceJob(CompressionCodecName.GZIP, WriteMapper.class, readSchema); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("bytesread").getValue() > 0L); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("bytestotal").getValue() > 0L); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("bytesread").getValue() - == readJob.getCounters().getGroup("parquet").findCounter("bytestotal").getValue()); - //not testing the time read counter since it could be zero due to the size of data is too small - } - - @Test - public void testReadWriteWithoutCounter() throws Exception { - Configuration conf = new Configuration(); - conf.set("parquet.benchmark.time.read", "false"); - conf.set("parquet.benchmark.bytes.total", "false"); - conf.set("parquet.benchmark.bytes.read", "false"); - runMapReduceJob(GZIP, conf, WriteMapper.class, readSchema); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("bytesread").getValue() == 0L); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("bytestotal").getValue() == 0L); - assertTrue(readJob.getCounters().getGroup("parquet").findCounter("timeread").getValue() == 0L); + runMapReduceJob(PartialWriteMapper.class, partialSchema); } private void waitForJob(Job job) throws InterruptedException, IOException {