Author: cutting
Date: Mon Nov 11 06:54:42 2013
New Revision: 1540620
URL: http://svn.apache.org/r1540620
Log:
AVRO-1373. Java: Add support for xz compresssion codec, using LZMA2.
Contributed by Nick White.
Added:
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java
(with props)
Modified:
avro/trunk/CHANGES.txt
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/CodecFactory.java
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/DataFileConstants.java
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFile.java
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFileConcat.java
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/AvroOutputFormat.java
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/tether/TetherOutputFormat.java
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapreduce/AvroOutputFormatBase.java
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/CreateRandomFileTool.java
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/DataFileWriteTool.java
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/FromTextTool.java
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/RecodecTool.java
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/Util.java
Modified: avro/trunk/CHANGES.txt
URL:
http://svn.apache.org/viewvc/avro/trunk/CHANGES.txt?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
--- avro/trunk/CHANGES.txt (original)
+++ avro/trunk/CHANGES.txt Mon Nov 11 06:54:42 2013
@@ -9,6 +9,9 @@ Trunk (not yet released)
AVRO-1388. Java: Add fsync support to DataFileWriter.
(Hari Shreedharan via cutting)
+ AVRO-1373. Java: Add support for "xz" compresssion codec, using LZMA2.
+ (Nick White via cutting)
+
IMPROVEMENTS
AVRO-1355. Java: Reject schemas with duplicate field
Modified:
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/CodecFactory.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/CodecFactory.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/CodecFactory.java
(original)
+++
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/CodecFactory.java
Mon Nov 11 06:54:42 2013
@@ -22,6 +22,7 @@ import java.util.Map;
import java.util.zip.Deflater;
import org.apache.avro.AvroRuntimeException;
+import org.tukaani.xz.LZMA2Options;
/** Encapsulates the ability to specify and configure a compression codec.
*
@@ -48,6 +49,12 @@ public abstract class CodecFactory {
return new DeflateCodec.Option(compressionLevel);
}
+ /** XZ codec, with specific compression.
+ * compressionLevel should be between 1 and 9, inclusive. */
+ public static CodecFactory xzCodec(int compressionLevel) {
+ return new XZCodec.Option(compressionLevel);
+ }
+
/** Snappy codec.*/
public static CodecFactory snappyCodec() {
return new SnappyCodec.Option();
@@ -67,23 +74,26 @@ public abstract class CodecFactory {
private static final Map<String, CodecFactory> REGISTERED =
new HashMap<String, CodecFactory>();
- private static final int DEFAULT_DEFLATE_LEVEL =
Deflater.DEFAULT_COMPRESSION;
+ public static final int DEFAULT_DEFLATE_LEVEL = Deflater.DEFAULT_COMPRESSION;
+ public static final int DEFAULT_XZ_LEVEL = LZMA2Options.PRESET_DEFAULT;
static {
addCodec("null", nullCodec());
addCodec("deflate", deflateCodec(DEFAULT_DEFLATE_LEVEL));
addCodec("snappy", snappyCodec());
addCodec("bzip2", bzip2Codec());
+ addCodec("xz", xzCodec(DEFAULT_XZ_LEVEL));
}
/** Maps a codec name into a CodecFactory.
*
- * Currently there are four codecs registered by default:
+ * Currently there are five codecs registered by default:
* <ul>
* <li>{@code null}</li>
* <li>{@code deflate}</li>
* <li>{@code snappy}</li>
* <li>{@code bzip2}</li>
+ * <li>{@code xz}</li>
* </ul>
*/
public static CodecFactory fromString(String s) {
Modified:
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/DataFileConstants.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/DataFileConstants.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/DataFileConstants.java
(original)
+++
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/DataFileConstants.java
Mon Nov 11 06:54:42 2013
@@ -38,5 +38,6 @@ public class DataFileConstants {
public static final String DEFLATE_CODEC = "deflate";
public static final String SNAPPY_CODEC = "snappy";
public static final String BZIP2_CODEC = "bzip2";
+ public static final String XZ_CODEC = "xz";
}
Added: avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java?rev=1540620&view=auto
==============================================================================
--- avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java
(added)
+++ avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java
Mon Nov 11 06:54:42 2013
@@ -0,0 +1,122 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.avro.file;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.nio.ByteBuffer;
+
+import org.apache.commons.compress.compressors.xz.XZCompressorInputStream;
+import org.apache.commons.compress.compressors.xz.XZCompressorOutputStream;
+import org.apache.commons.compress.utils.IOUtils;
+
+/** * Implements xz compression and decompression. */
+public class XZCodec extends Codec {
+
+ static class Option extends CodecFactory {
+ private int compressionLevel;
+
+ Option(int compressionLevel) {
+ this.compressionLevel = compressionLevel;
+ }
+
+ @Override
+ protected Codec createInstance() {
+ return new XZCodec(compressionLevel);
+ }
+ }
+
+ private ByteArrayOutputStream outputBuffer;
+ private int compressionLevel;
+
+ public XZCodec(int compressionLevel) {
+ this.compressionLevel = compressionLevel;
+ }
+
+ @Override
+ public String getName() {
+ return DataFileConstants.XZ_CODEC;
+ }
+
+ @Override
+ public ByteBuffer compress(ByteBuffer data) throws IOException {
+ ByteArrayOutputStream baos = getOutputBuffer(data.remaining());
+ OutputStream ios = new XZCompressorOutputStream(baos, compressionLevel);
+ writeAndClose(data, ios);
+ return ByteBuffer.wrap(baos.toByteArray());
+ }
+
+ @Override
+ public ByteBuffer decompress(ByteBuffer data) throws IOException {
+ ByteArrayOutputStream baos = getOutputBuffer(data.remaining());
+ InputStream bytesIn = new ByteArrayInputStream(
+ data.array(),
+ data.arrayOffset() + data.position(),
+ data.remaining());
+ InputStream ios = new XZCompressorInputStream(bytesIn);
+ try {
+ IOUtils.copy(ios, baos);
+ } finally {
+ ios.close();
+ }
+ return ByteBuffer.wrap(baos.toByteArray());
+ }
+
+ private void writeAndClose(ByteBuffer data, OutputStream to) throws
IOException {
+ byte[] input = data.array();
+ int offset = data.arrayOffset() + data.position();
+ int length = data.remaining();
+ try {
+ to.write(input, offset, length);
+ } finally {
+ to.close();
+ }
+ }
+
+ // get and initialize the output buffer for use.
+ private ByteArrayOutputStream getOutputBuffer(int suggestedLength) {
+ if (null == outputBuffer) {
+ outputBuffer = new ByteArrayOutputStream(suggestedLength);
+ }
+ outputBuffer.reset();
+ return outputBuffer;
+ }
+
+ @Override
+ public int hashCode() {
+ return compressionLevel;
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj)
+ return true;
+ if (getClass() != obj.getClass())
+ return false;
+ XZCodec other = (XZCodec)obj;
+ return (this.compressionLevel == other.compressionLevel);
+ }
+
+ @Override
+ public String toString() {
+ return getName() + "-" + compressionLevel;
+ }
+}
Propchange:
avro/trunk/lang/java/avro/src/main/java/org/apache/avro/file/XZCodec.java
------------------------------------------------------------------------------
svn:eol-style = native
Modified:
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFile.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFile.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
--- avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFile.java
(original)
+++ avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFile.java
Mon Nov 11 06:54:42 2013
@@ -67,6 +67,9 @@ public class TestDataFile {
r.add(new Object[] { CodecFactory.deflateCodec(9) });
r.add(new Object[] { CodecFactory.nullCodec() });
r.add(new Object[] { CodecFactory.snappyCodec() });
+ r.add(new Object[] { CodecFactory.xzCodec(0) });
+ r.add(new Object[] { CodecFactory.xzCodec(1) });
+ r.add(new Object[] { CodecFactory.xzCodec(6) });
return r;
}
@@ -81,7 +84,7 @@ public class TestDataFile {
"{\"type\": \"record\", \"name\": \"Test\", \"fields\": ["
+"{\"name\":\"stringField\", \"type\":\"string\"},"
+"{\"name\":\"longField\", \"type\":\"long\"}]}";
- private static final Schema SCHEMA = Schema.parse(SCHEMA_JSON);
+ private static final Schema SCHEMA = new Schema.Parser().parse(SCHEMA_JSON);
private File makeFile() {
return new File(DIR, "test-" + codec + ".avro");
Modified:
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFileConcat.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFileConcat.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFileConcat.java
(original)
+++
avro/trunk/lang/java/avro/src/test/java/org/apache/avro/TestDataFileConcat.java
Mon Nov 11 06:54:42 2013
@@ -65,6 +65,14 @@ public class TestDataFileConcat {
{ CodecFactory.deflateCodec(3), CodecFactory.nullCodec(), false });
r.add(new Object[]
{ CodecFactory.nullCodec(), CodecFactory.deflateCodec(6), false });
+ r.add(new Object[]
+ { CodecFactory.xzCodec(1), CodecFactory.xzCodec(2), false });
+ r.add(new Object[]
+ { CodecFactory.xzCodec(1), CodecFactory.xzCodec(2), true });
+ r.add(new Object[]
+ { CodecFactory.xzCodec(2), CodecFactory.nullCodec(), false });
+ r.add(new Object[]
+ { CodecFactory.nullCodec(), CodecFactory.xzCodec(2), false });
return r;
}
@@ -82,14 +90,14 @@ public class TestDataFileConcat {
","
+"{\"name\":\"longField\", \"type\":\"long\"}" +
"]}";
- private static final Schema SCHEMA = Schema.parse(SCHEMA_JSON);
+ private static final Schema SCHEMA = new Schema.Parser().parse(SCHEMA_JSON);
private File makeFile(String name) {
return new File(DIR, "test-" + name + ".avro");
}
@Test
- public void testConcateateFiles() throws IOException {
+ public void testConcatenateFiles() throws IOException {
System.out.println("SEED = "+SEED);
System.out.println("COUNT = "+COUNT);
for (int k = 0; k < 60; k++) {
Modified:
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/AvroOutputFormat.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/AvroOutputFormat.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/AvroOutputFormat.java
(original)
+++
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/AvroOutputFormat.java
Mon Nov 11 06:54:42 2013
@@ -40,6 +40,9 @@ import org.apache.avro.hadoop.file.Hadoo
import static org.apache.avro.file.DataFileConstants.DEFAULT_SYNC_INTERVAL;
import static org.apache.avro.file.DataFileConstants.DEFLATE_CODEC;
+import static org.apache.avro.file.DataFileConstants.XZ_CODEC;
+import static org.apache.avro.file.CodecFactory.DEFAULT_DEFLATE_LEVEL;
+import static org.apache.avro.file.CodecFactory.DEFAULT_XZ_LEVEL;
/**
* An {@link org.apache.hadoop.mapred.OutputFormat} for Avro data files.
@@ -57,12 +60,12 @@ public class AvroOutputFormat <T>
/** The configuration key for Avro deflate level. */
public static final String DEFLATE_LEVEL_KEY = "avro.mapred.deflate.level";
+ /** The configuration key for Avro XZ level. */
+ public static final String XZ_LEVEL_KEY = "avro.mapred.xz.level";
+
/** The configuration key for Avro sync interval. */
public static final String SYNC_INTERVAL_KEY = "avro.mapred.sync.interval";
- /** The default deflate level. */
- public static final int DEFAULT_DEFLATE_LEVEL = 1;
-
/** Enable output compression using the deflate codec and specify its
level.*/
public static void setDeflateLevel(JobConf job, int level) {
FileOutputFormat.setCompressOutput(job, true);
@@ -110,7 +113,8 @@ public class AvroOutputFormat <T>
CodecFactory factory = null;
if (FileOutputFormat.getCompressOutput(job)) {
- int level = job.getInt(DEFLATE_LEVEL_KEY, DEFAULT_DEFLATE_LEVEL);
+ int deflateLevel = job.getInt(DEFLATE_LEVEL_KEY, DEFAULT_DEFLATE_LEVEL);
+ int xzLevel = job.getInt(XZ_LEVEL_KEY, DEFAULT_XZ_LEVEL);
String codecName = job.get(AvroJob.OUTPUT_CODEC);
if (codecName == null) {
@@ -121,11 +125,13 @@ public class AvroOutputFormat <T>
job.set(AvroJob.OUTPUT_CODEC , avroCodecName);
return factory;
} else {
- return CodecFactory.deflateCodec(level);
+ return CodecFactory.deflateCodec(deflateLevel);
}
} else {
if ( codecName.equals(DEFLATE_CODEC)) {
- factory = CodecFactory.deflateCodec(level);
+ factory = CodecFactory.deflateCodec(deflateLevel);
+ } else if ( codecName.equals(XZ_CODEC)) {
+ factory = CodecFactory.xzCodec(xzLevel);
} else {
factory = CodecFactory.fromString(codecName);
}
Modified:
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/tether/TetherOutputFormat.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/tether/TetherOutputFormat.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/tether/TetherOutputFormat.java
(original)
+++
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapred/tether/TetherOutputFormat.java
Mon Nov 11 06:54:42 2013
@@ -58,7 +58,7 @@ class TetherOutputFormat
if (FileOutputFormat.getCompressOutput(job)) {
int level = job.getInt(AvroOutputFormat.DEFLATE_LEVEL_KEY,
- AvroOutputFormat.DEFAULT_DEFLATE_LEVEL);
+ CodecFactory.DEFAULT_DEFLATE_LEVEL);
writer.setCodec(CodecFactory.deflateCodec(level));
}
Modified:
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapreduce/AvroOutputFormatBase.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapreduce/AvroOutputFormatBase.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapreduce/AvroOutputFormatBase.java
(original)
+++
avro/trunk/lang/java/mapred/src/main/java/org/apache/avro/mapreduce/AvroOutputFormatBase.java
Mon Nov 11 06:54:42 2013
@@ -46,9 +46,12 @@ public abstract class AvroOutputFormatBa
protected static CodecFactory getCompressionCodec(TaskAttemptContext
context) {
if (FileOutputFormat.getCompressOutput(context)) {
// Default to deflate compression.
- int compressionLevel = context.getConfiguration().getInt(
+ int deflateLevel = context.getConfiguration().getInt(
org.apache.avro.mapred.AvroOutputFormat.DEFLATE_LEVEL_KEY,
- org.apache.avro.mapred.AvroOutputFormat.DEFAULT_DEFLATE_LEVEL);
+ CodecFactory.DEFAULT_DEFLATE_LEVEL);
+ int xzLevel = context.getConfiguration().getInt(
+ org.apache.avro.mapred.AvroOutputFormat.XZ_LEVEL_KEY,
+ CodecFactory.DEFAULT_XZ_LEVEL);
String outputCodec = context.getConfiguration()
.get(AvroJob.CONF_OUTPUT_CODEC);
@@ -60,10 +63,12 @@ public abstract class AvroOutputFormatBa
context.getConfiguration().set(AvroJob.CONF_OUTPUT_CODEC,
avroCodecName);
return HadoopCodecFactory.fromHadoopString(compressionCodec);
} else {
- return CodecFactory.deflateCodec(compressionLevel);
+ return CodecFactory.deflateCodec(deflateLevel);
}
} else if (DataFileConstants.DEFLATE_CODEC.equals(outputCodec)) {
- return CodecFactory.deflateCodec(compressionLevel);
+ return CodecFactory.deflateCodec(deflateLevel);
+ } else if (DataFileConstants.XZ_CODEC.equals(outputCodec)) {
+ return CodecFactory.xzCodec(xzLevel);
} else {
return CodecFactory.fromString(outputCodec);
}
Modified:
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/CreateRandomFileTool.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/CreateRandomFileTool.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/CreateRandomFileTool.java
(original)
+++
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/CreateRandomFileTool.java
Mon Nov 11 06:54:42 2013
@@ -26,7 +26,6 @@ import joptsimple.OptionSet;
import joptsimple.OptionSpec;
import org.apache.avro.Schema;
-import org.apache.avro.file.CodecFactory;
import org.apache.avro.file.DataFileWriter;
import org.apache.avro.generic.GenericDatumWriter;
import org.apache.trevni.avro.RandomData;
@@ -53,11 +52,8 @@ public class CreateRandomFileTool implem
p.accepts("count", "Record Count")
.withRequiredArg()
.ofType(Integer.class);
- OptionSpec<String> codec =
- p.accepts("codec", "Compression codec")
- .withRequiredArg()
- .defaultsTo("null")
- .ofType(String.class);
+ OptionSpec<String> codec = Util.compressionCodecOption(p);
+ OptionSpec<Integer> level = Util.compressionLevelOption(p);
OptionSpec<String> file =
p.accepts("schema-file", "Schema File")
.withOptionalArg()
@@ -87,7 +83,7 @@ public class CreateRandomFileTool implem
DataFileWriter<Object> writer =
new DataFileWriter<Object>(new GenericDatumWriter<Object>());
- writer.setCodec(CodecFactory.fromString(codec.value(opts)));
+ writer.setCodec(Util.codecFactory(opts, codec, level));
writer.create(schema, Util.fileOrStdout(args.get(0), out));
for (Object datum : new RandomData(schema, (int)count.value(opts)))
Modified:
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/DataFileWriteTool.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/DataFileWriteTool.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/DataFileWriteTool.java
(original)
+++
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/DataFileWriteTool.java
Mon Nov 11 06:54:42 2013
@@ -28,7 +28,7 @@ import joptsimple.OptionSet;
import joptsimple.OptionSpec;
import org.apache.avro.Schema;
-import org.apache.avro.file.CodecFactory;
+import org.apache.avro.file.DataFileConstants;
import org.apache.avro.file.DataFileWriter;
import org.apache.avro.generic.GenericDatumReader;
import org.apache.avro.generic.GenericDatumWriter;
@@ -54,11 +54,8 @@ public class DataFileWriteTool implement
List<String> args) throws Exception {
OptionParser p = new OptionParser();
- OptionSpec<String> codec =
- p.accepts("codec", "Compression codec")
- .withRequiredArg()
- .defaultsTo("null")
- .ofType(String.class);
+ OptionSpec<String> codec = Util.compressionCodecOption(p);
+ OptionSpec<Integer> level = Util.compressionLevelOption(p);
OptionSpec<String> file =
p.accepts("schema-file", "Schema File")
.withOptionalArg()
@@ -93,7 +90,7 @@ public class DataFileWriteTool implement
DataInputStream din = new DataInputStream(input);
DataFileWriter<Object> writer =
new DataFileWriter<Object>(new GenericDatumWriter<Object>());
- writer.setCodec(CodecFactory.fromString(codec.value(opts)));
+ writer.setCodec(Util.codecFactory(opts, codec, level,
DataFileConstants.NULL_CODEC));
writer.create(schema, out);
Decoder decoder = DecoderFactory.get().jsonDecoder(schema, din);
Object datum;
Modified:
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/FromTextTool.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/FromTextTool.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/FromTextTool.java
(original)
+++
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/FromTextTool.java
Mon Nov 11 06:54:42 2013
@@ -32,7 +32,6 @@ import org.apache.avro.Schema;
import org.apache.avro.file.CodecFactory;
import org.apache.avro.file.DataFileWriter;
import org.apache.avro.generic.GenericDatumWriter;
-import static org.apache.avro.file.DataFileConstants.DEFLATE_CODEC;
/** Reads a text file into an Avro data file.
*
@@ -57,11 +56,8 @@ public class FromTextTool implements Too
List<String> args) throws Exception {
OptionParser p = new OptionParser();
- OptionSpec<Integer> level = p.accepts("level", "compression level")
- .withOptionalArg().ofType(Integer.class);
-
- OptionSpec<String> codec = p.accepts("codec", "compression codec")
- .withOptionalArg().ofType(String.class);
+ OptionSpec<Integer> level = Util.compressionLevelOption(p);
+ OptionSpec<String> codec = Util.compressionCodecOption(p);
OptionSet opts = p.parse(args.toArray(new String[0]));
@@ -73,18 +69,8 @@ public class FromTextTool implements Too
return 1;
}
- int compressionLevel = 1; // Default compression level
- if (opts.hasArgument(level)) {
- compressionLevel = level.value(opts);
- }
+ CodecFactory codecFactory = Util.codecFactory(opts, codec, level);
- String codecName = opts.hasArgument(codec)
- ? codec.value(opts)
- : DEFLATE_CODEC;
- CodecFactory codecFactory = codecName.equals(DEFLATE_CODEC)
- ? CodecFactory.deflateCodec(compressionLevel)
- : CodecFactory.fromString(codecName);
-
BufferedInputStream inStream = Util.fileOrStdin(nargs.get(0), stdin);
BufferedOutputStream outStream = Util.fileOrStdout(nargs.get(1), out);
Modified:
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/RecodecTool.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/RecodecTool.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
---
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/RecodecTool.java
(original)
+++
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/RecodecTool.java
Mon Nov 11 06:54:42 2013
@@ -21,7 +21,6 @@ import java.io.InputStream;
import java.io.OutputStream;
import java.io.PrintStream;
import java.util.List;
-import java.util.zip.Deflater;
import joptsimple.OptionParser;
import joptsimple.OptionSet;
@@ -29,6 +28,7 @@ import joptsimple.OptionSpec;
import org.apache.avro.Schema;
import org.apache.avro.file.CodecFactory;
+import org.apache.avro.file.DataFileConstants;
import org.apache.avro.file.DataFileStream;
import org.apache.avro.file.DataFileWriter;
import org.apache.avro.generic.GenericDatumReader;
@@ -42,16 +42,8 @@ public class RecodecTool implements Tool
List<String> args) throws Exception {
OptionParser optParser = new OptionParser();
- OptionSpec<String> codecOpt = optParser
- .accepts("codec", "Compression codec")
- .withRequiredArg()
- .defaultsTo("null")
- .ofType(String.class);
- OptionSpec<String> levelOpt = optParser
- .accepts("level", "Compression level (only applies to deflate)")
- .withRequiredArg()
- .defaultsTo("" + Deflater.DEFAULT_COMPRESSION)
- .ofType(String.class);
+ OptionSpec<String> codecOpt = Util.compressionCodecOption(optParser);
+ OptionSpec<Integer> levelOpt = Util.compressionLevelOption(optParser);
OptionSet opts = optParser.parse(args.toArray(new String[0]));
List<String> nargs = opts.nonOptionArguments();
@@ -78,9 +70,8 @@ public class RecodecTool implements Tool
Schema schema = reader.getSchema();
DataFileWriter<GenericRecord> writer = new DataFileWriter<GenericRecord>(
new GenericDatumWriter<GenericRecord>());
- CodecFactory codec = opts.valueOf(codecOpt).equals("deflate")
- ? CodecFactory.deflateCodec(Integer.parseInt(levelOpt.value(opts)))
- : CodecFactory.fromString(codecOpt.value(opts));
+ // unlike the other Avro tools, we default to a null codec, not deflate
+ CodecFactory codec = Util.codecFactory(opts, codecOpt, levelOpt,
DataFileConstants.NULL_CODEC);
writer.setCodec(codec);
for (String key : reader.getMetaKeys()) {
if (!DataFileWriter.isReservedMeta(key)) {
Modified:
avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/Util.java
URL:
http://svn.apache.org/viewvc/avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/Util.java?rev=1540620&r1=1540619&r2=1540620&view=diff
==============================================================================
--- avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/Util.java
(original)
+++ avro/trunk/lang/java/tools/src/main/java/org/apache/avro/tool/Util.java Mon
Nov 11 06:54:42 2013
@@ -17,6 +17,8 @@
*/
package org.apache.avro.tool;
+import static org.apache.avro.file.DataFileConstants.DEFLATE_CODEC;
+
import java.io.BufferedInputStream;
import java.io.BufferedOutputStream;
import java.io.File;
@@ -26,8 +28,11 @@ import java.io.OutputStream;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.zip.Deflater;
import org.apache.avro.Schema;
+import org.apache.avro.file.CodecFactory;
+import org.apache.avro.file.DataFileConstants;
import org.apache.avro.file.DataFileReader;
import org.apache.avro.generic.GenericDatumReader;
import org.apache.avro.io.DecoderFactory;
@@ -37,6 +42,10 @@ import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
+import joptsimple.OptionSet;
+import joptsimple.OptionParser;
+import joptsimple.OptionSpec;
+
/** Static utility methods for tools. */
class Util {
/**
@@ -218,4 +227,36 @@ class Util {
}
}
+ static OptionSpec<String> compressionCodecOption(OptionParser optParser) {
+ return optParser
+ .accepts("codec", "Compression codec")
+ .withRequiredArg()
+ .ofType(String.class)
+ .defaultsTo("null");
+}
+
+ static OptionSpec<Integer> compressionLevelOption(OptionParser optParser) {
+ return optParser
+ .accepts("level", "Compression level (only applies to deflate and xz)")
+ .withRequiredArg()
+ .ofType(Integer.class)
+ .defaultsTo(Deflater.DEFAULT_COMPRESSION);
+ }
+
+ static CodecFactory codecFactory(OptionSet opts, OptionSpec<String> codec,
OptionSpec<Integer> level) {
+ return codecFactory(opts, codec, level, DEFLATE_CODEC);
+ }
+
+ static CodecFactory codecFactory(OptionSet opts, OptionSpec<String> codec,
OptionSpec<Integer> level, String defaultCodec) {
+ String codecName = opts.hasArgument(codec)
+ ? codec.value(opts)
+ : defaultCodec;
+ if(codecName.equals(DEFLATE_CODEC)) {
+ return CodecFactory.deflateCodec(level.value(opts));
+ } else if(codecName.equals(DataFileConstants.XZ_CODEC)) {
+ return CodecFactory.xzCodec(level.value(opts));
+ } else {
+ return CodecFactory.fromString(codec.value(opts));
+ }
+ }
}