Lets see how to emit double arrays from mapper and process them in reducer
DoubleArrayWritable class
public static class DoubleArrayWritable extends ArrayWritable { public DoubleArrayWritable() { super(DoubleWritable.class); } }
Driver()
job.setMapOutputKeyClass(IntWritable.class); job.setMapOutputValueClass(DoubleArrayWritable.class); job.setOutputKeyClass(NullWritable.class); job.setOutputValueClass(DoubleArrayWritable.class);
map()
import mywritable.DoubleArrayWritable;
public class MyMapper extends
Mapper<Object, Text, IntWritable, DoubleArrayWritable> {
public void map(Object key, Text value, Context context)
{
//Do something............
double[] arr = new double[size];
DoubleArrayWritable arrWritable = new DoubleArrayWritable();
DoubleWritable[] data = new DoubleWritable[size];
for (int k = 0; k < size; k++) {
data[k] = new DoubleWritable(arr[k]);
}
arrWritable.set(data);
context.write(mykey, arrWritable);
}
}
reduce()
import mywritable.DoubleArrayWritable;
public class MyReducer extends
Reducer<IntWritable, DoubleArrayWritable, NullWritable, Text> {
public void reduce(IntWritable key,
Iterable<DoubleArrayWritable> values, Context context){
double[] sum = new double[size];
for (DoubleArrayWritable c : values) {
Writable[] temp = new DoubleWritable[size];
temp = (c.get());
for (int i = 0; i < size; i++) {
sum[i] += Double.parseDouble(temp[i].toString());
}
//Do something and emit values ..................
context.write(out, new Text(emit));
}
}
}
Happy Hadooping.......