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));
+      }
+  }
 }


Reply via email to