【问题标题】:Hadoop MR2: Records with same key are processed independentlyHadoop MR2:具有相同键的记录独立处理
【发布时间】:2014-12-03 21:03:15
【问题描述】:

我有以下设置:映射器输出键类型为K1,值类型为V1K1WritableComparable 的记录。因此,组合器将K1Iterable<V1> 作为其输入。然后它进行聚合并准确输出一个K1, V1 记录。 reducer 从组合器获取输入,同样是K1, Iterable<V1>。据我了解,在 Reduce 阶段,每个人 K1 必须恰好存在一对 K1, Iterable<V1>。然后reducer 正好输出一个K2, V2K2 又是WritableComparable

我现在的问题是:我的输出文件中有多个K2, V2,即使在同一个文件中!我的关键类的比较方法是正确的,我仔细检查了它。这里出了什么问题?我还必须实现equals和hashCode吗?我认为相等是通过比较和检查比较结果是否为0来实现的。

还是我忘记了其他事情?

以下是关键实现:

密钥继承自的可写对象:

package somepackage;

import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;

import org.apache.hadoop.io.Writable;

public class SomeWritable implements Writable {

        private String _string1;
        private String _string2;

        public SomeWritable() {
                super();
        }

        public String getString1() {
                return _string1;
        }

        public void setString1(final String string1) {
                _string1 = string1;
        }

        public String getString2() {
                return _string2;
        }

        public void setString2(final String string2) {
                _string2 = string2;
        }

        @Override
        public void write(final DataOutput out) throws IOException {
                out.writeUTF(_string1);
                out.writeUTF(_string2);
        }

        @Override
        public void readFields(final DataInput in) throws IOException {
                _string1 = in.readUTF();
                _string2 = in.readUTF();
        }
}

我使用的钥匙:

package somepackage;

import static org.apache.commons.lang.ObjectUtils.compare;

import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;

import org.apache.hadoop.io.WritableComparable;

public class SomeKey extends SomeWritable implements
                WritableComparable<SomeKey> {

        private String _someOtherString;

        public String getSomeOtherString() {
                return _someOtherString;
        }

        public void setSomeOtherString(final String someOtherString) {
                _someOtherString = someOtherString;
        }

        @Override
        public void write(final DataOutput out) throws IOException {
                super.write(out);
                out.writeUTF(_someOtherString);
        }

        @Override
        public void readFields(final DataInput in) throws IOException {
                super.readFields(in);
                _someOtherString = in.readUTF();
        }

        @Override
        public int compareTo(final SomeKey o) {
                if (o == null) {
                        return 1;
                }
                if (o == this) {
                        return 0;
                }
                final int c1 = compare(_someOtherString, o._someOtherString);
                if (c1 != 0) {
                        return c1;
                }
                final int c2 = compare(getString1(), o.getString1());
                if (c2 != 0) {
                        return c2;
                }
                return compare(getString2(), o.getString2());
        }
}

【问题讨论】:

  • 向我们展示您的 Key 实现。
  • 不幸的是,我不允许这样做,因为它是公司代码。但是密钥扩展了另一个Writable,它不是WritableComparable,类似于public class K1 extends OtherK implements WritableComparable&lt;K1&gt;。这可能是问题吗?
  • 屏蔽这些东西,更改变量/名称/包等,没有它就粘贴我们无能为力。

标签: hadoop mapreduce


【解决方案1】:

我解决了这个问题:为确保始终将相同的 key 分配给同一个 reducer,hashCode() 的 key 必须基于 key 中的 current 值来实现。即使它们是可变的。有了这个,一切正常。

然后必须非常小心,不要在集合中使用这些类型或在地图等中使用这些类型。

【讨论】:

    猜你喜欢
    • 2018-12-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-12
    • 2022-06-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多