【问题标题】:Using List as Value in MapReduce Returns Identical Values在 MapReduce 中使用列表作为值返回相同的值
【发布时间】:2015-08-27 04:27:45
【问题描述】:

我有一个 MapReduce 作业,它输出一个 IntWritable 作为键和 Point(我创建的实现可写的对象)对象作为 map 函数的值。然后在 reduce 函数中,我使用 for-each 循环遍历 Points 的可迭代对象来创建列表:

@Override
public void reduce(IntWritable key, Iterable<Point> points, Context context) throws IOException, InterruptedException {

    List<Point> pointList = new ArrayList<>();
    for (Point point : points) {
        pointList.add(point);
    }
    context.write(key, pointList);
}

问题是这个列表的大小是正确的,但每个点都是完全相同的。我的 Point 类中的字段不是静态的,我在循环中单独打印了每个点,以确保这些点是唯一的(它们是唯一的)。此外,我创建了一个单独的类,它只创建几个点并将它们添加到列表中,这似乎有效,这意味着 MapReduce 做了一些我不知道的事情。

任何解决此问题的帮助将不胜感激。

更新: Mapper 类代码:

private final static IntWritable one = new IntWritable(1);
private Text word = new Text();
private IntWritable firstChar = new IntWritable();
private Point point = new Point();

@Override
public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
    String line = value.toString();
    StringTokenizer tokenizer = new StringTokenizer(line, " ");

    while(tokenizer.hasMoreTokens()) {
        String atts = tokenizer.nextToken();
        String cut = atts.substring(1, atts.length() - 1);
        String[] nums = cut.split(",");

        point.set(Double.parseDouble(nums[0]), Double.parseDouble(nums[1]), Double.parseDouble(nums[2]), Double.parseDouble(nums[3]));
        context.write(one, point);
    }
}

点类:

public class Point implements Writable {

public Double att1;
public Double att2;
public Double att3;
public Double att4;

public Point() {

}

public void set(Double att1, Double att2, Double att3, Double att4) {
    this.att1 = att1;
    this.att2 = att2;
    this.att3 = att3;
    this.att4 = att4;
}

@Override
public void write(DataOutput dataOutput) throws IOException {
    dataOutput.writeDouble(att1);
    dataOutput.writeDouble(att2);
    dataOutput.writeDouble(att3);
    dataOutput.writeDouble(att4);
}

@Override
public void readFields(DataInput dataInput) throws IOException {
    this.att1 = dataInput.readDouble();
    this.att2 = dataInput.readDouble();
    this.att3 = dataInput.readDouble();
    this.att4 = dataInput.readDouble();
}

@Override
public String toString() {
    String output = "{" + att1 + ", " + att2 + ", " + att3 + ", " + att4 + "}";
    return output;
}

【问题讨论】:

  • 请按照您在map中设置和在reduce中检索的方式添加map和reduce的代码。也是实现Writable的点类
  • 刚刚用 Point 和 Mapper 类更新了帖子。上面的所有代码都是每个类中的所有内容。
  • 尝试移动Point point = new Point();在地图内部,并采用 context.write(one, point);在while循环之外。

标签: java list hadoop mapreduce reduce


【解决方案1】:

问题出在你的减速器上。您不想将所有点存储在内存中。它们可能很大,Hadoop 为您解决了这个问题(尽管方式很尴尬)。

当循环通过给定的Iterable&lt;Points&gt; 时,每个Point 实例都会被重复使用,因此它只在给定时间保留一个实例。

这意味着当你拨打points.next()时,会发生以下两件事:

  1. Point实例被重用并设置下一个点数据
  2. Key 实例也是如此。

在您的情况下,您会在 List 中找到多次插入的 Point 的一个实例,并使用最后一个 Point 中的数据进行设置。

你不应该在你的 reducer 中保存 Writables 的实例,或者应该克隆它们。

您可以在此处阅读有关此问题的更多信息
https://cornercases.wordpress.com/2011/08/18/hadoop-object-reuse-pitfall-all-my-reducer-values-are-the-same/

【讨论】:

  • 问题是,当我应用它时,我想要将每个点与迭代中的其他点进行比较,所以我需要能够存储它们并返回它们。有没有办法做到这一点?
  • 您不想将它们存储在内存中。正如我所说,MapReduce 是大数据处理工具 - 值可能不适合内存。使用 Point 作为 key 怎么样?然后你会得到相同的点在reducer中分组和排序。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-04-10
  • 1970-01-01
  • 2019-04-09
  • 1970-01-01
  • 2015-03-17
相关资源
最近更新 更多