Author: omalley
Date: Tue Mar  8 06:01:03 2011
New Revision: 1079257

URL: http://svn.apache.org/viewvc?rev=1079257&view=rev
Log:
commit fb43edf954fb53cf8f9b86c207de6374707ee6fd
Author: Greg Roelofs <[email protected]>
Date:   Mon Jan 31 22:21:10 2011 -0800

    Fix : ubertasks fail if map outputs are compressed.  "codec" arg
    to Merger.merge() was never set...

Modified:
    
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/Merger.java
    
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/ReduceTask.java
    
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/UberTask.java

Modified: 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/Merger.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/Merger.java?rev=1079257&r1=1079256&r2=1079257&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/Merger.java
 (original)
+++ 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/Merger.java
 Tue Mar  8 06:01:03 2011
@@ -44,7 +44,7 @@ import org.apache.hadoop.util.Progress;
 import org.apache.hadoop.util.Progressable;
 
 /**
- * Merger is an utility class used by the Map and Reduce tasks for merging
+ * Merger is a utility class used by the Map and Reduce tasks for merging
  * both their memory and disk segments
  */
 @InterfaceAudience.Private
@@ -150,7 +150,7 @@ public class Merger {  
   }
 
   public static <K extends Object, V extends Object>
-    RawKeyValueIterator merge(Configuration conf, FileSystem fs,
+  RawKeyValueIterator merge(Configuration conf, FileSystem fs,
                             Class<K> keyClass, Class<V> valueClass,
                             List<Segment<K, V>> segments,
                             int mergeFactor, int inMemSegments, Path tmpDir,
@@ -180,14 +180,14 @@ public class Merger {  
                           Counters.Counter readsCounter,
                           Counters.Counter writesCounter,
                           Progress mergePhase)
-    throws IOException {
-  return new MergeQueue<K, V>(conf, fs, segments, comparator, reporter,
-                         sortSegments, codec).merge(keyClass, valueClass,
-                                             mergeFactor, inMemSegments,
-                                             tmpDir,
-                                             readsCounter, writesCounter,
-                                             mergePhase);
-}
+      throws IOException {
+    return new MergeQueue<K, V>(conf, fs, segments, comparator, reporter,
+                           sortSegments, codec).merge(keyClass, valueClass,
+                                               mergeFactor, inMemSegments,
+                                               tmpDir,
+                                               readsCounter, writesCounter,
+                                               mergePhase);
+  }
 
   public static <K extends Object, V extends Object>
   void writeFile(RawKeyValueIterator records, Writer<K, V> writer, 
@@ -230,7 +230,7 @@ public class Merger {  
     public Segment(Configuration conf, FileSystem fs, Path file,
                    CompressionCodec codec, boolean preserve,
                    Counters.Counter mergedMapOutputsCounter)
-  throws IOException {
+    throws IOException {
       this(conf, fs, file, 0, fs.getFileStatus(file).getLen(), codec, 
preserve, 
            mergedMapOutputsCounter);
     }
@@ -416,7 +416,9 @@ public class Merger {  
       this.reporter = reporter;
       
       for (Path file : inputs) {
-        LOG.debug("MergeQ: adding: " + file);
+        if (LOG.isDebugEnabled()) {
+          LOG.debug("MergeQ: adding: " + file);
+        }
         segments.add(new Segment<K, V>(conf, fs, file, codec, !deleteInputs, 
                                        (file.toString().endsWith(
                                            Task.MERGED_OUTPUT_PREFIX) ? 

Modified: 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/ReduceTask.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/ReduceTask.java?rev=1079257&r1=1079256&r2=1079257&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/ReduceTask.java
 (original)
+++ 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/ReduceTask.java
 Tue Mar  8 06:01:03 2011
@@ -134,7 +134,7 @@ public class ReduceTask extends Task {
                                            getCounters());
   }
 
-  private CompressionCodec initCodec() {
+  static CompressionCodec initCodec(JobConf conf) {
     // check if map-outputs are to be compressed
     if (conf.getCompressMapOutput()) {
       Class<? extends CompressionCodec> codecClass =
@@ -390,7 +390,7 @@ public class ReduceTask extends Task {
     }
     
     // Initialize the codec
-    codec = initCodec();
+    codec = initCodec(conf);
     RawKeyValueIterator rIter = null;
     boolean isLocal = "local".equals(job.get(JTConfig.JT_IPC_ADDRESS, 
"local"));
     if (!isLocal) {

Modified: 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/UberTask.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/UberTask.java?rev=1079257&r1=1079256&r2=1079257&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/UberTask.java
 (original)
+++ 
hadoop/mapreduce/branches/yahoo-merge/src/java/org/apache/hadoop/mapred/UberTask.java
 Tue Mar  8 06:01:03 2011
@@ -375,7 +375,7 @@ class UberTask extends Task {
     final FileSystem rfs = FileSystem.getLocal(job).getRaw();
     RawKeyValueIterator rIter =
         Merger.merge(job, rfs, job.getMapOutputKeyClass(),
-                     job.getMapOutputValueClass(), null,   // no codec
+                     job.getMapOutputValueClass(), reduce.initCodec(localConf),
                      ReduceTask.getMapFiles(reduce, rfs, true),
                      !conf.getKeepFailedTaskFiles(),
                      job.getInt(JobContext.IO_SORT_FACTOR, 100),


Reply via email to